【发布时间】:2020-03-20 15:13:44
【问题描述】:
在 kafka 流中定义拓扑时,可以添加全局状态存储。它需要一个源主题以及一个ProcessorSupplier。
处理器接收记录并可以在将它们添加到存储之前从理论上对其进行转换。但是在恢复的情况下,记录会直接从源主题(更改日志)插入到全局状态存储中,跳过处理器中完成的最终转换。
+-------------+ +-------------+ +---------------+
| | | | | global |
|source topic -------------> processor +--------------> state |
|(changelog) | | | | store |
+-------------+ +-------------+ +---------------+
| ^
| |
+---------------------------------------------------------+
record directly inserted during restoration
StreamsBuilder#addGlobalStore(StoreBuilder storeBuilder, String topic, Consumed consumed, ProcessorSupplier stateUpdateSupplier) 将全局 StateStore 添加到拓扑。
根据文档
注意:您不应使用处理器将转换后的记录插入到全局状态存储中。该存储使用源主题作为更改日志,并且在恢复期间将插入直接来自源的记录。这个 ProcessorNode 应该用来保持 StateStore 是最新的。
与此同时,主要错误目前在 kafka 错误跟踪器上打开:KAFKA-7663 Custom Processor supplied on addGlobalStore is not used when restoring state from topic 准确解释了文档中的说明,但似乎是一个公认的错误。
我想知道 KAFKA-7663 是否确实是一个错误。根据文档,它似乎是这样设计的,在这种情况下,我很难理解用例。
有人能解释一下这个低级 API 的主要用例吗?我唯一能想到的就是处理副作用,例如在处理器中执行一些日志操作。
额外问题:如果源主题作为全局存储的更新日志,当一条记录因为保留期已过期而从主题中删除时,它会从全局状态存储中删除吗?还是只有在从更改日志中完全恢复商店后才会在商店中进行删除。
【问题讨论】:
-
请注意,旧文档没有指出问题,我们只是将文档更新为“中间修复”。
标签: java apache-kafka apache-kafka-streams