【发布时间】:2018-06-18 21:09:55
【问题描述】:
我正在通过 Kafka 流接收消息。它们由用户 ID 键入。他们进来时会得到一个序列号和时间戳。消息在 15 分钟后“过期”。用户可以根据给定的时间(最多 15 分钟)或顺序请求新消息。
我最初拥有的是这样的:
` StreamsBuilder streamsBuilder = new StreamsBuilder();
KStream<String, Message> inboundStream = streamsBuilder.stream("incoming.topic");
messageSupplier = Stores.persistentKeyValueStore("user.messages");
KTable<String, MessageCache> messageTable = inboundStream
.filter(this::userExists)
.peek(this::recordInboundMessage)
.map(this::markMessage) // add sequence/timestamp
.groupByKey()
.aggregate(this::createMessageCache,
this::addMessageToMessageCache,
Materialized.as(messageSupplier));
// ---> Some other setup stuff, then start the streams
`
MessageCache 保存消息列表(当我们将消息添加到缓存时删除过期消息)。当我收到消息请求时,我会浏览列表并过滤掉相应的消息。
我在想我可以使用其中一种窗口策略,但找不到实际持久化消息列表的示例。
这是最好的方法吗?还是我错过了一些能让这更容易/更好的东西?
【问题讨论】:
标签: java stream apache-kafka apache-kafka-streams