【问题标题】:TopologyTestDriver sending incorrect message on KTable aggregationsTopologyTestDriver 在 KTable 聚合上发送错误消息
【发布时间】:2019-01-25 20:08:45
【问题描述】:

我有一个在 KTable 上聚合的拓扑。 这是我创建的一种通用方法,用于在我拥有的不同主题上构建此拓扑。

public static <A, B, C> KTable<C, Set<B>> groupTable(KTable<A, B> table, Function<B, C> getKeyFunction,
        Serde<C> keySerde, Serde<B> valueSerde, Serde<Set<B>> aggregatedSerde) {
    return table
            .groupBy((key, value) -> KeyValue.pair(getKeyFunction.apply(value), value),
                    Serialized.with(keySerde, valueSerde))
            .aggregate(() -> new HashSet<>(), (key, newValue, agg) -> {
                agg.remove(newValue);
                agg.add(newValue);
                return agg;
            }, (key, oldValue, agg) -> {
                agg.remove(oldValue);
                return agg;
            }, Materialized.with(keySerde, aggregatedSerde));
}

这在使用 Kafka 时效果很好,但在通过 `TopologyTestDriver` 进行测试时却不行。

在这两种情况下,当我获得更新时,首先调用subtractor,然后调用adder。问题是当使用TopologyTestDriver 时,会发送两条消息进行更新:一条在subtractor 调用之后,另一条在adder 调用之后。更不用说subrtractor之后和adder之前发送的消息处于错误的阶段。

其他人可以确认这是一个错误吗?我已经针对 Kafka 版本 2.0.1 和 2.1.0 对此进行了测试。

编辑:
我在 github 中创建了一个测试用例来说明这个问题:https://github.com/mulho/topology-testcase

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    有两条输出记录(一条“减号”记录和一条“加号”记录)是预期行为。理解它的工作原理有点棘手,所以让我试着解释一下。

    假设你有以下输入表:

     key |  value
    -----+---------
      A  |  <10,2>
      B  |  <10,3>
      C  |  <11,4>
    

    KTable#groupBy() 上,您将值的第一部分提取为新键(即1011),然后对第二部分求和(即234)在聚合中。因为AB 记录都将10 作为新键,所以您可以对2+3 求和,也可以对4 求和作为新键11。结果表将是:

     key |  value
    -----+---------
      10 |  5
      11 |  4
    

    现在假设一条更新记录&lt;B,&lt;11,5&gt;&gt;将原来的输入KTable改为:

     key |  value
    -----+---------
      A  |  <10,2>
      B  |  <11,5>
      C  |  <11,4>
    

    因此,新结果表应将5+411210 相加:

     key |  value
    -----+---------
      10 |  2
      11 |  9
    

    如果您将第一个结果表与第二个结果表进行比较,您可能会注意到 两行 行都得到了更新。从10|5 中减去旧的B|&lt;10,3&gt; 记录得到10|2,并将新的B|&lt;11,5&gt; 记录添加到11|4 得到11|9

    这正是您看到的两条输出记录。第一条输出记录(在执行减法之后)更新第一行(它减去不再属于聚合结果的旧值),而第二条记录将新值添加到聚合结果中。在我们的例子中,减记录是&lt;10,&lt;null,&lt;10,3&gt;&gt;&gt;,加记录是&lt;11,&lt;&lt;11,5&gt;,null&gt;&gt;(这些记录的格式是&lt;key, &lt;plus,minus&gt;&gt;(注意减记录只设置minus部分,而加记录只设置plus 部分)。

    最后说明:加号和减号记录不能放在一起,因为加号和减号记录的键可以不同(在我们的示例中为1110),因此可能进入不同的分区.这意味着加号和减号操作可能由不同的机器执行,因此不可能只发出一条同时包含加号和减号的记录。

    【讨论】:

    • 非常感谢您的解释。我只是想知道为什么我在运行实际应用程序时看不到那些,而不是使用 TopologyTestDriver。在我的情况下,由于我没有接触新键(即作者姓名永远不会改变),实际实现是否有优化以知道新键与以前的记录一致,因此此更新不需要作为 2 条单独的消息发送?
    • 应该没有任何区别,两种情况都应该有两条消息(我不记得这种情况的任何内部优化)。你说“我没有看到那些”——如何“观察”输出流?
    • 我的测试用例非常简单,我有一个输入主题,使用 KTable 读取,该 KTable 将结果分组和聚合并将结果发送到另一个主题。我只是使用命令行来打印来自输出主题的消息。另外,为了确保它不是由于“提交间隔”,我在聚合加法器和减法器上放置了一个断点以在两者上添加延迟,但在减法器之后和加法器之前仍然没有发布消息。跨度>
    • 会不会和缓存有关? TopologyTestDriver flushes 在每个输入记录之后。您是否尝试禁用缓存?
    猜你喜欢
    • 2020-02-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-02-21
    相关资源
    最近更新 更多