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

如何使用MockProducer编写Spring Kafka生产者事务的单元测试

基于MockProducer测试Spring Kafka生产者事务

可行性结论

完全可以使用MockProducer编写单元测试,覆盖生产者事务的提交成功和失败回滚场景。以下是具体实现方案,包括KafkaTemplate的初始化方法以及完整测试用例。


核心步骤:初始化事务型KafkaTemplate

要让KafkaTemplate配合MockProducer支持事务,需要自定义ProducerFactory,让它返回我们的MockProducer实例,并开启事务配置:

// 1. 初始化支持事务的MockProducer
// 第一个参数为true,表示启用事务模式
MockProducer<String, Data> mockProducer = new MockProducer<>(
        true, 
        new StringSerializer(), 
        new JsonSerializer<>()
);

// 2. 构建自定义ProducerFactory,返回MockProducer
ProducerFactory<String, Data> producerFactory = new DefaultKafkaProducerFactory<>(
        // 配置事务ID前缀,与生产环境一致
        Map.of(ProducerConfig.TRANSACTION_ID_PREFIX_CONFIG, "tx-"),
        new StringSerializer(),
        new JsonSerializer<>()
) {
    @Override
    public Producer<String, Data> createProducer() {
        return mockProducer;
    }
};

// 3. 初始化KafkaTemplate并开启事务支持
KafkaTemplate<String, Data> kafkaTemplate = new KafkaTemplate<>(producerFactory);
kafkaTemplate.setTransactional(true);

测试用例1:事务提交场景

验证事务提交后,消息被正确写入MockProducer的历史记录:

@Test
void testTransactionalCommit() {
    // 初始化MockProducer和KafkaTemplate
    MockProducer<String, Data> mockProducer = new MockProducer<>(true, new StringSerializer(), new JsonSerializer<>());
    ProducerFactory<String, Data> producerFactory = new DefaultKafkaProducerFactory<>(
            Map.of(ProducerConfig.TRANSACTION_ID_PREFIX_CONFIG, "tx-"),
            new StringSerializer(),
            new JsonSerializer<>()
    ) {
        @Override
        public Producer<String, Data> createProducer() {
            return mockProducer;
        }
    };
    KafkaTemplate<String, Data> kafkaTemplate = new KafkaTemplate<>(producerFactory);
    kafkaTemplate.setTransactional(true);

    // 模拟事务流程
    kafkaTemplate.initTransaction();
    kafkaTemplate.beginTransaction();

    // 发送测试消息
    ProducerRecord<String, Data> testRecord = new ProducerRecord<>("test-topic", new Data("test-content"));
    kafkaTemplate.send(testRecord);

    // 提交前:事务未完成,MockProducer无历史记录
    assertTrue(mockProducer.history().isEmpty());

    // 提交事务
    kafkaTemplate.commitTransaction();

    // 提交后:验证消息已写入历史记录
    assertEquals(1, mockProducer.history().size());
    assertEquals(testRecord.value(), mockProducer.history().get(0).value());
}

测试用例2:事务回滚场景

验证事务回滚后,消息不会被写入MockProducer的历史记录:

@Test
void testTransactionalRollback() {
    // 初始化MockProducer和KafkaTemplate
    MockProducer<String, Data> mockProducer = new MockProducer<>(true, new StringSerializer(), new JsonSerializer<>());
    ProducerFactory<String, Data> producerFactory = new DefaultKafkaProducerFactory<>(
            Map.of(ProducerConfig.TRANSACTION_ID_PREFIX_CONFIG, "tx-"),
            new StringSerializer(),
            new JsonSerializer<>()
    ) {
        @Override
        public Producer<String, Data> createProducer() {
            return mockProducer;
        }
    };
    KafkaTemplate<String, Data> kafkaTemplate = new KafkaTemplate<>(producerFactory);
    kafkaTemplate.setTransactional(true);

    // 模拟事务流程
    kafkaTemplate.initTransaction();
    kafkaTemplate.beginTransaction();

    // 发送测试消息
    ProducerRecord<String, Data> testRecord = new ProducerRecord<>("test-topic", new Data("test-content"));
    kafkaTemplate.send(testRecord);

    // 回滚事务
    kafkaTemplate.rollbackTransaction();

    // 回滚后:验证MockProducer无历史记录
    assertTrue(mockProducer.history().isEmpty());
}

测试自定义业务方法(sendData)

如果要测试你自己的@Transactional注解修饰的sendData方法,可以结合Spring的TransactionTemplate来控制事务:

@Test
void testSendDataOnCommit() {
    // 初始化MockProducer和KafkaTemplate
    MockProducer<String, Data> mockProducer = new MockProducer<>(true, new StringSerializer(), new JsonSerializer<>());
    ProducerFactory<String, Data> producerFactory = new DefaultKafkaProducerFactory<>(
            Map.of(ProducerConfig.TRANSACTION_ID_PREFIX_CONFIG, "tx-"),
            new StringSerializer(),
            new JsonSerializer<>()
    ) {
        @Override
        public Producer<String, Data> createProducer() {
            return mockProducer;
        }
    };
    KafkaTemplate<String, Data> kafkaTemplate = new KafkaTemplate<>(producerFactory);
    kafkaTemplate.setTransactional(true);

    // 实例化你的业务类,注入KafkaTemplate
    YourBusinessService service = new YourBusinessService();
    service.setKafkaTemplate(kafkaTemplate);

    // 使用TransactionTemplate执行事务
    TransactionTemplate transactionTemplate = new TransactionTemplate(transactionManager);
    transactionTemplate.execute(status -> {
        boolean result = service.sendData("test-topic", new Data("test-content"));
        assertTrue(result);
        return null;
    });

    // 验证事务提交后消息已写入
    assertEquals(1, mockProducer.history().size());
}

@Test
void testSendDataOnRollback() {
    // 初始化MockProducer和KafkaTemplate
    MockProducer<String, Data> mockProducer = new MockProducer<>(true, new StringSerializer(), new JsonSerializer<>());
    ProducerFactory<String, Data> producerFactory = new DefaultKafkaProducerFactory<>(
            Map.of(ProducerConfig.TRANSACTION_ID_PREFIX_CONFIG, "tx-"),
            new StringSerializer(),
            new JsonSerializer<>()
    ) {
        @Override
        public Producer<String, Data> createProducer() {
            return mockProducer;
        }
    };
    KafkaTemplate<String, Data> kafkaTemplate = new KafkaTemplate<>(producerFactory);
    kafkaTemplate.setTransactional(true);

    // 实例化你的业务类,注入KafkaTemplate
    YourBusinessService service = new YourBusinessService();
    service.setKafkaTemplate(kafkaTemplate);

    // 模拟事务回滚(抛出异常触发)
    TransactionTemplate transactionTemplate = new TransactionTemplate(transactionManager);
    assertThrows(RuntimeException.class, () -> {
        transactionTemplate.execute(status -> {
            service.sendData("test-topic", new Data("test-content"));
            // 手动触发回滚
            status.setRollbackOnly();
            throw new RuntimeException("Force transaction rollback");
        });
    });

    // 验证事务回滚后无消息写入
    assertTrue(mockProducer.history().isEmpty());
}

注意事项

  • MockProducer构造器的第一个参数必须设为true,否则不支持事务
  • 自定义ProducerFactory时,必须确保createProducer返回同一个MockProducer实例
  • KafkaTemplate必须调用setTransactional(true)才能启用事务支持
  • 测试回滚场景时,需确保事务被正确触发回滚(如抛出异常、手动设置setRollbackOnly)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 23:27:50