AWS EMR Spark作业执行S3操作遇SdkInterruptedException问题排查
AWS EMR Spark作业中S3 getObject抛出SdkInterruptedException的问题分析与解决
问题描述
在AWS EMR集群上部署Spark作业时,执行以下代码行时抛出com.amazonaws.http.timers.client.SdkInterruptedException异常,该异常被包装在com.amazonaws.AbortedException中:
S3Object response = s3client.getObject(new GetObjectRequest(s3URI.getBucket(), s3URI.getKey()));
版本信息
- Java版本:1.8
- AWS SDK版本:1.11.1026
依赖
com.amazonaws:aws-java-sdk-bundle:jar:1.11.1026
相关代码
public static ByteArrayOutputStream getS3Object(String path) { AmazonS3 s3client = null; try { AmazonS3URI s3URI = new AmazonS3URI(path); s3client = AmazonS3ClientBuilder.defaultClient(); S3Object response = s3client.getObject(new GetObjectRequest(s3URI.getBucket(), s3URI.getKey())); ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream(); byte[] buffer = new byte[1024]; int bytesRead; while ((bytesRead = response.getObjectContent().read(buffer)) != -1) { byteArrayOutputStream.write(buffer, 0, bytesRead); } response.close(); return byteArrayOutputStream; } catch (AbortedException | IOException e) { e.printStackTrace(); System.err.println("AWS call interrupted retrying " + e); } finally { if (s3client != null) { s3client.shutdown(); } } throw new RuntimeException("s3 object is null"); }
堆栈跟踪
com.amazonaws.AbortedException: at com.amazonaws.http.AmazonHttpClient$RequestExecutor.handleInterruptedException(AmazonHttpClient.java:880) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.execute(AmazonHttpClient.java:757) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.access$500(AmazonHttpClient.java:715) at com.amazonaws.http.AmazonHttpClient$RequestExecutionBuilderImpl.execute(AmazonHttpClient.java:697) at com.amazonaws.http.AmazonHttpClient.execute(AmazonHttpClient.java:561) at com.amazonaws.http.AmazonHttpClient.execute(AmazonHttpClient.java:541) at com.amazonaws.services.s3.AmazonS3Client.invoke(AmazonS3Client.java:5456) at com.amazonaws.services.s3.AmazonS3Client.invoke(AmazonS3Client.java:5403) at com.amazonaws.services.s3.AmazonS3Client.getObject(AmazonS3Client.java:1524) com.example.S3Utils.getS3Object(S3Utils.java:78) Caused by: com.amazonaws.http.timers.client.SdkInterruptedException at com.amazonaws.http.AmazonHttpClient$RequestExecutor.checkInterrupted(AmazonHttpClient.java:935) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.checkInterrupted(AmazonHttpClient.java:921) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.executeHelper(AmazonHttpClient.java:1115) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.doExecute(AmazonHttpClient.java:814) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.executeWithTimer(AmazonHttpClient.java:781) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.execute(AmazonHttpClient.java:755) ... 53 more
原因分析
- 请求超时触发中断:AWS SDK内置的请求执行计时器超时,主动中断当前请求线程,抛出
SdkInterruptedException并最终包装为AbortedException,默认超时配置可能不匹配EMR环境下的S3访问需求。 - Spark任务线程被中断:Spark作业任务可能因YARN资源回收、任务超时设置过短等原因被强制终止,导致S3请求线程被中断。
- S3客户端频繁创建销毁:代码中每次调用方法都新建S3客户端并在finally块shutdown,导致连接池无法复用,增加请求延迟和失败概率,易触发超时中断。
- 缺乏有效重试机制:catch块仅打印日志,标注"retrying"但未实现实际重试逻辑,无法应对临时网络波动或S3服务端延迟。
解决方案
1. 调整AWS SDK超时配置
通过ClientConfiguration设置合理的超时参数适配EMR环境:
ClientConfiguration config = new ClientConfiguration() .withConnectionTimeout(5000) // 连接超时5秒 .withSocketTimeout(30000) // Socket读取超时30秒 .withRequestTimeout(60000); // 整体请求执行超时60秒 AmazonS3 s3client = AmazonS3ClientBuilder.standard() .withClientConfiguration(config) .build();
2. 复用S3客户端
避免每次调用创建新客户端,可将客户端实例设为静态变量或通过Spark广播变量复用,减少连接开销:
private static final AmazonS3 s3client = AmazonS3ClientBuilder.standard() .withClientConfiguration(new ClientConfiguration() .withConnectionTimeout(5000) .withSocketTimeout(30000)) .build(); public static ByteArrayOutputStream getS3Object(String path) { try { AmazonS3URI s3URI = new AmazonS3URI(path); S3Object response = s3client.getObject(new GetObjectRequest(s3URI.getBucket(), s3URI.getKey())); // 后续读取逻辑不变 } catch (AbortedException | IOException e) { // 重试逻辑 } // 无需在finally中shutdown客户端 }
3. 实现有效重试机制
添加带退避策略的重试逻辑,比如手动实现循环重试:
public static ByteArrayOutputStream getS3Object(String path) { int retryCount = 3; // 重试3次 while (retryCount > 0) { try { AmazonS3URI s3URI = new AmazonS3URI(path); AmazonS3 s3client = AmazonS3ClientBuilder.defaultClient(); S3Object response = s3client.getObject(new GetObjectRequest(s3URI.getBucket(), s3URI.getKey())); ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream(); byte[] buffer = new byte[1024]; int bytesRead; while ((bytesRead = response.getObjectContent().read(buffer)) != -1) { byteArrayOutputStream.write(buffer, 0, bytesRead); } response.close(); s3client.shutdown(); return byteArrayOutputStream; } catch (AbortedException | IOException e) { retryCount--; System.err.println("AWS call failed, retrying " + retryCount + " times left: " + e); try { Thread.sleep(1000 * (3 - retryCount)); // 指数退避 } catch (InterruptedException ie) { Thread.currentThread().interrupt(); } } } throw new RuntimeException("Failed to get S3 object after retries"); }
4. 检查EMR与Spark配置
- 调整Spark任务超时:设置
spark.executor.heartbeatInterval、spark.network.timeout等参数,避免任务因心跳超时被终止。 - 确保EMR集群资源充足:避免因节点资源耗尽导致任务被YARN杀死。
5. 升级AWS SDK版本
当前使用的1.11.1026版本较旧,升级到1.11.x系列最新版本或迁移到AWS SDK v2,可修复已知的超时和线程中断相关bug。
内容的提问来源于stack exchange,提问作者Vasanth Subramanian
相关产品推荐
相关产品推荐

