From c7810090da66d64475af39d317ad30a7c63523d1 Mon Sep 17 00:00:00 2001 From: fengzhiqiang Date: Mon, 7 Jan 2019 14:29:04 +0800 Subject: [PATCH] =?UTF-8?q?=E6=9C=8D=E5=8A=A1=E5=99=A8=E5=8F=82=E6=95=B0?= =?UTF-8?q?=EF=BC=8C=E5=AD=97=E7=AC=A6=E4=B8=B2=E8=BD=AC=E6=8D=A2=E4=B8=BA?= =?UTF-8?q?=E6=95=B0=E7=BB=84?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- serv.go | 14 +++++--------- 1 file changed, 5 insertions(+), 9 deletions(-) diff --git a/serv.go b/serv.go index 20c96ea..35939b8 100644 --- a/serv.go +++ b/serv.go @@ -6,6 +6,7 @@ import ( "net/http" "os" "os/signal" + "strings" "syscall" "time" @@ -33,7 +34,6 @@ const ( var ( newline = []byte{'\n'} - space = []byte{' '} addr = flag.String("addr", ":9000", "http service address") bootstrapServers = flag.String("bootstrap-servers", "localhost:9092", "kafka bootstrap servers") partition = flag.Int("partition", 0, "kafka tipic partion") @@ -76,21 +76,17 @@ func (c *Client) readPump(topic string) { log.WithFields(log.Fields{ "topic": topic, "bootstrap servers": *bootstrapServers, + "partition": partition, "group": group}).Info("create kafka client") - // consumer, err := kafka.NewConsumer(&kafka.ConfigMap{ - // "bootstrap.servers": *bootstrapServers, - // "group.id": group, - // "session.timeout.ms": 6000, - // "auto.offset.reset": "latest"}) reader := kafka.NewReader(kafka.ReaderConfig{ - Brokers: []string{*bootstrapServers}, + Brokers: strings.Split(*bootstrapServers, ","), Topic: topic, - Partition: 0, + Partition: *partition, MinBytes: 10e1, MaxBytes: 10e6, // 10MB }) - reader.SetOffset(-1) //lastoffset + reader.SetOffset(-1) //latest offset log.WithFields(log.Fields{ "reader": reader,