看来我错了,内存状态存储缓存工作正常。我将简要说明我是如何测试它的,也许有人会发现它很有用。我制作了一个非常基本的Kafka Streams应用程序,它只读取抽象为KTable的主题。
public class Main {
public static void main(String[] args) {
StreamsBuilder builder = new StreamsBuilder();
Logger logger = LoggerFactory.getLogger(Main.class);
builder.table("inputTopic", Materialized.as(Stores.inMemoryKeyValueStore("myStore")).withCachingEnabled())
.toStream()
.foreach((k, v) -> logger.info("Result: {} - {}", k, v));
new KafkaStreams(builder.build(), getProperties()).start();
}
private static Properties getProperties() {
Properties properties = new Properties();
properties.put(APPLICATION_ID_CONFIG, "testApp");
properties.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
properties.put(COMMIT_INTERVAL_MS_CONFIG, 10000);
properties.put(CACHE_MAX_BYTES_BUFFERING_CONFIG, 10 * 1024 * 1024L);
properties.put(DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
properties.put(DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
return properties;
}
}
然后我运行了卡夫卡的控制台制作人:
/kafka-console-producer.sh --broker-list localhost:9092 --topic inputTopic --property "parse.key=true" --property "key.separator=:"
并发送了几条消息:a:a、a:b、a:c。应用程序中只能看到最后一条消息,因此缓存工作正常。
2018-03-06 21:21:57信息主:26-结果:a-c
我还稍微更改了流以检查
aggregate
方法
builder.stream("inputTopic")
.groupByKey()
.aggregate(() -> "", (k, v, a) -> a + v, Materialized.as(Stores.inMemoryKeyValueStore("aggregate")))
.toStream()
.foreach((k, v) -> logger.info("Result: {} - {}", k, v));
我用同一个键连续快速发送了几条消息,只收到了一个结果,因此数据没有立即发送到下游-完全按照预期。