Kafka:Sarama、幂等性和 transactional.id

Tim*_*rib 7 go apache-kafka kafka-producer-api sarama

Shopify/sarama是否提供类似于transactional.idJVM API的选项?

该库支持幂等(Config.Producer.Idemponent,类似于enable.idempotence),但我不明白如何在没有transactional.id.

如果我错了,请纠正我,Sarama 中缺少有关这些选项的文档。但是根据 JVM 文档,没有标识符的幂等性将受到单个生产者会话的限制。换句话说,当生产者失败并重新启动时,我们将失去保证。

我在源代码和一些测试(例如)中找到了相关属性,但不明白如何在外部使用它们。

nip*_*una 3

Shopify/sarama为 Kafka Exactly Once(幂等性)提供幂等启用的生产者。但为此,需要进行以下配置设置。

\n

来自Shopify/sarama/config.go

\n
    if c.Producer.Idempotent {\n        if !c.Version.IsAtLeast(V0_11_0_0) {\n            return ConfigurationError("Idempotent producer requires Version >= V0_11_0_0")\n        }\n        if c.Producer.Retry.Max == 0 {\n            return ConfigurationError("Idempotent producer requires Producer.Retry.Max >= 1")\n        }\n        if c.Producer.RequiredAcks != WaitForAll {\n            return ConfigurationError("Idempotent producer requires Producer.RequiredAcks to be WaitForAll")\n        }\n        if c.Net.MaxOpenRequests > 1 {\n            return ConfigurationError("Idempotent producer requires Net.MaxOpenRequests to be 1")\n        }\n    }\n
Run Code Online (Sandbox Code Playgroud)\n

Shopify/sarama中,他们是如何做到这一点的, \'sproducerEpoch中有一个ID 。您可以参考Shopify/sarama/async_ Producer.go中的文件。该 Id 通过生产者初始化进行初始化,并在成功生成每条消息时递增。读取函数以查看文件中的内容。AsyncProducertransactionManagerbumpEpoch()async_producer.go

\n

这是生产者与代理会话的序列 ID,它随每条消息一起发送。消息发布成功时递增。

\n

阅读这个例子。它描述了幂等性的工作原理。

\n

您对制作人会议事实的看法是正确的。这正是曾经为单一制作人会议所承诺的。当序列失败后重新声明生产者时,可能会出现重复。

\n
\n

当生产者重新启动时,会分配新的 PID。因此,仅对单个生产者会话承诺幂等性。即使生产者在失败时重试请求,每条消息也会在日志中保留一次。根据生产者获取数据的来源,仍然可能存在重复项。Kafka 不会\xe2\x80\x99t 处理生产者收到的重复数据。因此,在某些情况下,您可能需要额外的重复数据删除系统。

\n
\n