【问题标题】:Is 'exactly once' only for streams (topic1 -> app -> topic2)?“恰好一次”仅适用于流(主题 1 -> 应用程序 -> 主题 2)?
【发布时间】: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


    【解决方案1】:

    如果您使用的是 Kafka 的 Streams API(或其他支持使用 Kafka 进行精确一次处理的工具),那么 Kafka 的精确一次语义 (EOS) 将涵盖所有应用程序:

    topic A --> App 1 --> topic B --> App 2 --> topic C
    

    在您的用例中,一个问题是初始 CDC 步骤是否也支持 EOS。也就是说,你必须要问一个问题:涉及到哪些步骤,EOS涵盖了所有步骤?

    在以下示例中,当(且仅当)初始 CDC 步骤也支持 EOS 时,端到端支持 EOS,就像数据流的其余部分一样。

    SQL --CDC--> topic A --> App 1 --> topic B --> App 2 --> topic C
    

    如果你在CDC步骤中使用Kafka Connect,那么你必须检查你使用的连接器是否支持EOS。

    我正在阅读的所有内容以及我发现的有关事务性 API 的每个示例都指向“流式传输”。

    Kafka 生产者/消费者客户端的事务 API 为 EOS 处理提供了原语。位于生产者/消费者客户端之上的 Kafka Streams 使用此功能来实现 EOS,开发人员只需几行代码即可轻松使用它(例如在应用程序需要时自动处理状态管理)进行有状态的操作,如聚合或连接)。也许生产者/消费者之间的关系 Kafka Streams 是您阅读文档后的困惑?

    当然,您也可以在开发应用程序时使用底层的 Kafka 生产者和消费者客户端(使用事务性 API)“构建自己的”,但这需要更多的工作。

    我正在努力思考如何在 Kafka 主题中使用完全一次,或者如果 Kafka 的完全一次是为非“流式”用例而构建的。我们是否必须建立自己的重复数据删除和容错能力?

    不确定您所说的“非流式传输”用例是什么意思。如果你的意思是,“如果我们不想使用 Kafka Streams 或 KSQL(或其他可以从 Kafka 读取数据来处理数据的现有工具),我们需要做什么来在我们的应用程序中实现 EOS?”,那么答案是“是的,在这种情况下,您必须直接使用 Kafka 生产者/客户端,并确保您对它们所做的任何事情都能正确实现 EOS 处理。” (而且因为后者比较困难,所以将这个 EOS 功能添加到 Kafka Streams 中。)

    希望对你有帮助。

    【讨论】:

    • 谢谢。当我说非流媒体时,我指的是 app1 -> topic -> app2 而不是 topic1 -> app -> topic2。我最关心的只是 app1 -> topic -> app2 的 EOS,但没有人在这种情况下谈论“恰好一次”,这让我感到困惑。是否有示例说明如何使用生产者和消费者 api 来确保 app1 -> 主题 -> app2 只发生一次?
    • 我想知道对于我的情况,忽略我的架构中的 sql 和 http 部分,在生产者上启用幂等性并在消费者上启用 read_commited 是否就足够了。
    • app1 -> topic -> app2 也被支持。它与 topic1 -> app1 -> topic2 -> app2 没有什么不同。例如,App1 应该使用 Kafka Streams 或 Kafka 生产者客户端(幂等生产者)的事务 API。 App2 应酌情使用 Kafka Streams 或 Kafka 消费者客户端(具有 EOS 相关功能)。
    猜你喜欢
    • 2016-12-23
    • 2013-09-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-01-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多