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

Spring Batch读取AWS S3大CSV时遇连接重置问题求助

Spring Batch读取S3大文件时遭遇连接重置异常

我们有一个Spring Batch应用,需从AWS S3读取单份200MB的大型CSV文件并写入数据库。适配新CSV schema做小幅更新后,在LocalStack中运行导入流程一切正常,但切换到真实部署环境时出现连接重置异常,具体报错如下:

java.net.SocketException: Connection reset
    at sun.nio.ch.NioSocketImpl.implRead(NioSocketImpl.java:328)
    at sun.nio.ch.NioSocketImpl.read(NioSocketImpl.java:355)
    at sun.nio.ch.NioSocketImpl$1.read(NioSocketImpl.java:808)
    at java.net.Socket$SocketInputStream.read(Socket.java:966)
    at sun.security.ssl.SSLSocketInputRecord.read(SSLSocketInputRecord.java:484)
    at sun.security.ssl.SSLSocketInputRecord.readFully(SSLSocketInputRecord.java:467)
    at sun.security.ssl.SSLSocketInputRecord.decodeInputRecord(SSLSocketInputRecord.java:243)
    at sun.security.ssl.SSLSocketInputRecord.decode(SSLSocketInputRecord.java:181)
    at sun.security.ssl.SSLTransport.decode(SSLTransport.java:111)
    at sun.security.ssl.SSLSocketImpl.decode(SSLSocketImpl.java:1513)
    at sun.security.ssl.SSLSocketImpl.readApplicationRecord(SSLSocketImpl.java:1484)
    at sun.security.ssl.SSLSocketImpl$AppInputStream.read(SSLSocketImpl.java:1069)
    at org.apache.http.impl.io.SessionInputBufferImpl.streamRead(SessionInputBufferImpl.java:137)
    at org.apache.http.impl.io.SessionInputBufferImpl.read(SessionInputBufferImpl.java:197)
    at org.apache.http.impl.io.ContentLengthInputStream.read(ContentLengthInputStream.java:176)
    at org.apache.http.conn.EofSensorInputStream.read(EofSensorInputStream.java:135)
    at com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
    at com.amazonaws.event.ProgressInputStream.read(ProgressInputStream.java:180)
    at com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
    at com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
    at com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
    at com.amazonaws.event.ProgressInputStream.read(ProgressInputStream.java:180)
    at com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
    at com.amazonaws.util.LengthCheckInputStream.read(LengthCheckInputStream.java:107)
    at com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
    at com.amazonaws.services.s3.internal.S3AbortableInputStream.read(S3AbortableInputStream.java:125)
    at com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
    at sun.nio.cs.StreamDecoder.readBytes(StreamDecoder.java:270)
    at sun.nio.cs.StreamDecoder.implRead(StreamDecoder.java:313)
    at sun.nio.cs.StreamDecoder.read(StreamDecoder.java:188)
    at java.io.InputStreamReader.read(InputStreamReader.java:177)
    at java.io.BufferedReader.fill(BufferedReader.java:162)
    at java.io.BufferedReader.readLine(BufferedReader.java:329)
    at java.io.BufferedReader.readLine(BufferedReader.java:396)
    at org.springframework.batch.item.file.FlatFileItemReader.readLine(FlatFileItemReader.java:207)
    ... 33 common frames omitted
Wrapped by: org.springframework.batch.item.file.NonTransientFlatFileException: Unable to read from resource: [Amazon s3 resource [bucket='bucket-name-here' and object='import/test.csv']]
    at org.springframework.batch.item.file.FlatFileItemReader.readLine(FlatFileItemReader.java:226)
    at org.springframework.batch.item.file.FlatFileItemReader.doRead(FlatFileItemReader.java:178)
    at org.springframework.batch.item.support.AbstractItemCountingItemStreamItemReader.read(AbstractItemCountingItemStreamItemReader.java:93)
    at org.springframework.batch.item.support.SynchronizedItemStreamReader.read(SynchronizedItemStreamReader.java:57)
    at org.springframework.batch.item.support.SynchronizedItemStreamReader$$FastClassBySpringCGLIB$$987ea09.invoke(<generated>)
    at org.springframework.cglib.proxy.MethodProxy.invoke(MethodProxy.java:218)
    at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.invokeJoinpoint(CglibAopProxy.java:793)
    at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:163)
    at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.proceed(CglibAopProxy.java:763)
    at org.springframework.aop.support.DelegatingIntroductionInterceptor.doProceed(DelegatingIntroductionInterceptor.java:137)
    at org.springframework.aop.support.DelegatingIntroductionInterceptor.invoke(DelegatingIntroductionInterceptor.java:124)
    at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:186)
    at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.proceed(CglibAopProxy.java:763)
    at org.springframework.aop.framework.CglibAopProxy$DynamicAdvisedInterceptor.intercept(CglibAopProxy.java:708)
    at org.springframework.batch.item.support.SynchronizedItemStreamReader$$EnhancerBySpringCGLIB$$2a78c69a.read(<generated>)
    at org.springframework.batch.core.step.item.SimpleChunkProvider.doRead(SimpleChunkProvider.java:99)
    at org.springframework.batch.core.step.item.FaultTolerantChunkProvider.read(FaultTolerantChunkProvider.java:87)
    ... 17 common frames omitted
