如何避免Spring Integration多实例重复处理同一SFTP文件?
Spring Integration多实例文件抓取冲突解决方案
问题背景
当前使用Spring Integration实现文件传输、Spring Batch处理文件,最后通过Spring Integration标记文件完成。现有配置包括:
- 采用
JdbcMetadataStore存储元数据 - 使用
SftpSimplePatternFileListFilter和SftpPersistentAcceptOnceFileListFilter进行文件过滤 - 数据库为PostgreSQL
多实例部署时,出现多个实例同时抓取同一文件导致流程中断的问题,尝试过固定延迟和cron调度仍无法解决,需要实现可靠的锁机制保证实例间文件分发无同步问题。
当前Spring Integration Flow配置
IntegrationFlow.from( Sftp.inboundAdapter(sftpSessionFactory) .remoteDirectory(".") .deleteRemoteFiles(false) .filter(remoteFileFilter) .localDirectory(new File(localDirectory)) .autoCreateLocalDirectory(true) .maxFetchSize(1), configurer -> configurer.poller(Pollers.cron("some-cron-here") .advice(rotatingServerAdvice))) .transform(transformerService, "toJobLaunchRequest") .channel(jobLaunchingChannel) .get();
事务配置下的错误栈
2024-09-13T10:50:02.499-04:00 ERROR 16323 --- [some-service] [ scheduling-1] o.s.integration.handler.LoggingHandler : org.springframework.messaging.MessagingException: Problem occurred while synchronizing '/remote-file-path' to local directory at org.springframework.integration.file.remote.synchronizer.AbstractInboundFileSynchronizer.synchronizeToLocalDirectory(AbstractInboundFileSynchronizer.java:345) at org.springframework.integration.file.remote.synchronizer.AbstractInboundFileSynchronizingMessageSource.doReceive(AbstractInboundFileSynchronizingMessageSource.java:266) at org.springframework.integration.file.remote.synchronizer.AbstractInboundFileSynchronizingMessageSource.doReceive(AbstractInboundFileSynchronizingMessageSource.java:69) at org.springframework.integration.endpoint.AbstractFetchLimitingMessageSource.doReceive(AbstractFetchLimitingMessageSource.java:47) at org.springframework.integration.endpoint.AbstractMessageSource.receive(AbstractMessageSource.java:142) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:568) at org.springframework.aop.support.AopUtils.invokeJoinpointUsingReflection(AopUtils.java:355) at org.springframework.aop.framework.ReflectiveMethodInvocation.invokeJoinpoint(ReflectiveMethodInvocation.java:196) at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:163) at org.springframework.integration.aop.ReceiveMessageAdvice.invoke(ReceiveMessageAdvice.java:56) at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:184) at org.springframework.aop.framework.JdkDynamicAopProxy.invoke(JdkDynamicAopProxy.java:223) at jdk.proxy2/jdk.proxy2.$Proxy118.receive(Unknown Source) at org.springframework.integration.endpoint.SourcePollingChannelAdapter.receiveMessage(SourcePollingChannelAdapter.java:220) at org.springframework.integration.endpoint.AbstractPollingEndpoint.doPoll(AbstractPollingEndpoint.java:450) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:568) at org.springframework.aop.support.AopUtils.invokeJoinpointUsingReflection(AopUtils.java:355) at org.springframework.aop.framework.ReflectiveMethodInvocation.invokeJoinpoint(ReflectiveMethodInvocation.java:196) at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:163) at org.springframework.transaction.interceptor.TransactionAspectSupport.invokeWithinTransaction(TransactionAspectSupport.java:379) at org.springframework.transaction.interceptor.TransactionInterceptor.invoke(TransactionInterceptor.java:119) at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:184) at org.springframework.aop.framework.JdkDynamicAopProxy.invoke(JdkDynamicAopProxy.java:223) at jdk.proxy2/jdk.proxy2.$Proxy117.call(Unknown Source) at org.springframework.integration.endpoint.AbstractPollingEndpoint.pollForMessage(AbstractPollingEndpoint.java:419) at org.springframework.integration.endpoint.AbstractPollingEndpoint.lambda$createPoller$4(AbstractPollingEndpoint.java:355) at org.springframework.integration.util.ErrorHandlingTaskExecutor.lambda$execute$0(ErrorHandlingTaskExecutor.java:56) at org.springframework.core.task.SyncTaskExecutor.execute(SyncTaskExecutor.java:50) at org.springframework.integration.util.ErrorHandlingTaskExecutor.execute(ErrorHandlingTaskExecutor.java:54) at org.springframework.integration.endpoint.AbstractPollingEndpoint.lambda$createPoller$5(AbstractPollingEndpoint.java:348) at org.springframework.scheduling.support.DelegatingErrorHandlingRunnable.run(DelegatingErrorHandlingRunnable.java:54) at org.springframework.scheduling.concurrent.ReschedulingRunnable.run(ReschedulingRunnable.java:96) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:842) Caused by: org.springframework.messaging.MessagingException: Failed to execute on session at org.springframework.integration.file.remote.RemoteFileTemplate.execute(RemoteFileTemplate.java:459) at org.springframework.integration.file.remote.synchronizer.AbstractInboundFileSynchronizer.synchronizeToLocalDirectory(AbstractInboundFileSynchronizer.java:338) ... 43 more Caused by: org.springframework.jdbc.UncategorizedSQLException: PreparedStatementCallback; uncategorized SQLException for SQL [SELECT METADATA_VALUE FROM INT_METADATA_STORE WHERE METADATA_KEY=? AND REGION=? ]; SQL state [25P02]; error code [0]; ERROR: current transaction is aborted, commands ignored until end of transaction block at org.springframework.jdbc.core.JdbcTemplate.translateException(JdbcTemplate.java:1549) at org.springframework.jdbc.core.JdbcTemplate.execute(JdbcTemplate.java:677) at org.springframework.jdbc.core.JdbcTemplate.query(JdbcTemplate.java:723) at org.springframework.jdbc.core.JdbcTemplate.query(JdbcTemplate.java:754) at org.springframework.jdbc.core.JdbcTemplate.query(JdbcTemplate.java:767) at org.springframework.jdbc.core.JdbcTemplate.queryForObject(JdbcTemplate.java:889) at org.springframework.jdbc.core.JdbcTemplate.queryForObject(JdbcTemplate.java:916) at org.springframework.integration.jdbc.metadata.JdbcMetadataStore.putIfAbsent(JdbcMetadataStore.java:230) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:568) at org.springframework.aop.support.AopUtils.invokeJoinpointUsingReflection(AopUtils.java:355) at org.springframework.aop.framework.ReflectiveMethodInvocation.invokeJoinpoint(ReflectiveMethodInvocation.java:196) at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:163) at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.proceed(CglibAopProxy.java:768) at org.springframework.transaction.interceptor.TransactionAspectSupport.invokeWithinTransaction(TransactionAspectSupport.java:379) at org.springframework.transaction.interceptor.TransactionInterceptor.invoke(TransactionInterceptor.java:119) at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:184) at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.proceed(CglibAopProxy.java:768) at org.springframework.aop.framework.CglibAopProxy$DynamicAdvisedInterceptor.intercept(CglibAopProxy.java:720) at org.springframework.integration.jdbc.metadata.JdbcMetadataStore$$SpringCGLIB$$0.putIfAbsent(<generated>) at org.springframework.integration.file.filters.AbstractPersistentAcceptOnceFileListFilter.accept(AbstractPersistentAcceptOnceFileListFilter.java:83) at org.springframework.integration.file.filters.ChainFileListFilter.accept(ChainFileListFilter.java:68) at org.springframework.integration.file.remote.synchronizer.AbstractInboundFileSynchronizer.transferFilesFromRemoteToLocal(AbstractInboundFileSynchronizer.java:374) at org.springframework.integration.file.remote.synchronizer.AbstractInboundFileSynchronizer.lambda$synchronizeToLocalDirectory$0(AbstractInboundFileSynchronizer.java:339) at org.springframework.integration.file.remote.RemoteFileTemplate.execute(RemoteFileTemplate.java:450) ... 44 more Caused by: org.postgresql.util.PSQLException: ERROR: current transaction is aborted, commands ignored until end of transaction block at org.postgresql.core.v3.QueryExecutorImpl.receiveErrorResponse(QueryExecutorImpl.java:2733) at org.postgresql.core.v3.QueryExecutorImpl.processResults(QueryExecutorImpl.java:2420) at org.postgresql.core.v3.QueryExecutorImpl.execute(QueryExecutorImpl.java:372) at org.postgresql.jdbc.PgStatement.executeInternal(PgStatement.java:517) at org.postgresql.jdbc.PgStatement.execute(PgStatement.java:434) at org.postgresql.jdbc.PgPreparedStatement.executeWithFlags(PgPreparedStatement.java:194) at org.postgresql.jdbc.PgPreparedStatement.executeQuery(PgPreparedStatement.java:137) at com.zaxxer.hikari.pool.ProxyPreparedStatement.executeQuery(ProxyPreparedStatement.java:52) at com.zaxxer.hikari.pool.HikariProxyPreparedStatement.executeQuery(HikariProxyPreparedStatement.java) at org.springframework.jdbc.core.JdbcTemplate$1.doInPreparedStatement(JdbcTemplate.java:732) at org.springframework.jdbc.core.JdbcTemplate.execute(JdbcTemplate.java:658) ... 69 more Caused by: org.postgresql.util.PSQLException: ERROR: duplicate key value violates unique constraint "int_metadata_store_pk" Detail: Key (metadata_key, region)=(remote:SimpleTestFile (94).txt, DEFAULT) already exists. at org.postgresql.core.v3.QueryExecutorImpl.receiveErrorResponse(QueryExecutorImpl.java:2733) at org.postgresql.core.v3.QueryExecutorImpl.processResults(QueryExecutorImpl.java:2420) at org.postgresql.core.v3.QueryExecutorImpl.execute(QueryExecutorImpl.java:372) at org.postgresql.jdbc.PgStatement.executeInternal(PgStatement.java:517) at org.postgresql.jdbc.PgStatement.execute(PgStatement.java:434) at org.postgresql.jdbc.PgPreparedStatement.executeWithFlags(PgPreparedStatement.java:194) at org.postgresql.jdbc.PgPreparedStatement.executeUpdate(PgPreparedStatement.java:155) at com.zaxxer.hikari.pool.ProxyPreparedStatement.executeUpdate(ProxyPreparedStatement.java:61) at com.zaxxer.hikari.pool.HikariProxyPreparedStatement.executeUpdate(HikariProxyPreparedStatement.java) at org.springframework.jdbc.core.JdbcTemplate.lambda$update$2(JdbcTemplate.java:975) at org.springframework.jdbc.core.JdbcTemplate.execute(JdbcTemplate.java:658) at org.springframework.jdbc.core.JdbcTemplate.update(JdbcTemplate.java:970) at org.springframework.jdbc.core.JdbcTemplate.update(JdbcTemplate.java:1014) at org.springframework.integration.jdbc.metadata.JdbcMetadataStore.tryToPutIfAbsent(JdbcMetadataStore.java:241) at org.springframework.integration.jdbc.metadata.JdbcMetadataStore.putIfAbsent(JdbcMetadataStore.java:222) ... 63 more
解决方案与最佳实践
1. 修复JdbcMetadataStore事务问题
错误栈显示,当多个实例同时插入同一元数据记录时,唯一键冲突导致事务中断,后续SQL无法执行。
解决方法:
- 为
JdbcMetadataStore配置独立事务,避免与轮询的全局事务绑定。可通过TransactionProxyFactoryBean为其单独设置事务属性,确保putIfAbsent操作在独立事务中执行,异常不会扩散到全局轮询事务。 - 升级Spring Integration到5.5.15+或6.0.5+版本,这些版本优化了
JdbcMetadataStore的putIfAbsent错误处理逻辑,避免事务中断影响后续操作。
2. 引入分布式锁实现实例互斥
方式一:Spring Integration原生锁注册表
配置基于PostgreSQL的分布式锁注册表:
@Bean public LockRegistry jdbcLockRegistry(DataSource dataSource) { return new JdbcLockRegistry(dataSource); }
创建锁通知,确保轮询操作互斥:
@Bean public ReceiveMessageAdvice lockAdvice(LockRegistry lockRegistry) { LockReceiverAdvice advice = new LockReceiverAdvice(lockRegistry); advice.setLockKey("sftp-file-poll-global-lock"); return advice; }
修改轮询器配置,添加锁通知:
configurer -> configurer.poller(Pollers.cron("some-cron-here") .advice(lockAdvice)))
方式二:数据库行级锁增强文件过滤
自定义过滤器,通过PostgreSQL的FOR UPDATE SKIP LOCKED实现无等待锁,确保只有一个实例能获取文件:
public class LockingSftpFileFilter extends SftpPersistentAcceptOnceFileListFilter { private final JdbcTemplate jdbcTemplate; public LockingSftpFileFilter(MetadataStore metadataStore, JdbcTemplate jdbcTemplate) { super(metadataStore); this.jdbcTemplate = jdbcTemplate; } @Override public boolean accept(ChannelSftp.LsEntry file) { String fileName = file.getFilename(); // 尝试获取文件锁,已被锁定的文件直接跳过 Integer exists = jdbcTemplate.queryForObject( "SELECT 1 FROM file_processing_lock WHERE file_name = ? FOR UPDATE SKIP LOCKED", Integer.class, fileName); if (exists == null) { try { jdbcTemplate.update( "INSERT INTO file_processing_lock(file_name, locked_at) VALUES (?, NOW())", fileName); return super.accept(file); } catch (DuplicateKeyException e) { return false; } } return false; } // 文件处理完成后释放锁 public void releaseLock(String fileName) { jdbcTemplate.update("DELETE FROM file_processing_lock WHERE file_name = ?", fileName); } }
在Spring Batch作业完成后,调用releaseLock释放文件锁。
3. 优化轮询与筛选策略
- 保持
maxFetchSize=1,每次轮询仅抓取一个文件,降低冲突概率。 - 替换
SftpSimplePatternFileListFilter为SftpRegexPatternFileListFilter,精准匹配目标文件,减少不必要的扫描。 - 为轮询器配置带重试的任务执行器,针对锁冲突、事务异常自动重试:
@Bean public ErrorHandlingTaskExecutor errorHandlingTaskExecutor() { SimpleAsyncTaskExecutor executor = new SimpleAsyncTaskExecutor(); return new ErrorHandlingTaskExecutor(executor, t -> { if (t instanceof UncategorizedSQLException || t instanceof DuplicateKeyException) { throw new RuntimeException("Retryable conflict", t); } }); }
更新轮询器配置:
configurer -> configurer.poller(Pollers.cron("some-cron-here") .advice(lockAdvice) .taskExecutor(errorHandlingTaskExecutor()) .maxMessagesPerPoll(1)))
4. 原子化标记文件状态
- 远程文件处理完成后,使用SFTP原子重命名操作将文件移动到
completed目录,避免状态变更冲突。 - 数据库维护文件状态表,使用
UPDATE ... WHERE语句原子更新状态,确保状态一致性。
相关产品推荐
相关产品推荐

