我正在尝试为Kafka Streams开发一个交互式查询应用程序。这是一个简单的基于count()的状态存储。但我看到的问题是,一旦我将应用程序扩展到多个实例,我就开始获得一些键的空值
KStream<String, String> inputStream = builder.stream(INPUT_TOPIC, Consumed.with(Serdes.String(), Serdes.String())); //key: foo, value:bar
inputStream.groupByKey(Grouped.with(Serdes.String(), Serdes.String()))
.count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as(STATE_STORE_NAME)
.withKeySerde(Serdes.String())
.withValueSerde(Serdes.Long()));就测试基于DSL的管道而言,这差不多就是它了。我有一个用于交互式查询的REST端点
KafkaStreams streams = ...;
ReadOnlyKeyValueStore<String, Long> averageStore = streams.store(storeName, QueryableStoreTypes.<String, Long>keyValueStore());
Long count = averageStore.get(word);count为空-此行为仅适用于某些键。并且这与密钥是否存在于本地无关
发布于 2019-12-28 07:29:32
当您扩展Kafka Streams应用程序时,只有全局表在所有实例上都是可见的。对于常规KTable,整个数据集只有一部分可用。
您需要查找关键的元数据,并将REST调用重定向到相应的实例,如here所述。
https://stackoverflow.com/questions/57974923
复制相似问题