【发布时间】:2021-01-21 11:25:15
【问题描述】:
我们将 KTable 实现为 Internal-State-Store。
a.) 我如何以及在哪里可以指定,这个 Internal-State-Store 应该是 Persistent 并自动备份到另一个 kafka 主题?
b.) 我们如何指定这个 Internal-State-Store 应该是全局的,这样我的任何流任务都应该能够引用它?
c.) 传入的 messageRecords 是否有频率写入 Internal-State-Store ?会不会发生这样的情况,一个特定的 MessageRecord 被 Stream-processor 处理,存储在 KTable 中,然后我的 stream-processor 死了,它无法进入 Internal-State-Store !
在我们现在使用的 sn-p 之下 :-
KTable<String, String> KT0 = streamsBuilder.table(AppConfigs.topicName, Materialized.as(AppConfigs.stateStoreName)));
任何回复都将受到高度赞赏!
【问题讨论】:
标签: apache-kafka apache-kafka-streams ktable