【发布时间】:2020-10-02 17:31:48
【问题描述】:
stateStore.get() 在 KStream 上从 transform() 使用时返回不一致的结果。它返回 null,即使对应的键值已 put() 进入存储区。
谁能解释KeyValueStore<>的这种行为?
@Component
public class StreamProcessor {
@StreamListener
public void process(@Input(KStreamBindings.INPUT_STREAM) KStream<String, JsonNode> inputStream) {
KStream<String, JsonNode> joinedEvents = inputStream
.selectKey((key, value) -> computeKey(value))
.transform(
() -> new SelfJoinTransformer((v1, v2) -> join(v1, v2), "join_store"),
"join_store"
);
joinedEvents
.foreach((key, value) -> System.out.format("%s,joined=%b\n",key, value.has("right")));
}
private JsonNode join(JsonNode left, JsonNode right) {
((ObjectNode) left).set("right", right);
return left;
}
}
public class SelfJoinTransformer implements Transformer<String, JsonNode, KeyValue<String, JsonNode>> {
private KeyValueStore<String, JsonNode> stateStore;
private ValueJoiner<JsonNode, JsonNode, JsonNode> valueJoiner;
private String storeName;
public SelfJoinTransformer(ValueJoiner<JsonNode, JsonNode, JsonNode> valueJoiner, String storeName) {
this.storeName = storeName;
this.valueJoiner = valueJoiner;
}
@Override
public void init(ProcessorContext context) {
this.stateStore = (KeyValueStore<String, JsonNode>) context.getStateStore(storeName);
}
@Override
public KeyValue<String, JsonNode> transform(String key, JsonNode value) {
JsonNode oldValue = stateStore.get(key);
if (oldValue != null) { //this condition rarely holds true
stateStore.delete(key);
System.out.format("%s,joined\n", key);
return KeyValue.pair(key, valueJoiner.apply(oldValue, value));
}
stateStore.put(key, value);
return null;
}
}
【问题讨论】:
-
您能否添加示例消息流,何时以及发生什么?
-
在调用 transform(...) 之前是否使用
KStream::map更改密钥? -
在调用 transform 之前使用了 selectKey()
-
@BartoszWardziński 我很快就会添加流程
标签: apache-kafka apache-kafka-streams rocksdb spring-cloud-stream-binder-kafka