doushenyu2537 2018-12-19 06:19
浏览 495

Golang segmentio / kafka-go消费者无法正常工作

I am using segmentio/kafka-go to connect to Kafka.

// to produce messages
topic := "my-topic"
partition := 0

conn, _ := kafka.DialLeader(context.Background(), "tcp", "localhost:9092", topic, partition)

conn.SetWriteDeadline(time.Now().Add(10*time.Second))
conn.WriteMessages(
    kafka.Message{Value: []byte("one!")},
    kafka.Message{Value: []byte("two!")},
    kafka.Message{Value: []byte("three!")},
)

conn.Close()

I am able to produce into my Kafka server using this code.

// to consume messages
topic := "my-topic"
partition := 0

conn, _ := kafka.DialLeader(context.Background(), "tcp", "localhost:9092", topic, partition)

conn.SetReadDeadline(time.Now().Add(10*time.Second))
batch := conn.ReadBatch(10e3, 1e6) // fetch 10KB min, 1MB max

b := make([]byte, 10e3) // 10KB max per message
for {
    _, err := batch.Read(b)
    if err != nil {
        // err -> "invalid codec"
        break
    }
    fmt.Println(string(b))
}

batch.Close()
conn.Close()

But I am unable to consume using the above code. I am getting the error invalid codec. What can be the reason?

In case relevant, I tweaked the minimum batch size to 1 so that it tries to consume something.

  • 写回答

1条回答 默认 最新

  • dsfhe34889 2019-09-04 10:56
    关注

    Just a guess: try adding an import to load compression codecs, in case your topics use compression.

    import _ "github.com/segmentio/kafka-go/snappy"

    评论

报告相同问题?

悬赏问题

  • ¥20 有关区间dp的问题求解
  • ¥15 多电路系统共用电源的串扰问题
  • ¥15 slam rangenet++配置
  • ¥15 有没有研究水声通信方面的帮我改俩matlab代码
  • ¥15 对于相关问题的求解与代码
  • ¥15 ubuntu子系统密码忘记
  • ¥15 信号傅里叶变换在matlab上遇到的小问题请求帮助
  • ¥15 保护模式-系统加载-段寄存器
  • ¥15 电脑桌面设定一个区域禁止鼠标操作
  • ¥15 求NPF226060磁芯的详细资料