【发布时间】:2019-05-09 17:57:24
【问题描述】:
我有一个架构,其中我们有两个独立的应用程序。原始来源是一个sql数据库。 App1 侦听 CDC 表以跟踪对该数据库中表的更改、规范化和序列化这些更改。它获取这些序列化消息并将它们发送到 Kafka 主题。 App2 监听该主题,将消息调整为不同的格式,并通过 HTTP 将这些调整后的消息发送到各自的目的地。
所以我们的流式架构如下所示:
SQL(CDC 事件)-> App1(规范化事件)-> Kafka -> App2(使事件适应端点)-> 各种端点
我们希望在失败的情况下添加错误处理,并且不能容忍重复事件、丢失事件或更改顺序。鉴于上述架构,我们真正关心的是,Exactly-once 应用于从 App1 到 App2(我们独立的生产者和消费者)的消息
我正在阅读的所有内容以及我发现的有关事务性 API 的每个示例都指向“流式传输”。看起来 Kafka 流 api 是为单个应用程序设计的,该应用程序从 Kafka 主题中获取输入,进行处理,然后将其输出到另一个 Kafka 主题,这似乎不适用于我们对 Kafka 的使用。这是Confluent's docs的摘录:
现在,流处理只不过是一个读-处理-写操作 关于 Kafka 主题;消费者从 Kafka 主题中读取消息,一些 处理逻辑转换这些消息或修改状态 由处理器维护,生产者写入结果 消息到另一个 Kafka 主题。正是在流处理 简单地执行一个读-处理-写操作的能力 一度。在这种情况下,“得到正确的答案”意味着不会错过 任何输入消息或产生任何重复输出。这是 用户期望从恰好一次流处理器获得的行为。
我正在努力思考如何在 Kafka 主题中使用精确一次,或者如果 Kafka 的精确一次甚至是为非“流”用例构建的。我们是否必须建立自己的重复数据删除和容错能力?
【问题讨论】:
标签: java apache-kafka apache-kafka-streams