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

使用Spring Integration+MySQL实现Outbox模式时遭遇锁等待超时

问题描述

我尝试用Spring Integration实现Outbox模式,配置了相关Bean、ApplicationRunner在启动时模拟发送2封邮件(邮件处理逻辑会主动抛出异常),期望Spring Integration每隔1秒重试发送,且数据库保留消息记录,重启程序后能继续重试。但运行时90%概率出现锁等待超时异常,临时给轮询器添加初始延迟后问题暂时消失,但生产环境仍可能复现,求根本解决方案。


相关配置代码

配置类代码

@Configuration
public class SpringIntegrationTestApplicationConfiguration {
    private static final Logger LOGGER = LoggerFactory.getLogger(SpringIntegrationTestApplicationConfiguration.class);

    public static final String CONCURRENT_METADATA_STORE_PREFIX = "_spring_integration_";

    @MessagingGateway
    public interface EmailGateway {
        @Gateway(requestChannel = "mailbox")
        void sendEmail(String mailBody,
                       @Header String target);
    }

    @Bean
    public JdbcChannelMessageStore messageStore(DataSource dataSource) {
        JdbcChannelMessageStore jdbcChannelMessageStore = new JdbcChannelMessageStore(dataSource);
        jdbcChannelMessageStore.setTablePrefix(CONCURRENT_METADATA_STORE_PREFIX);
        jdbcChannelMessageStore.setChannelMessageStoreQueryProvider(
                new MySqlChannelMessageStoreQueryProvider());
        return jdbcChannelMessageStore;
    }

    @Bean
    ConcurrentMetadataStore concurrentMetadataStore(DataSource dataSource) {
        JdbcMetadataStore jdbcMetadataStore = new JdbcMetadataStore(dataSource);
        jdbcMetadataStore.setTablePrefix(CONCURRENT_METADATA_STORE_PREFIX);
        return jdbcMetadataStore;
    }


    @Bean
    MessageHandler sendEmailMessageHandler() {
        return new MessageHandler() {
            @Override
            public void handleMessage(Message<?> message) throws MessagingException {
                String target = (String) message.getHeaders().get("target");
                LOGGER.info("not sending email with body: {} {}", message, target);
                throw new RuntimeException("");
            }
        };
    }

    @Bean
    QueueChannel mailboxChannel(JdbcChannelMessageStore jdbcChannelMessageStore) {
        return MessageChannels.queue(jdbcChannelMessageStore, "mailbox").getObject();
    }

    @Bean
    public IntegrationFlow buildFlow(ChannelMessageStore channelMessageStore,
                                     MessageHandler sendEmailMessageHandler) {
        return IntegrationFlow.from("mailbox")
                              .routeToRecipients(routes -> {
                                  routes
                                          .transactional()
                                          .recipientFlow(flow -> flow
                                                  .channel(channels -> channels.queue(channelMessageStore, "outbox"))
                                                  .handle(sendEmailMessageHandler, e -> e.poller(poller -> poller.fixedDelay(1000).transactional()))
                                          );
                              }).get();
    }

}

ApplicationRunner代码

@Component
public class Runner implements ApplicationRunner {

    private static final Logger LOGGER = LoggerFactory.getLogger(Runner.class);

    private final SpringIntegrationTestApplicationConfiguration.EmailGateway emailGateway;

    public Runner(SpringIntegrationTestApplicationConfiguration.EmailGateway emailGateway) {
        this.emailGateway = emailGateway;
    }

    @Override
    public void run(ApplicationArguments args) throws Exception {
        LOGGER.info("Sending 1");
        emailGateway.sendEmail("This is my body", "target");
        LOGGER.info("Sending 2");
        emailGateway.sendEmail("This is my body2", "target2");
    }
}

Docker Compose配置

