mirror of
https://github.com/buger/goreplay.git
synced 2024-04-21 12:32:02 +00:00
+2
-7
@@ -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
|
||||
|
||||
+2
-2
@@ -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)
|
||||
|
||||
+2
-7
@@ -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
|
||||
|
||||
@@ -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"))
|
||||
|
||||
|
||||
Reference in New Issue
Block a user