【发布时间】:2021-06-19 11:05:55
【问题描述】:
我有两个话题:
- 1 个带有事件数据的主题(
EventData,比如说 5 个分区)- 该主题的日志使用 CustomerID 作为键。 - 1 个 compact 主题,包含丰富的数据(
EnrichmentKVs,比如说 3 个分区)- 该主题的日志使用相同的 CustomerID 作为键。
目标是将 EnrichmentKV 保存在 Faust 表中,当 EventData 日志流入时,它们会使用该表中的数据进行丰富并发布到新的流/主题。
所以我有两个 Faust (python) 应用程序,每个应用程序都有自己运行的实例数量:
- App1(N-instances running)使用 key=CustomerId 发布到 EventData 主题
- App2(M 实例运行)执行以下操作:
- 更新浮士德表 (
EnrichmentKVsTable) 以获取来自 EnrichmentKVs 主题的值 - 来自 EventData 主题的流式输入,并将来自浮士德表的数据与来自
Eventdata的数据流“连接”起来
- 更新浮士德表 (
我的理解是,App2 的每个实例都只会有一部分基于分区键的 EnrichmentKV 表。要使“JOIN”工作,EventData(key="1234") 的任何日志都必须与 EnrichmentKVsTable(key="1234") 的日志转到相同的 App2 instance
当两个输入主题的分区不同,并且每个应用程序的实例数量也可能不同时,浮士德如何确保这一点?还是我处理这个问题有误?
【问题讨论】:
标签: python apache-kafka faust