【发布时间】:2019-10-04 15:29:27
【问题描述】:
我试图了解 Kafka 的事务 API。 This link定义原子读-进程-写周期如下:
首先,让我们考虑一下原子读-进程-写周期的含义。简而言之,这意味着如果应用程序在某个主题分区 tp0 的偏移量 X 处消费消息 A,并在对消息 A 进行一些处理使得 B = F(A) 后将消息 B 写入主题分区 tp1,则仅当消息 A 和 B 被视为成功使用并一起发布或根本不一起发布时,读取-处理-写入周期才是原子的。
它进一步说:
使用为至少一次交付语义配置的普通 Kafka 生产者和消费者,流处理应用程序可能会以下列方式丢失恰好一次处理语义:
producer.send() 可能会由于内部重试而导致重复写入消息 B。这由幂等生产者解决,不是本文其余部分的重点。
我们可能会重新处理输入消息 A,从而导致重复的 B 消息被写入输出,从而违反了仅处理一次的语义。如果流处理应用程序在写入 B 之后但在将 A 标记为已使用之前崩溃,则可能会发生重新处理。因此当它恢复时,它会再次消耗A并再次写入B,导致重复。
最后,在分布式环境中,应用程序将崩溃,或者——更糟糕的是!——暂时失去与系统其余部分的连接。通常,新实例会自动启动以替换那些被认为丢失的实例。通过这个过程,我们可能有多个实例处理相同的输入主题并写入相同的输出主题,从而导致重复输出并违反恰好一次处理的语义。我们称之为“僵尸实例”问题。
我们在 Kafka 中设计了事务 API 来解决第二个和第三个问题。事务通过使这些循环原子化并促进僵尸防护,从而在读取-处理-写入循环中实现一次性处理。
疑问:
上面的第 2 点和第 3 点描述了何时可能发生消息重复,这些消息重复使用事务 API 进行处理。事务 API 是否也有助于在任何情况下避免消息丢失?
-
大多数在线(例如,here 和 here)的 Kafka 事务 API 示例涉及:
while (true) { ConsumerRecords records = consumer.poll(Long.MAX_VALUE); producer.beginTransaction(); for (ConsumerRecord record : records) producer.send(producerRecord(“outputTopic”, record)); producer.sendOffsetsToTransaction(currentOffsets(consumer), group); producer.commitTransaction(); }这基本上是读-处理-写循环。那么事务 API 是否只在读-写-写循环中有用?
-
This 文章给出了非读写场景下的事务 API 示例:
producer.initTransactions(); try { producer.beginTransaction(); producer.send(record1); producer.send(record2); producer.commitTransaction(); } catch(ProducerFencedException e) { producer.close(); } catch(KafkaException e) { producer.abortTransaction(); }上面写着:
这允许生产者将一批消息发送到多个分区,以便批处理中的所有消息最终对任何消费者可见,或者对消费者不可见。
这个例子是否正确,并展示了另一种使用事务 API 的方式,不同于读-处理-写循环? (请注意,它也不会向事务提交偏移量。)
-
在我的应用程序中,我只是使用来自 kafka 的消息,进行处理并将它们记录到数据库中。那是我的整个管道。
一个。所以,我猜这不是 read-process-write 循环。 Kafka 事务 API 对我的场景有用吗?
b.我还需要确保每条消息只处理一次。我想在 producer 中设置
idempotent=true就足够了,我不需要事务 API,对吧?c。我可能会运行多个管道实例,但我不会将处理输出写入 Kafka。所以我想这永远不会涉及僵尸(重复的生产者写信给卡夫卡)。所以,我想事务 API 不会帮助我避免重复处理场景,对吧? (我可能必须在同一个数据库事务中将偏移量和处理输出保存到数据库中,并在生产者重启期间读取偏移量以避免重复处理。)
【问题讨论】:
标签: apache-kafka kafka-producer-api kafka-transactions-api