【问题标题】:Unit testing a kafka topology that's using kstream joins对使用 kstream 连接的 kafka 拓扑进行单元测试
【发布时间】:2019-07-30 12:20:06
【问题描述】:

我有一个执行两个 kstream 连接的拓扑,我面临的问题是尝试使用 TopologyTestDriver 进行单元测试时发送几个带有 pipeInput 然后 readOutput 的 ConsumerRecord。它似乎不起作用。

我认为这可能是因为连接使用的是我们在测试中没有使用的实际 kafka 中的内部 Rocksdb。

所以我一直在寻找解决方案,但找不到任何解决方案。

注意:这种测试方法在删除 kstream-kstream 连接时效果很好。

【问题讨论】:

    标签: java unit-testing join apache-kafka apache-kafka-streams


    【解决方案1】:

    我有一个进行两个 kstream 连接的拓扑,我面临的问题是尝试使用 TopologyTestDriver 进行单元测试时发送几个带有 pipeInput 然后 readOutput 的 ConsumerRecord。它似乎不起作用。

    按照设计,但不幸的是,在您的情况下,TopologyTestDriver 并不是 Kafka Streams 引擎在运行时如何工作的 100% 准确模型。值得注意的是,新传入事件的处理顺序存在一些差异。

    这确实会在尝试测试某些连接时引起问题,例如,某些连接,因为这些操作依赖于某个处理顺序(例如,在流表连接中,表之前应该已经有一个键 'alice' 的条目'alice' 的流端事件到达,否则流端 'alice' 的连接输出将不包含任何表端数据)。

    所以我一直在寻找解决方案,但找不到任何解决方案。

    我的建议是使用启动嵌入式 Kafka 集群的测试,然后使用“真正的”Kafka Streams 引擎(即,不是 TopologyTestDriver)针对该集群运行测试。实际上,这意味着您正在将测试从单元测试更改为集成/系统测试:您的测试将启动一个成熟的 Kafka Streams 拓扑,该拓扑与与您的测试在同一台机器上运行的嵌入式 Kafka 集群进行通信。

    请参阅 Apache Kafka 项目中的 Kafka Streams 集成测试,其中 EmbeddedKafkaClusterIntegrationTestUtils 是工具的核心部分。连接的具体测试示例是StreamTableJoinIntegrationTest(有一些与连接相关的集成测试)及其父级AbstractJoinIntegrationTest。 (值得一提的是,https://github.com/confluentinc/kafka-streams-examples#examples-integration-tests 有进一步的集成测试示例,其中包括使用 Apache Avro 作为数据格式时还涵盖 Confluent Schema Registry 的测试等)

    但是,除非我弄错了,否则集成测试及其工具包含在 Kafka Streams 的 test utilities 工件(即 org.apache.kafka:kafka-streams-test-utils)中。所以你必须做一些复制粘贴到你自己的代码库中。

    【讨论】:

      【解决方案2】:

      您看过 Kafka Streams 单元测试 [1] 吗?这是关于输入数据并使用模拟处理器检查最终结果。

      例如对于以下流连接:

              stream1 = builder.stream(topic1, consumed);
              stream2 = builder.stream(topic2, consumed);
              joined = stream1.outerJoin(
                  stream2,
                  MockValueJoiner.TOSTRING_JOINER,
                  JoinWindows.of(ofMillis(100)),
                  StreamJoined.with(Serdes.Integer(), Serdes.String(), Serdes.String()));
              joined.process(supplier);
      

      然后您可以开始将输入项通过管道传输到第一个或第二个主题中,并检查每个连续的输入管道,处理器可以检查什么:

      // push two items to the primary stream; the other window is empty
                  // w1 = {}
                  // w2 = {}
                  // --> w1 = { 0:A0, 1:A1 }
                  //     w2 = {}
                  for (int i = 0; i < 2; i++) {
                      inputTopic1.pipeInput(expectedKeys[i], "A" + expectedKeys[i]);
                  }
                  processor.checkAndClearProcessResult(EMPTY);
      
                  // push two items to the other stream; this should produce two items
                  // w1 = { 0:A0, 1:A1 }
                  // w2 = {}
                  // --> w1 = { 0:A0, 1:A1 }
                  //     w2 = { 0:a0, 1:a1 }
                  for (int i = 0; i < 2; i++) {
                      inputTopic2.pipeInput(expectedKeys[i], "a" + expectedKeys[i]);
                  }
                  processor.checkAndClearProcessResult(new KeyValueTimestamp<>(0, "A0+a0", 0),
                      new KeyValueTimestamp<>(1, "A1+a1", 0));
      

      我希望这会有所帮助。

      参考资料: [1]https://github.com/apache/kafka/blob/trunk/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KStreamKStreamJoinTest.java#L279

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2019-04-15
        • 2017-06-09
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-08-30
        相关资源
        最近更新 更多