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

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

原因分析

  1. 请求超时触发中断:AWS SDK内置的请求执行计时器超时,主动中断当前请求线程,抛出SdkInterruptedException并最终包装为AbortedException,默认超时配置可能不匹配EMR环境下的S3访问需求。
  2. Spark任务线程被中断:Spark作业任务可能因YARN资源回收、任务超时设置过短等原因被强制终止,导致S3请求线程被中断。
  3. S3客户端频繁创建销毁:代码中每次调用方法都新建S3客户端并在finally块shutdown,导致连接池无法复用,增加请求延迟和失败概率,易触发超时中断。
  4. 缺乏有效重试机制: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 15:45:55