diff --git a/input_kafka.go b/input_kafka.go index 5420340..48c890f 100644 --- a/input_kafka.go +++ b/input_kafka.go @@ -8,17 +8,18 @@ import ( "github.com/Shopify/sarama" ) -// KafkaInput should make consumer client. +// KafkaInput is used for recieving Kafka messages and +// transforming them into HTTP payloads. type KafkaInput struct { config *KafkaConfig consumers []sarama.PartitionConsumer messages chan *sarama.ConsumerMessage } -// NewKafkaInput constructor for KafkaInput +// NewKafkaInput creates instance of kafka consumer client. func NewKafkaInput(address string, config *KafkaConfig) *KafkaInput { c := sarama.NewConfig() - // c.Consumer. + // Configuration options go here con, err := sarama.NewConsumer([]string{config.host}, c) if err != nil { diff --git a/output_kafka.go b/output_kafka.go index 77c4a35..35a0d67 100644 --- a/output_kafka.go +++ b/output_kafka.go @@ -11,7 +11,7 @@ import ( "github.com/buger/gor/proto" ) -// KafkaOutput should make producer client. +// KafkaOutput is used for sending payloads to kafka in JSON format. type KafkaOutput struct { config *KafkaConfig producer sarama.AsyncProducer