【发布时间】:2019-12-28 07:29:32
【问题描述】:
我正在尝试为 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 为 null - 此行为仅适用于某些键。这与本地是否存在密钥无关
【问题讨论】:
-
这个运气好吗?我有一个 kafka 流应用程序,当我从状态存储中获取拓扑中另一个处理器中的键的值时,它为我之前放置的相同键提供了空值。
标签: apache-kafka apache-kafka-streams