【发布时间】: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是三个不同的数据。
这是扫描仪:
-
在将批次结束记录
B1E推送到主题之前,我们知道B1批次经过验证是无效的。 -
在这种情况下,
B1批次的所有数据从BS到BE都应该转到特定主题。
所以在上面的例子中,b1 批处理应该转到主题 T1,b2 批处理应该转到 T2。
我如何使用 Kafka 做到这一点?
【问题讨论】:
-
使用一些窗口机制来存储批处理数据,直到处理完批处理结束记录
BE。验证所有批次后,将记录发布到相关主题。 -
你能分享一些例子吗
标签: apache-kafka confluent-platform