【发布时间】:2020-08-28 02:44:15
【问题描述】:
我想知道是否指定了流拓扑处理消息的顺序。
例子:
// read input messages
KStream<String, String> inputMessages = builder.stream("demo_input_topic_1");
inputMessages = inputMessages.peek((k, v) -> System.out.println("TECHN. NEW MESSAGE: key: " + k + ", value: " + v));
// check if message was already processed
KTable<String, Long> alreadyProcessedMessages = inputMessages.groupByKey().count();
KStream<String, String> newMessages =
inputMessages.leftJoin(alreadyProcessedMessages, (streamValue, tableValue) -> getMessageValueOrNullIfKnownMessage(streamValue, tableValue));
KStream<String, String> filteredNewMessages =
newMessages.filter((key, val) -> val != null).peek((k, v) -> System.out.println("FUNC. NEW MESSAGE: key: " + k + ", value: " + v));
// process the message
filteredNewMessages.map((key, value) -> KeyValue.pair(key, "processed message: " + value))
.peek((k, v) -> System.out.println("PROCESSED MESSAGE: key: " + k + ", value: " + v)).to("demo_output_topic_1");
与getMessageValueOrNullIfKnownMessage(...):
private static String getMessageValueOrNullIfKnownMessage(String newMessageValue, Long messageCounter) {
if (messageCounter > 1) {
return null;
}
return newMessageValue;
}
所以示例中只有一个输入和一个输出主题。
输入主题在alreadyProcessedMessages 中被计数(因此创建了一个本地状态)。此外,输入主题与计数表alreadyProcessedMessages 连接,连接的结果是流newMessages(如果消息计数> 1,则此流中的消息值为null,否则为消息的原始值)。
然后,过滤newMessages 的消息(过滤掉null 的值)并将结果写入输出主题。
那么这个最小流的作用是:它将输入主题中的所有消息写入具有新键(以前未处理过的键)的输出主题。
在流工作的测试中。但我认为这不能保证。它只是有效的,因为消息在加入之前首先由计数节点处理。
但是该订单有保证吗?
据我在所有文档中看到的,无法保证此处理顺序。因此,如果有新消息到达,也可能发生这种情况:
- 消息由“加入节点”处理。
- 消息由“计数节点”处理。
这当然会产生不同的结果(因此在这种情况下,如果具有相同键的消息第二次出现,它仍然会与原始值连接,因为它还没有被计算在内)。
那么在某处指定了处理顺序吗?
我知道在新版本的 Kafka 中,KStream-KTable 连接是根据输入分区中消息的时间戳完成的。但这在这里没有帮助,因为拓扑使用相同的输入分区(因为它的消息相同)。
谢谢
【问题讨论】:
标签: apache-kafka stream apache-kafka-streams