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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 01:48:28