From a65256d70731b58ba5c42d1f3a5c2f13c6416409 Mon Sep 17 00:00:00 2001 From: Ubuntu Date: Tue, 27 Feb 2018 12:25:29 +0000 Subject: [PATCH] Fixed kafka client run of brokers and index out of range issues --- emitter.go | 3 ++- input_kafka.go | 4 +++- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/emitter.go b/emitter.go index 8e4180c..9a59735 100644 --- a/emitter.go +++ b/emitter.go @@ -61,7 +61,8 @@ func CopyMulty(src io.Reader, writers ...io.Writer) (err error) { if nr > 0 && len(buf) > nr { payload := buf[:nr] meta := payloadMeta(payload) - requestID := string(meta[1]) + // requestID := string(meta[1]) + requestID := string(meta[0]) _maxN := nr if nr > 500 { diff --git a/input_kafka.go b/input_kafka.go index e03fdbb..343cd11 100644 --- a/input_kafka.go +++ b/input_kafka.go @@ -3,6 +3,7 @@ package main import ( "encoding/json" "log" + "strings" "github.com/Shopify/sarama" "github.com/Shopify/sarama/mocks" @@ -27,7 +28,8 @@ func NewKafkaInput(address string, config *KafkaConfig) *KafkaInput { con = config.consumer } else { var err error - con, err = sarama.NewConsumer([]string{config.host}, c) + //con, err = sarama.NewConsumer([]string{config.host}, c) + con, err = sarama.NewConsumer(strings.Split(config.host,","), c) if err != nil { log.Fatalln("Failed to start Sarama(Kafka) consumer:", err)