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

