【问题标题】:KeyValueStore.get() returns inconsistent resultsKeyValueStore.get() 返回不一致的结果
【发布时间】: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


【解决方案1】:

似乎消息正在消失的原因(假设标点符号没有删除它们)是您使用 KStream::selectKey(...),它更改密钥,但不进行重新分区 您可能会在错误的分区中查找密钥。

看下面的场景:

  • Msg1:k1v1 (partition0)
  • 消息2:k2v2 (partition1)

假设消息被放在不同的分区中(因为键) selectKey 后:k1 -&gt; k,k2 -&gt; k

  • 消息1:kv1
  • 消息2:kv2

操作selectKey 是无状态的,因此消息不会发送到下游(主题)并且不会发生重新分区。 对于第一条消息:将值放在存储区 (partition0) 中的键 - k 当第二条消息到达时:对于key - k 没有消息,因为它是不同的分区(partition1)

【讨论】:

    猜你喜欢
    • 2016-03-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-09-08
    相关资源
    最近更新 更多