【发布时间】:2015-04-02 10:58:12
【问题描述】:
所以我正在尝试将 Kafka 用于我的应用程序,该应用程序有一个生产者将操作记录到 Kafka MQ 中,而消费者从 MQ 中读取它。由于我的应用程序在 Go 中,我正在使用 Shopify Sarama 使这成为可能.
现在,我可以读取 MQ 并使用 a
打印消息内容fmt.Printf
不过,我真的希望错误处理比控制台打印更好,我愿意加倍努力。
现在用于消费者连接的代码:
mqCfg := sarama.NewConfig()
master, err := sarama.NewConsumer([]string{brokerConnect}, mqCfg)
if err != nil {
panic(err) // Don't want to panic when error occurs, instead handle it
}
以及消息的处理:
go func() {
defer wg.Done()
for message := range consumer.Messages() {
var msgContent Message
_ = json.Unmarshal(message.Value, &msgContent)
fmt.Printf("Reading message of type %s with id : %d\n", msgContent.Type, msgContent.ContentId) //Don't want to print it
}
}()
我的问题(我是测试 Kafka 的新手,也是一般的 kafka 新手):
上面的程序哪里可能出现错误,以便我处理?任何示例代码对我来说都是很好的开始。我能想到的错误情况是 msgContent 在 JSON 中实际上不包含任何 ContentId 类型的字段。
在 kafka 中,是否存在消费者试图读取当前偏移量但由于某种原因无法读取的情况(即使 JSON 格式正确)?我的消费者是否可以回溯说在失败的偏移读取之上的 x 步并重新处理偏移?还是有更好的方法来做到这一点?再说一遍,这些情况会是什么?
我乐于阅读和尝试事物。
【问题讨论】:
-
json.Unmarshal 可能会导致 err ,如果您不想引起恐慌......请不要:)
-
哈。谢谢。关于我如何做 #2 的任何想法?
标签: go shopify apache-kafka