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

如何避免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语句原子更新状态,确保状态一致性。
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 12:56:00