【问题标题】:Kafka KTable Materialized-State-Store controlKafka KTable Materialized-State-Store 控制
【发布时间】: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


    【解决方案1】:

    a) 如果您有一个状态存储的自定义实现,您可以通过Materialized.as(KeyValueStoreSupplier) 传递它。

    b) 对于全局商店用例,您可以使用builder.globalKTable()

    c) 写入发生在每条记录的基础上,但可以缓存在内存中。在提交输入主题偏移量之前,状态存储将被刷新,因此您永远不会错过任何数据。默认情况下,KafkaStreams 提供至少一次处理语义。

    【讨论】:

      猜你喜欢
      • 2018-09-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-08-13
      • 2020-04-24
      • 2022-10-13
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多