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