version: "3.9"
services:
  db:
    image: mysql:8-oracle
    environment:
      MYSQL_ROOT_PASSWORD: 'root'
      MYSQL_ALLOW_EMPTY_PASSWORD: 1
      MYSQL_ROOT_HOST: "%"
      MYSQL_DATABASE: 'sidb'
    ports:
      - "3306:3306"
    healthcheck:
      test: [ "CMD", "mysqladmin" ,"ping", "-h", "localhost" ]
      timeout: 10s
      interval: 5s
      retries: 10

数据库配置(application.properties)

spring.datasource.url=jdbc:mysql://localhost:3306/sidb
spring.datasource.username=root
spring.datasource.password=root
spring.datasource.driver-class-name=com.mysql.cj.jdbc.Driver

异常信息

org.springframework.dao.CannotAcquireLockException: PreparedStatementCallback; SQL [INSERT into _spring_integration_CHANNEL_MESSAGE(
    MESSAGE_ID,
    GROUP_KEY,
    REGION,
    CREATED_DATE,
    MESSAGE_PRIORITY,
    MESSAGE_BYTES)
values (?, ?, ?, ?, ?, ?)
]; Lock wait timeout exceeded; try restarting transaction
    at org.springframework.jdbc.support.SQLExceptionSubclassTranslator.doTranslate(SQLExceptionSubclassTranslator.java:78) ~[spring-jdbc-6.1.1.jar:6.1.1]
    at org.springframework.jdbc.support.AbstractFallbackSQLExceptionTranslator.translate(AbstractFallbackSQLExceptionTranslator.java:107) ~[spring-jdbc-6.1.1.jar:6.1.1]

解决方案

1. 拆分事务边界,避免跨通道锁竞争

当前配置中routeToRecipients和下游轮询器都开启了transactional(),导致消息从mailbox写入outbox的事务,与轮询器从outbox取消息的事务同时竞争数据库锁。移除routeToRecipients上的事务配置,让消息写入outbox的操作独立,轮询器事务仅覆盖消息处理阶段:

@Bean
public IntegrationFlow buildFlow(ChannelMessageStore channelMessageStore,
                                 MessageHandler sendEmailMessageHandler) {
    return IntegrationFlow.from("mailbox")
                          .routeToRecipients(routes -> {
                              routes
                                      .recipientFlow(flow -> flow
                                              .channel(channels -> channels.queue(channelMessageStore, "outbox"))
                                              .handle(sendEmailMessageHandler, e -> e.poller(poller -> poller.fixedDelay(1000).transactional()))
                                      );
                          }).get();
}

2. 配置JDBC消息存储的锁超时时间

显式设置JdbcChannelMessageStore的锁超时,避免事务持有锁过久:

@Bean
public JdbcChannelMessageStore messageStore(DataSource dataSource) {
    JdbcChannelMessageStore jdbcChannelMessageStore = new JdbcChannelMessageStore(dataSource);
    jdbcChannelMessageStore.setTablePrefix(CONCURRENT_METADATA_STORE_PREFIX);
    jdbcChannelMessageStore.setChannelMessageStoreQueryProvider(new MySqlChannelMessageStoreQueryProvider());
    // 设置锁超时时间(单位:毫秒)
    jdbcChannelMessageStore.setLockTimeout(5000);
    return jdbcChannelMessageStore;
}

3. 最小化事务粒度

确保轮询器事务仅包含必要的操作:消息取出、处理、失败回滚。如果邮件处理逻辑中有耗时操作,将其移出事务边界,避免锁持有时间过长。

4. 调整MySQL锁超时参数

在MySQL配置文件中修改innodb_lock_wait_timeout参数(默认50秒),根据业务场景适当延长,但不建议设置过大:

# my.cnf 配置示例
innodb_lock_wait_timeout = 60

5. 结合本地重试优化重试逻辑

在消息处理器上添加本地重试,减少轮询器的锁竞争频率:

.handle(sendEmailMessageHandler, e -> e
        .poller(poller -> poller.fixedDelay(1000).transactional())
        .retry(retry -> retry.maxAttempts(3).backOff(backOff -> backOff.fixedDelay(1000)))
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 17:07:34