【问题标题】:How to get a sorted KeyValueStore from a KTable?如何从 KTable 中获取已排序的 KeyValueStore?
【发布时间】:2019-09-02 07:20:09
【问题描述】:

我想从 KStream 中具体化一个 KTable,并且我希望 KeyValueStore 按 Key 排序。

我尝试查找 KTable API 规范 (https://kafka.apache.org/20/javadoc/org/apache/kafka/streams/kstream/KTable.html),但不存在“排序”方法。我还查阅了这篇文章 (https://dzone.com/articles/how-to-order-streamed-dataframes),该文章建议通过处理器 API 实现排序。但是,我正在检查是否可以通过其他方式实现?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams rocksdb


    【解决方案1】:

    KafkaStream 允许您物化可查询的状态存储。 然后,您可以通过调用方法kafkaStream#store() 获得对存储的只读访问权限。

    如果您定义持久存储,KafkaStreams 将使用 RocksDB 来存储您的数据。返回的 KeyValueIterator 实例将使用 RocksDB 迭代器,它允许您以排序方式迭代键值Rocks Iterator-Implementation

    例子:

        KafkaStreams streams = new KafkaStreams(topology, props);
        ReadOnlyKeyValueStore<Object, Object> store = streams.store("storeName", QueryableStoreTypes.keyValueStore());
        KeyValueIterator<Object, Object> iterator = store.all();
    

    【讨论】:

      【解决方案2】:

      使用密钥将事件添加到 StateStore。 StateStore 返回的 KeyValueIterator 以有序的方式导航 KeyValue。

      public class SortProcessor extends AbstractProcessor<String, Event> {
      
          private static Logger LOG = LoggerFactory.getLogger(SortProcessor.class);
          private final String stateStore;
          private final Long bufferIntervalInSeconds;
      
          // Why not use a simple Java NavigableMap? Check out my answer at : https://stackoverflow.com/a/62677079/2256618
          private KeyValueStore<String, Event> keyValueStore;
      
          public SortProcessor(String stateStore, Long bufferIntervalInSeconds) {
              this.stateStore = stateStore;
              this.bufferIntervalInSeconds = bufferIntervalInSeconds;
          }
      
          @Override
          public void init(ProcessorContext processorContext) {
              super.init(processorContext);
              keyValueStore = (KeyValueStore) context().getStateStore(stateStore);
              context().schedule(Duration.ofSeconds(bufferIntervalInSeconds), PunctuationType.WALL_CLOCK_TIME, this::punctuate);
          }
      
          void punctuate(long timestamp) {
              LOG.info("Punctuator invoked...");
              try (KeyValueIterator<String, Event> iterator = keyValueStore.all()) {
                  while (iterator.hasNext()) {
                      KeyValue<String, Event> next = iterator.next();
                      if (next.value == null) {
                          continue;
                      }
                      LOG.info("Sending {}", next.key);
                      context().forward(null, next.value);
                      keyValueStore.delete(next.key);
                  }
              }
          }
      
          @Override
          public void process(String key, Event value) {
              Event event = Event.builder(value).payload(value.getPayload().toUpperCase()).build();
              keyValueStore.put(event.getEventType().name() + " " + event.getId(), event);
          }
      
          public static String getName() {
              return "sort-processor";
          }
      }
      

      可执行代码是here。我在这里使用了一个简单的内存状态存储。如果您预计会在短时间内发生大量事件,则可以使用其他答案中已经建议的持久状态存储。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2011-05-17
        • 1970-01-01
        • 2011-11-12
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多