使用KStream.processValues()时FixedKeyProcessor获取状态存储为null的问题
我构建了如下拓扑,使用processValues()方法将Streams DSL与Processor API结合,并添加了状态存储。
KStream<String, SecurityCommand> securityCommands = builder.stream( "security-command", Consumed.with(Serdes.String(), JsonSerdes.securityCommand())); StoreBuilder<KeyValueStore<String, UserAccountSnapshot>> storeBuilder = Stores.keyValueStoreBuilder( Stores.persistentKeyValueStore("user-account-snapshot"), Serdes.String(), JsonSerdes.userAccountSnapshot()); builder.addStateStore(storeBuilder); securityCommands.processValues(() -> new SecurityCommandProcessor(), Named.as("security-command-processor"), "user-account-snapshot") .processValues(() -> new UserAccountSnapshotUpdater(), Named.as("user-snapshot-updater"), "user-account-snapshot") .to("security-event", Produced.with( Serdes.String(), JsonSerdes.userAccountEvent()));
SecurityCommandProcessor代码如下:
class SecurityCommandProcessor implements FixedKeyProcessor<String, SecurityCommand, UserAccountEvent> { private KeyValueStore<String, UserAccountSnapshot> kvStore; private FixedKeyProcessorContext context; @Override public void init(FixedKeyProcessorContext context) { this.kvStore = (KeyValueStore<String, UserAccountSnapshot>) context.getStateStore("user-account-snapshot"); this.context = context; } ... }
问题在于context.getStateStore("user-account-snapshot")返回null。我尝试使用已废弃的transformValues()实现几乎相同的逻辑时,可以正常获取状态存储,只有使用processValues()时出现此问题,请问我哪里操作有误?
问题原因与解决方案
问题出在processValues()的参数传递方式上——你当前把状态存储名称作为单独参数传递,但实际上processValues()的重载方法中,状态存储名称需要通过Named对象来关联,而不是作为独立参数传入。
正确的写法应该是通过Named.withProcessorStateStores()方法来指定当前处理器要使用的状态存储:
securityCommands.processValues( () -> new SecurityCommandProcessor(), Named.as("security-command-processor") .withProcessorStateStores("user-account-snapshot") ) .processValues( () -> new UserAccountSnapshotUpdater(), Named.as("user-snapshot-updater") .withProcessorStateStores("user-account-snapshot") ) .to("security-event", Produced.with(Serdes.String(), JsonSerdes.userAccountEvent()));
为什么之前的写法不行?
processValues()方法中,如果你直接把状态存储名称作为最后一个参数传入,这个参数对应的是父处理器的状态存储(用于状态迁移等场景),而不是当前处理器要绑定的状态存储。而废弃的transformValues()方法的参数逻辑和processValues()不同,它直接接受状态存储名称作为绑定参数,所以之前用transformValues()是正常的。
通过Named.withProcessorStateStores()明确指定当前处理器要关联的状态存储后,FixedKeyProcessorContext就能正确获取到对应的存储实例了。
内容的提问来源于stack exchange,提问作者Markus

