如何使用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
相关产品推荐
相关产品推荐

