【发布时间】: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