You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka Streams中Join操作单元测试失败问题求助

Kafka Streams Join操作单元测试失败排查

问题背景

业务场景为将Employee流与Insurance表进行关联,输出包含保险公司名称的EnrichedEmployee,关联依据是Employee对象中的insuranceIdx字段。编写的EnrichEmployeesTopoplogy在生产Kafka环境运行正常,但使用kafka-streams-test-utils编写的单元测试断言始终失败,怀疑是测试时序问题。

拓扑实现代码

public class EnrichEmployeesTopoplogy {


    public void createTopology(final StreamsBuilder builder) {

        final Serde<EnrichedEmployee> enrichedEmployeeSerde = SerdeFactory.serdeForEnrichedEmployee();
        final Serde<Insurance> insuranceSerde = SerdeFactory.serdeForInsurance();
        final Serde<Employee> employeeSerde = SerdeFactory.serdeForEmployee();

        // 读取输入的Employee事件流
        final KStream<String, Employee> employeeStream =
                builder.stream("employees", Consumed.with(Serdes.String(), employeeSerde));

        // 重新设置流的key,使其与Insurance表的key一致
        final KStream<String, Employee> rekeyedEmployeeStream
                = employeeStream.selectKey((k, v) -> String.valueOf(v.getInsuranceIdx()));


        // 读取Insurance事件表
        final KTable<String, Insurance> insuranceTable = builder.table("insurances",
                Consumed.with(Serdes.String(), insuranceSerde));

        // 定义值关联器
        ValueJoiner<Employee, Insurance, EnrichedEmployee> employeeJoiner =
                (employee, insurance) ->
                        EnrichedEmployee.builder()
                                .idx(employee.getIdx())
                                .email(employee.getEmail())
                                .insuranceName(insurance.getName())
                                .build();

        // 配置关联参数
        final Joined<String, Employee, Insurance> joined = Joined.with(Serdes.String(), employeeSerde, insuranceSerde);


        rekeyedEmployeeStream.join(insuranceTable, employeeJoiner, joined)
                .to("enriched-employees", Produced.with(Serdes.String(), enrichedEmployeeSerde));


    }
}

单元测试代码

class EnrichedEmployeeTopologyTest {


    private TopologyTestDriver testDriver;
    private TestInputTopic<String, Employee> inputTopicEmployees;
    private TestInputTopic<String, Insurance> inputTopicInsurances;
    private TestOutputTopic<String, EnrichedEmployee> outputTopicEnrichedEmployees;

    @BeforeEach
    void setUp() {

        final Serde<EnrichedEmployee> enrichedEmployeeSerde = SerdeFactory.serdeForEnrichedEmployee();
        final Serde<Insurance> insuranceSerde = SerdeFactory.serdeForInsurance();
        final Serde<Employee> employeeSerde = SerdeFactory.serdeForEmployee();

        final Properties streamsConfiguration = new Properties();
        streamsConfiguration.put(StreamsConfig.APPLICATION_ID_CONFIG, "test01");
        streamsConfiguration.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "dummy:9092");

        final EnrichEmployeesTopoplogy enrichEmployeesTopoplogy = new EnrichEmployeesTopoplogy();
        final StreamsBuilder builder = new StreamsBuilder();

        enrichEmployeesTopoplogy.createTopology(builder);

        final Topology topology = builder.build();

        testDriver = new TopologyTestDriver(topology, streamsConfiguration);
        inputTopicEmployees = testDriver.createInputTopic("employees", Serdes.String().serializer(), employeeSerde.serializer());
        inputTopicInsurances = testDriver.createInputTopic("insurances", Serdes.String().serializer(), insuranceSerde.serializer());
        outputTopicEnrichedEmployees = testDriver.createOutputTopic("enriched-employees", Serdes.String().deserializer(), enrichedEmployeeSerde.deserializer());

    }

    @AfterEach
    void tearDown() {
        testDriver.close();
    }

    @Test
    void createTopology() {
        final String employeeIdx = UUID.randomUUID().toString();
        final String insuranceIdx = UUID.randomUUID().toString();

        final Employee employee = Employee.builder().idx(employeeIdx).email("foo1@bar.de").insuranceIdx(insuranceIdx).build();
        final Insurance insurance = Insurance.builder().idx(insuranceIdx).name("insurance 01").build();


        inputTopicEmployees.pipeInput(employee.getIdx(), employee);
        inputTopicInsurances.pipeInput(insurance.getIdx(), insurance);

        Assertions.assertEquals(true, !outputTopicEnrichedEmployees.isEmpty());
    }
}

问题原因与解决方法

核心原因

KTable在测试环境中需要先完成数据加载初始化,才能与后续流入的KStream事件完成关联。当前测试先发送Employee事件,此时KTable中尚无对应Insurance数据,事件会直接被跳过,无法生成关联结果。

解决步骤

  1. 调整数据发送顺序:先发送Insurance数据到KTable,让测试驱动完成表数据的加载初始化,再发送Employee事件。
  2. 优化断言逻辑:替换空判断为具体字段断言,提升测试准确性。

修改后的测试方法代码:

@Test
void createTopology() {
    final String employeeIdx = UUID.randomUUID().toString();
    final String insuranceIdx = UUID.randomUUID().toString();

    final Employee employee = Employee.builder().idx(employeeIdx).email("foo1@bar.de").insuranceIdx(insuranceIdx).build();
    final Insurance insurance = Insurance.builder().idx(insuranceIdx).name("insurance 01").build();

    // 先发送Insurance数据,确保KTable完成初始化
    inputTopicInsurances.pipeInput(insurance.getIdx(), insurance);
    // 再发送Employee事件触发关联
    inputTopicEmployees.pipeInput(employee.getIdx(), employee);

    // 断言输出不为空并验证字段正确性
    Assertions.assertFalse(outputTopicEnrichedEmployees.isEmpty());
    EnrichedEmployee result = outputTopicEnrichedEmployees.readValue();
    Assertions.assertEquals("insurance 01", result.getInsuranceName());
    Assertions.assertEquals(employee.getEmail(), result.getEmail());
}

额外优化建议

  • 测试中尽量使用具体字段断言,而非仅判断输出是否为空,能快速定位数据匹配问题。
  • 若涉及窗口或时间相关逻辑,需调用testDriver.advanceWallClockTime()推进测试时间触发处理逻辑。

内容的提问来源于stack exchange,提问作者Thomas Mueller

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.29 02:17:58