【问题标题】:Is the processing order of a Kafka Streams topology specified?是否指定了 Kafka Streams 拓扑的处理顺序?
【发布时间】: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


    【解决方案1】:

    没有保证。即使在当前实现中,使用了List 的子节点:https://github.com/apache/kafka/blob/trunk/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java#L203-L206——但是,不能保证子节点按照它们在 DSL 中指定的顺序附加到此列表中(因为是中间的翻译层,可以以不同的顺序添加节点)。此外,实施可能会随时更改。

    我能想到的唯一解决方法(相当昂贵)可能是在 repartiton 主题中发送流端数据:

    KStream<String, String> newMessages =
       inputMessages.through(...) // note: as of 2.6.0 release, you could use `repartition()` instead of `through()`
                    .leftJoin(alreadyProcessedMessages, ...);
    

    这样,KTable 将在执行连接之前更新,因为需要先读回记录。但是,由于在回读记录时您没有任何保证,因此在连接完成之前可能会对表进行多次更新,这可能会使您处于与以前类似的情况。 (此外,通过其他主题重新路由数据有点昂贵。)

    使用处理器 API,您将拥有移动控制权,您可以调用 context.forward(..., To.child(...))。但是,对于这种情况,您还需要手动实现聚合和连接:

    KStream routing = inputMessages.transform(...);
    routing.groupByKey(...);
    routing.leftJoin(...);
    

    对于这种情况,您会在 transform() 之后获得要避免的重新分区主题:

    KStream routing = inputMessages.transform(...);
    routing.transform(...); // implement the aggregation
    routing.transform(...); // implement the join
    

    连续的transform() 不会触发自动重新分区。

    【讨论】:

    • 谢谢!如果 Confluent 工程师回答,那我别无选择,只能接受 ;-)
    【解决方案2】:

    这只是缩小未决问题的部分答案:

    在 (Confluent's Stream Architecture overview) 中声明使用“深度优先处理策略”来遍历拓扑。没有提到在多个路径上可以通过相同输入到达的节点同步。 (但是,在1 的细节级别上,基于此将其排除在外是很牵强的。)

    关于DFS遍历取分支的顺序,我没有找到明确的说法。然而,在这个Confluent documenation on namings within the topology 中,一些示例显示了“拓扑中的运算符顺序”。现在可以假设这个顺序。这似乎是由源代码中 DSL 运算符的顺序给出的,也是执行顺序。这将提供您所要求的保证。但是,我无法证实任何其他来源的假设。

    剩下的两个问题可以通过在 PAPI 实现中找到相关的源代码来回答。

    1. 真的只是没有同步点的普通 DFS 遍历吗?
    2. DFS 中的分支顺序真的是2 中定义的运算符顺序吗?如果不是,那是什么?

    【讨论】:

      猜你喜欢
      • 2017-06-09
      • 1970-01-01
      • 2019-10-23
      • 1970-01-01
      • 1970-01-01
      • 2020-05-04
      • 2018-06-25
      • 1970-01-01
      相关资源
      最近更新 更多