mirror of
https://github.com/buger/goreplay.git
synced 2024-04-21 12:32:02 +00:00
Fixed kafka client run of brokers and index out of range issues
This commit is contained in:
+2
-1
@@ -61,7 +61,8 @@ func CopyMulty(src io.Reader, writers ...io.Writer) (err error) {
|
|||||||
if nr > 0 && len(buf) > nr {
|
if nr > 0 && len(buf) > nr {
|
||||||
payload := buf[:nr]
|
payload := buf[:nr]
|
||||||
meta := payloadMeta(payload)
|
meta := payloadMeta(payload)
|
||||||
requestID := string(meta[1])
|
// requestID := string(meta[1])
|
||||||
|
requestID := string(meta[0])
|
||||||
|
|
||||||
_maxN := nr
|
_maxN := nr
|
||||||
if nr > 500 {
|
if nr > 500 {
|
||||||
|
|||||||
+3
-1
@@ -3,6 +3,7 @@ package main
|
|||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"log"
|
"log"
|
||||||
|
"strings"
|
||||||
|
|
||||||
"github.com/Shopify/sarama"
|
"github.com/Shopify/sarama"
|
||||||
"github.com/Shopify/sarama/mocks"
|
"github.com/Shopify/sarama/mocks"
|
||||||
@@ -27,7 +28,8 @@ func NewKafkaInput(address string, config *KafkaConfig) *KafkaInput {
|
|||||||
con = config.consumer
|
con = config.consumer
|
||||||
} else {
|
} else {
|
||||||
var err error
|
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 {
|
if err != nil {
|
||||||
log.Fatalln("Failed to start Sarama(Kafka) consumer:", err)
|
log.Fatalln("Failed to start Sarama(Kafka) consumer:", err)
|
||||||
|
|||||||
Reference in New Issue
Block a user