【发布时间】:2020-01-24 05:52:21
【问题描述】:
Kafka 文档说,幂等生产者可以使用相同的生产者会话,我无法理解这一点。
说,Kafka为每条消息添加序列号,最后一个序列号在Kafka中维护(不确定它在哪里维护)。
它是如何生成序列号的,它保存在哪里?
为什么生产者崩溃再上来时不能保持顺序?
如何在生产者会话之间实现真正的幂等性?
【问题讨论】:
标签: apache-kafka
Kafka 文档说,幂等生产者可以使用相同的生产者会话,我无法理解这一点。
说,Kafka为每条消息添加序列号,最后一个序列号在Kafka中维护(不确定它在哪里维护)。
它是如何生成序列号的,它保存在哪里?
为什么生产者崩溃再上来时不能保持顺序?
如何在生产者会话之间实现真正的幂等性?
【问题讨论】:
标签: apache-kafka
幂等生产者只在生产者进程的生命周期内有保证。如果崩溃,新的 Idempotent Producer 将具有不同的 ProducerId 并开始自己的序列。
序列号只是从 0 开始,每条记录单调递增。如果一条记录未能交付,它会使用其现有的序列号再次发送,以便代理可以对其进行重复数据删除(如果需要)。序列号是每个生产者和每个分区的。
目前 Kafka 不提供“继续” Idempotent Producer 会话的方法。每次启动它都会获得一个新的唯一 ProducerId(由集群生成)
【讨论】:
这是 Kafka 严重缺失的一个功能,我没有看到一个优雅而有效的方法来解决它而不修改 Kafka 本身。
作为初步,如果您希望在任何故障(生产者或代理)中实现真正的幂等性,那么您 absolutely positively need 某种业务层中的id(而不是较低级别的传输层) .
您可以在 Kafka 中使用这样的 id 来执行以下操作:您的生产者写入一个主题至少一次,然后您有一个 Kafka Streams 进程使用您的业务从该主题删除重复的消息layer id 并将剩余的唯一消息发布到另一个主题。为了提高效率,您应该使用单调递增的 id,也就是 sequence number,否则您将不得不保留(并坚持)您见过的每个 id,这相当于内存泄漏,除非您将重复数据删除功能限制为最近 x 天/小时/分钟,并且仅保留最新的 ID。
或者,您试试 Apache Pulsar,除了解决 Kafka 的其他痛点(必须进行昂贵的手动和容易出错的重新平衡以扩展主题,仅举几例)一)has this feature built in.
【讨论】:
“幂等”配置只在生产者不崩溃的情况下有效。
但是,使用事务,您可以在不同的分区上只发送一次数据。 您使用生产者 ID 设置交易 ID(自动创建)。 如果新的生产者 id 到达时具有相同的事务 id,则意味着您遇到了问题。 然后,记录将只写入一次。
【讨论】: