From 358cb78a536cd9c6a1cc9880fbd9c8a5d262cf3e Mon Sep 17 00:00:00 2001 From: David Fradin Date: Mon, 19 Oct 2020 09:56:07 +0200 Subject: [PATCH] Fix #835: kafka bug (#836) Correction of error introduce in PR #800 --- input_kafka.go | 9 ++------- input_kafka_test.go | 4 ++-- output_kafka.go | 9 ++------- output_kafka_test.go | 4 ++-- 4 files changed, 8 insertions(+), 18 deletions(-) diff --git a/input_kafka.go b/input_kafka.go index d168b7a..10f426f 100644 --- a/input_kafka.go +++ b/input_kafka.go @@ -17,13 +17,8 @@ type KafkaInput struct { messages chan *sarama.ConsumerMessage } -// NewKafkaInput creates instance of kafka consumer client. -func NewKafkaInput(address string, config *InputKafkaConfig) *KafkaInput { - return NewKafkaInputWithTLS(address, config, nil) -} - -// NewKafkaInputWithTLS creates instance of kafka consumer client with TLS -func NewKafkaInputWithTLS(address string, config *InputKafkaConfig, tlsConfig *KafkaTLSConfig) *KafkaInput { +// NewKafkaInput creates instance of kafka consumer client with TLS config +func NewKafkaInput(address string, config *InputKafkaConfig, tlsConfig *KafkaTLSConfig) *KafkaInput { c := NewKafkaConfig(tlsConfig) var con sarama.Consumer diff --git a/input_kafka_test.go b/input_kafka_test.go index 2c6be17..544da6b 100644 --- a/input_kafka_test.go +++ b/input_kafka_test.go @@ -20,7 +20,7 @@ func TestInputKafkaRAW(t *testing.T) { consumer: consumer, Topic: "test", UseJSON: false, - }) + },nil) buf := make([]byte, 1024) n, err := input.Read(buf) @@ -47,7 +47,7 @@ func TestInputKafkaJSON(t *testing.T) { consumer: consumer, Topic: "test", UseJSON: true, - }) + },nil) buf := make([]byte, 1024) n, err := input.Read(buf) diff --git a/output_kafka.go b/output_kafka.go index ecf46b6..61e7f30 100644 --- a/output_kafka.go +++ b/output_kafka.go @@ -23,13 +23,8 @@ type KafkaOutput struct { // KafkaOutputFrequency in milliseconds const KafkaOutputFrequency = 500 -// NewKafkaOutput creates instance of kafka producer client. -func NewKafkaOutput(address string, config *OutputKafkaConfig) io.Writer { - return NewKafkaOutputWithTLS(address, config, nil) -} - -// NewKafkaOutputWithTLS creates instance of kafka producer client. -func NewKafkaOutputWithTLS(address string, config *OutputKafkaConfig, tlsConfig *KafkaTLSConfig) io.Writer { +// NewKafkaOutput creates instance of kafka producer client with TLS config +func NewKafkaOutput (address string, config *OutputKafkaConfig, tlsConfig *KafkaTLSConfig) io.Writer { c := NewKafkaConfig(tlsConfig) var producer sarama.AsyncProducer diff --git a/output_kafka_test.go b/output_kafka_test.go index 87a98be..67fb464 100644 --- a/output_kafka_test.go +++ b/output_kafka_test.go @@ -17,7 +17,7 @@ func TestOutputKafkaRAW(t *testing.T) { producer: producer, Topic: "test", UseJSON: false, - }) + },nil) output.Write([]byte("1 2 3\nGET / HTTP1.1\r\nHeader: 1\r\n\r\n")) @@ -40,7 +40,7 @@ func TestOutputKafkaJSON(t *testing.T) { producer: producer, Topic: "test", UseJSON: true, - }) + }, nil) output.Write([]byte("1 2 3\nGET / HTTP1.1\r\nHeader: 1\r\n\r\n"))