【问题标题】:Timeout for aggregated records in Kafka table?Kafka表中聚合记录的超时?
【发布时间】:2019-01-08 21:56:54
【问题描述】:

我使用 Kafka 来处理消息。消息可以分为几个部分(这是一个复合消息)。因此,在流中,我可以拥有例如一条复合消息,该消息分为三个部分。换句话说,它将是 Kafka 流中的三条记录,但它是一条大消息。我想使用 Kafka 表在一个 Kafka 记录中合并部分复合消息。合并后,一条消息将插入数据库(Postgres)。每个零件都有零件的数量和总数。例如,如果我在流中有一条消息的三个部分(三个 Kafka 记录) - 每个部分的字段总数为 3。

我的理解是,在积极的情况下,任务很简单:聚合表中的部分,从表中创建流并过滤具有等于聚合部分大小和部分总数的记录,在一条合并消息中过滤映射并将其插入数据库中( Postgres)。

但消极的情况也是可能的。在极少数情况下,其中一个部件根本无法插入到 Kafka 中(或者它会在超时后很久才插入)。因此,例如在流中,只有一个复合消息的三个中的两个部分会出现。在这种情况下,我必须在数据库(Postgres)中插入未完全构造的消息(它将仅包含两部分,而不是三部分)。如何在 Kafka 中实现这种负面场景?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    我建议查看标点符号:https://docs.confluent.io/current/streams/developer-guide/processor-api.html#defining-a-stream-processor

    另请注意,您可以混合搭配处理器 API 和 DSL:https://docs.confluent.io/current/streams/developer-guide/dsl-api.html#applying-processors-and-transformers-processor-api-integration

    如果您为 KTable 聚合提供商店名称,则可以将商店连接到注册标点符号的自定义处理器。总体而言,为整个应用程序使用处理器 API 而不是 DSL 可能会更好。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2023-03-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-04-26
      • 2017-09-03
      • 2021-07-21
      相关资源
      最近更新 更多