【问题标题】:Perform a batch validation in Kafka and sent to corresponding topic在 Kafka 中执行批量验证并发送到相应的主题
【发布时间】:2021-12-17 12:52:58
【问题描述】:

我在 Kafka 主题中存储了以下批处理格式:

Data generated -->  B2E, T3, T2, T1, B2S | B1E, T3, T2, T1, B1S  --> Data Consumed

这里BS表示批次开始,BE表示批次结束,t1,t2,t3是三个不同的数据。

这是扫描仪:

  1. 在将批次结束记录B1E推送到主题之前,我们知道B1批次经过验证是无效的。

  2. 在这种情况下,B1 批次的所有数据从BSBE 都应该转到特定主题。

所以在上面的例子中,b1 批处理应该转到主题 T1,b2 批处理应该转到 T2。

我如何使用 Kafka 做到这一点?

【问题讨论】:

  • 使用一些窗口机制来存储批处理数据,直到处理完批处理结束记录BE 。验证所有批次后,将记录发布到相关主题。
  • 你能分享一些例子吗

标签: apache-kafka confluent-platform


【解决方案1】:

根据上图,逐条读取和验证每条消息,并在内存窗口中更新一些(带有批处理数据的内存映射),直到收到相关的批处理结束BE。批处理结束后BE从窗口收到读取相关批处理数据,并将所有批处理记录发布到从处理(验证)结果中选择的主题。

要窗口化,您可以使用像 或 Kafka 流式处理 state store 这样的内存映射。

如果您需要类似 KStream 的解决方案,它将是如下所示的流

【讨论】:

  • 这里的内存窗口是简单的Map??
  • 是的。您可以使用简单的内存映射数据结构或一些存储实现。
  • 我正在寻找基于 Kafka 的解决方案
  • 数据消费->处理->根据处理结果发布。因此,如果您需要等待批处理结束记录,则需要将它们窗口化到某个状态存储中。这是 Kafka Streaming 可能的解决方案。您可以使用 Stream State Store 或 KTable 之类的东西。
  • 能否请您推荐 Kafka 解决方案?
猜你喜欢
  • 2017-06-12
  • 2013-09-28
  • 2022-11-03
  • 2019-09-21
  • 2019-04-30
  • 2019-06-28
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多