如何在Spring Integration流中测试JPA适配器步骤?
问题
我有一个集成流(简化版本如下),包含多个带JPA适配器的步骤:第一步根据ID从数据库获取现有记录,第二步将更新后的实体保存到数据库。
@Autowired private EntityManager entityManager; @Bean public IntegrationFlow flow() { return IntegrationFlows .from( Kafka.messageDrivenChannelAdapter(consumerFactory(), "inTopic") .id("InboundKafkaAdapter")) .transform(transform()) .filter(filter()) .handle( Jpa.retrievingGateway(this.entityManager) .idExpression("headers['" + Headers.ID + "']") .entityClass(TestEntity.class), s -> s.advice(interceptForResult()).requiresReply(true)) .filter(secondFilter()) .transform(transformUpdatedEntity()) .handle(jpaAdapter(), ConsumerEndpointSpec::transactional) .handle( Kafka.outboundChannelAdapter(kafkaTemplate()) .topic("outTopic")) .get(); } private JpaUpdatingOutboundEndpointSpec jpaAdapter() { return Jpa.updatingGateway(this.entityManager) .entityClass(TestEntity.class) .flush(true) .persistMode(PersistMode.MERGE); }
测试时我希望Mock所有外部依赖,仅测试流本身。尝试用IntegrationFlowContext的代码如下:
@SpringBootTest @SpringIntegrationTest public class SampleFlowTest { @Mock private SampleFlow sampleFlow; @Test public void testSampleFlow() throws IOException { IntegrationFlow originalFlow = sampleFlow.flow(); IntegrationFlowContext.IntegrationFlowRegistration flowRegistration = integrationFlowContext.registration(originalFlow).register(); final Message<?> request = MessageBuilder.withPayload("somePayload") .setHeader(KafkaHeaders.TOPIC, "inTopic") .setHeader(KafkaHeaders.MESSAGE_KEY, "1") .setHeader(KafkaHeaders.RECEIVED_MESSAGE_KEY, "1") .build(); Message<?> response = flowRegistration.getMessagingTemplate().sendAndReceive(request); flowRegistration.destroy(); } @Configuration @EnableIntegration public static class Config { // 一些配置Bean } }
运行测试时,流执行到第一个JPA步骤就出问题,因为interceptForResult()逻辑无法正常工作。移除第一个JPA步骤后,第二个步骤又抛出事务管理器Bean缺失的异常。
我还尝试用substituteMessageHandlerFor模拟这两个处理器(已设置ID),但再次出现找不到指定ID的Bean的异常,我认为是因为原流类被Mock了。
恳请帮助我在测试中Mock JPA步骤!
更新 -- 测试实现:
@DirtiesContext @SpringIntegrationTest (noAutoStartup = {"InboundKafkaAdapter"}) @SpringBootTest( webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT, classes = Application.class) @ExtendWith({SpringExtension.class}) public class FlowTest { @Autowired private ApplicationContext applicationContext; @Autowired private MockIntegrationContext mockIntegrationContext; @Autowired private MyIntegrationFlow integrationFlow; @Autowired private QueueChannel testChannel; @Autowired @Qualifier("flow.channel#0") private MessageChannel flow; @Test public void test() { Message request = generateMessageForEvent(); MessageHandler mockMessageHandler = mockMessageHandler().handleNextAndReply(Function.identity()); this.mockIntegrationContext.substituteMessageHandlerFor( "JpaRetrievingGateway", mockMessageHandler); this.mockIntegrationContext.substituteMessageHandlerFor( "JpaUpdatingGateway", mockMessageHandler); this.mockIntegrationContext.substituteMessageHandlerFor( "OutboundKafkaAdapter", mockMessageHandler); flow.send(request); Message<String> reply = (Message<String>) testChannel.receive(0); Assert.assertNotNull("reply should not be null", reply); } @Configuration @EnableIntegration public static class Config { @Bean public QueueChannel testChannel(){ return new QueueChannel(); } } }
解决方案
步骤1:给JPA网关和Kafka出站适配器设置明确ID
原流中JPA检索网关、更新网关以及Kafka出站适配器都没有设置ID(仅入站Kafka适配器有ID),这是导致substituteMessageHandlerFor找不到目标的核心原因。修改原流代码,给这些处理器添加ID:
@Autowired private EntityManager entityManager; @Bean public IntegrationFlow flow() { return IntegrationFlows .from( Kafka.messageDrivenChannelAdapter(consumerFactory(), "inTopic") .id("InboundKafkaAdapter")) .transform(transform()) .filter(filter()) .handle( Jpa.retrievingGateway(this.entityManager) .id("JpaRetrievingGateway") // 添加明确ID .idExpression("headers['" + Headers.ID + "']") .entityClass(TestEntity.class), s -> s.advice(interceptForResult()).requiresReply(true)) .filter(secondFilter()) .transform(transformUpdatedEntity()) .handle(jpaAdapter(), ConsumerEndpointSpec::transactional) .handle( Kafka.outboundChannelAdapter(kafkaTemplate()) .id("OutboundKafkaAdapter") // 添加明确ID .topic("outTopic")) .get(); } private JpaUpdatingOutboundEndpointSpec jpaAdapter() { return Jpa.updatingGateway(this.entityManager) .id("JpaUpdatingGateway") // 添加明确ID .entityClass(TestEntity.class) .flush(true) .persistMode(PersistMode.MERGE); }
步骤2:配置测试用事务管理器
第二个JPA步骤标记了transactional,测试环境需要提供一个事务管理器Bean。可以添加一个测试用的Mock事务管理器:
在测试类的内部Config类中添加:
@Bean public PlatformTransactionManager transactionManager() { return new MockPlatformTransactionManager(); }
注:
MockPlatformTransactionManager是Spring提供的测试类,位于org.springframework.transaction.support包下。
步骤3:修正测试类逻辑
- 不要Mock原
MyIntegrationFlow类,直接通过@Autowired注入真实的流Bean,否则流中的处理器不会被Spring容器管理,导致找不到ID。 - 针对不同的Mock需求,为每个处理器创建对应的Mock逻辑,避免共用同一个
mockMessageHandler。
修改后的测试类示例:
@DirtiesContext @SpringIntegrationTest(noAutoStartup = {"InboundKafkaAdapter"}) @SpringBootTest(classes = Application.class) public class FlowTest { @Autowired private MockIntegrationContext mockIntegrationContext; @Autowired @Qualifier("flow.channel#0") private MessageChannel flowInputChannel; @Autowired private QueueChannel testChannel; @Test public void testFlow() { // 1. 构建测试请求消息 Message<String> request = MessageBuilder.withPayload("testPayload") .setHeader(Headers.ID, "1") .build(); // 2. Mock JPA检索网关:返回模拟的TestEntity TestEntity mockEntity = new TestEntity(); mockEntity.setId("1"); mockEntity.setName("MockedEntity"); MessageHandler jpaRetrieveMock = mockMessageHandler() .handleNextAndReply(msg -> mockEntity); // 3. Mock JPA更新网关:返回更新后的实体 MessageHandler jpaUpdateMock = mockMessageHandler() .handleNextAndReply(msg -> { TestEntity updated = (TestEntity) msg.getPayload(); updated.setName("UpdatedMockedEntity"); return updated; }); // 4. Mock Kafka出站适配器:将消息转发到测试通道 MessageHandler kafkaOutMock = mockMessageHandler() .handleNext(msg -> testChannel.send(msg)); // 5. 替换真实处理器为Mock mockIntegrationContext.substituteMessageHandlerFor("JpaRetrievingGateway", jpaRetrieveMock); mockIntegrationContext.substituteMessageHandlerFor("JpaUpdatingGateway", jpaUpdateMock); mockIntegrationContext.substituteMessageHandlerFor("OutboundKafkaAdapter", kafkaOutMock); // 6. 发送测试消息到流的输入通道 flowInputChannel.send(request); // 7. 验证测试通道收到消息 Message<?> reply = testChannel.receive(5000); Assert.assertNotNull("未收到回复消息", reply); TestEntity payload = (TestEntity) reply.getPayload(); Assert.assertEquals("UpdatedMockedEntity", payload.getName()); } @Configuration @EnableIntegration public static class TestConfig { @Bean public QueueChannel testChannel() { return new QueueChannel(); } @Bean public PlatformTransactionManager transactionManager() { return new MockPlatformTransactionManager(); } } }
关键说明
- 禁止Mock原Flow类:
@Mock private MyIntegrationFlow integrationFlow会导致Spring容器无法加载真实的流定义,所有处理器Bean都不会被注册,自然找不到指定ID。 - 给处理器设置明确ID:这是使用
substituteMessageHandlerFor的前提,Spring Integration会根据ID识别对应的处理器。 - 提供事务管理器:标记
transactional的端点必须有可用的事务管理器Bean,测试环境可以用Mock版本避免依赖真实数据库。
内容的提问来源于stack exchange,提问作者LisaIce1
相关产品推荐
相关产品推荐