Wrapped by: org.springframework.batch.core.step.skip.NonSkippableReadException: Non-skippable exception during read
    at org.springframework.batch.core.step.item.FaultTolerantChunkProvider.read(FaultTolerantChunkProvider.java:105)
    at org.springframework.batch.core.step.item.SimpleChunkProvider$1.doInIteration(SimpleChunkProvider.java:126)
    at org.springframework.batch.repeat.support.RepeatTemplate.getNextResult(RepeatTemplate.java:375)
    at org.springframework.batch.repeat.support.RepeatTemplate.executeInternal(RepeatTemplate.java:215)
    at org.springframework.batch.repeat.support.RepeatTemplate.iterate(RepeatTemplate.java:145)
    at org.springframework.batch.core.step.item.SimpleChunkProvider.provide(SimpleChunkProvider.java:118)
    at org.springframework.batch.core.step.item.ChunkOrientedTasklet.execute(ChunkOrientedTasklet.java:71)
    at org.springframework.batch.core.step.tasklet.TaskletStep$ChunkTransactionCallback.doInTransaction(TaskletStep.java:407)
    at org.springframework.batch.core.step.tasklet.TaskletStep$ChunkTransactionCallback.doInTransaction(TaskletStep.java:331)
    at org.springframework.transaction.support.TransactionTemplate.execute(TransactionTemplate.java:140)
    at org.springframework.batch.core.step.tasklet.TaskletStep$2.doInChunkContext(TaskletStep.java:273)
    at org.springframework.batch.core.scope.context.StepContextRepeatCallback.doInIteration(StepContextRepeatCallback.java:82)
    at org.springframework.batch.repeat.support.TaskExecutorRepeatTemplate$ExecutingRunnable.run(TaskExecutorRepeatTemplate.java:262)
    at org.springframework.security.concurrent.DelegatingSecurityContextRunnable.run(DelegatingSecurityContextRunnable.java:82)
    at io.github.jhipster.async.ExceptionHandlingAsyncTaskExecutor.lambda$createWrappedRunnable$1(ExceptionHandlingAsyncTaskExecutor.java:78)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
    at java.lang.Thread.run(Thread.java:840)

当前配置

AmazonS3客户端配置

@Bean
@Primary
open fun amazonS3(awsCredentialsProvider: AWSCredentialsProvider): AmazonS3 {
        return AmazonS3ClientBuilder.standard()
            .withClientConfiguration(
                ClientConfiguration()
                    .withTcpKeepAlive(true)
                    .withRequestTimeout(MAX_TIMEOUT)
                    .withSocketTimeout(MAX_TIMEOUT)
            )
            .withCredentials(awsCredentialsProvider)
            .build()
}

S3文件获取方式

resourceLoader.getResource("s3://%s".formatted(filename));

已尝试无效的方案

  • 开启TcpKeepAlive并延长S3客户端超时时间
  • 自定义Spring Batch重试读取器实现重试逻辑

可行解决方案

1. 启用AWS SDK原生重试策略

AWS SDK内置的重试策略比自定义逻辑更适配S3服务特性,修改客户端配置添加重试参数:

@Bean
@Primary
open fun amazonS3(awsCredentialsProvider: AWSCredentialsProvider): AmazonS3 {
    val clientConfig = ClientConfiguration()
        .withTcpKeepAlive(true)
        .withRequestTimeout(MAX_TIMEOUT)
        .withSocketTimeout(MAX_TIMEOUT)
        .withMaxErrorRetry(5) // 设置最大重试次数
        .withRetryPolicy(PredefinedRetryPolicies.DEFAULT_RETRY_POLICY) // 用SDK默认重试策略,也可自定义
        .withConnectionTTL(300000) // 设置连接存活时间(5分钟),避免闲置被断开

    return AmazonS3ClientBuilder.standard()
        .withClientConfiguration(clientConfig)
        .withCredentials(awsCredentialsProvider)
        .build()
}

2. 使用S3分段下载

200MB文件可分段读取,避免单连接长时间占用被网络设备断开:

// 替代resourceLoader,直接用S3客户端分段获取对象
val getObjectRequest = GetObjectRequest(bucketName, fileName)
    .withRange(0, 1024*1024) // 每次读取1MB,可根据网络情况调整分段大小
val s3Object = amazonS3.getObject(getObjectRequest)
val inputStream = s3Object.objectContent
// 将inputStream传入FlatFileItemReader进行读取

也可结合Spring Batch的断点续读机制,实现断开后从上次分段位置继续读取。

3. 调整FlatFileItemReader缓冲配置

默认逐行读取可能导致连接长时间空闲,增大缓冲减少网络交互的空闲时长:

@Bean
fun flatFileItemReader(s3Resource: Resource): FlatFileItemReader<YourDto> {
    return FlatFileItemReaderBuilder<YourDto>()
        .resource(s3Resource)
        .name("csvReader")
        .delimited()
        .names(*yourColumnNames)
        .targetType(YourDto::class.java)
        .bufferedReaderFactory { inputStream ->
            BufferedReader(InputStreamReader(inputStream), 1024*1024) // 设置1MB缓冲
        }
        .build()
}

4. 排查网络环境限制

真实部署环境的防火墙、NAT网关或负载均衡可能存在连接超时限制(比如5分钟空闲断开),联系运维确认网络设备超时时间,调整客户端的socketTimeout和connectionTTL与之匹配。

5. 升级AWS SDK版本

旧版本AWS SDK可能存在连接管理bug,尝试升级到最新稳定版,或迁移到AWS SDK v2(v2的连接池和重试机制更高效)。

内容的提问来源于stack exchange,提问作者d-viso

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 15:45:54