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

使用Spark Java+hadoop-aws写入Cassandra数据集到S3时遇AWSBadRequestException

问题:Spark通过hadoop-aws写入Amazon S3触发AWSBadRequestException排查

版本信息

  • hadoop-aws: 3.3.1
  • spark: 3.2.1

代码示例

SparkConf sparkConf = new SparkConf();

sparkConf.setMaster("local[*]")
        .set("spark.hadoop.fs.s3a.access.key", "my access key")
        .set("spark.hadoop.fs.s3a.secret.key", "my secret key")
        .set("spark.hadoop.fs.s3a.endpoint", "s3.us-west-2.amazonaws.com")
        .set(CassandraOptions.SPARK_CONF_HOST, getHost(env))
        .set(CassandraOptions.SPARK_CONF_PORT, getPort(env))
        .set(CassandraOptions.SPARK_CONF_LOCALDC, getDatacenter(env))
        .set(CassandraOptions.TABLE, table)
        .set(CassandraOptions.KEYSPACE, keyspace);
this.sparkSession = SparkSession.builder().config(sparkConf).getOrCreate();

// 获取Dataset<Row>数据集
// Dataset<Row> ds = dfr.load().where(filterCondition);
// ds.show();

String s3Path = "s3a://" + bucket_name + "/" + path;
ds.write().parquet(s3Path);

错误堆栈

Exception in thread "main" org.apache.hadoop.fs.s3a.AWSBadRequestException: getFileStatus on s3a://my-test-bucket/data/configuration-migration/cassandra/Oregon/shared_services/orgid_appurl/00Dx0000000JBysEAG: com.amazonaws.services.s3.model.AmazonS3Exception: Bad Request (Service: Amazon S3; Status Code: 400; Error Code: 400 Bad Request; Request ID: SAJPPB0RXB43CJ8Y; S3 Extended Request ID: vMUJ5utvguWgzUmBVGw80FAOPP0OBU9a5QFjRDEbJtaHMSl7qZm4+LZzloflAyzSh3Z6maEX6n8=; Proxy: null), S3 Extended Request ID: vMUJ5utvguWgzUmBVGw80FAOPP0OBU9a5QFjRDEbJtaHMSl7qZm4+LZzloflAyzSh3Z6maEX6n8=:400 Bad Request: Bad Request (Service: Amazon S3; Status Code: 400; Error Code: 400 Bad Request; Request ID: SAJPPB0RXB43CJ8Y; S3 Extended Request ID: vMUJ5utvguWgzUmBVGw80FAOPP0OBU9a5QFjRDEbJtaHMSl7qZm4+LZzloflAyzSh3Z6maEX6n8=; Proxy: null)
    at org.apache.hadoop.fs.s3a.S3AUtils.translateException(S3AUtils.java:243)
    at org.apache.hadoop.fs.s3a.S3AUtils.translateException(S3AUtils.java:170)
    at org.apache.hadoop.fs.s3a.S3AFileSystem.s3GetFileStatus(S3AFileSystem.java:3286)
    at org.apache.hadoop.fs.s3a.S3AFileSystem.innerGetFileStatus(S3AFileSystem.java:3185)
    at org.apache.hadoop.fs.s3a.S3AFileSystem.getFileStatus(S3AFileSystem.java:3053)
    at org.apache.hadoop.fs.FileSystem.exists(FileSystem.java:1760)
    at org.apache.hadoop.fs.s3a.S3AFileSystem.exists(S3AFileSystem.java:4263)
    at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand.run(InsertIntoHadoopFsRelationCommand.scala:117)
    at org.apache.spark.sql.execution.command.DataWritingCommandExec.sideEffectResult$lzycompute(commands.scala:113)
    at org.apache.spark.sql.execution.command.DataWritingCommandExec.sideEffectResult(commands.scala:111)
    at org.apache.spark.sql.execution.command.DataWritingCommandExec.executeCollect(commands.scala:125)
    at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.$anonfun$applyOrElse$1(QueryExecution.scala:110)
    at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$5(SQLExecution.scala:103)
    at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:163)
    at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:90)
    at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775)
    at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:64)
    at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:110)
    at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:106)
    at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:481)
    at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:82)
    at org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:481)
    at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDownWithPruning(LogicalPlan.scala:30)
    at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning(AnalysisHelper.scala:267)
    at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning$(AnalysisHelper.scala:263)
    at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:30)
    at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:30)
    at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:457)
    at org.apache.spark.sql.execution.QueryExecution.eagerlyExecuteCommands(QueryExecution.scala:106)
    at org.apache.spark.sql.execution.QueryExecution.commandExecuted$lzycompute(QueryExecution.scala:93)
    at org.apache.spark.sql.execution.QueryExecution.commandExecuted(QueryExecution.scala:91)
    at org.apache.spark.sql.execution.QueryExecution.assertCommandExecuted(QueryExecution.scala:128)
排查方案
  • 验证凭证有效性:确认配置的spark.hadoop.fs.s3a.access.key和spark.hadoop.fs.s3a.secret.key无多余空格或错误字符,可用AWS CLI测试:aws s3 ls s3://my-test-bucket --region us-west-2。
  • 调整Endpoint配置:AWS S3标准区域无需显式设置fs.s3a.endpoint,SDK会自动推断。若保留设置,确保us-west-2的Endpoint格式正确,或直接移除该配置项测试。
  • 检查桶权限:确保凭证对应的IAM用户/角色拥有目标桶的s3:PutObject、s3:GetObject、s3:ListBucket权限,策略资源范围需覆盖目标桶及子路径。
  • 配置签名版本:hadoop-aws 3.3.1默认签名版本可能与区域不兼容,添加配置:
    sparkConf.set("spark.hadoop.fs.s3a.signing-algorithm", "AWS4SignerType")
    
  • 验证路径合法性:确认S3路径无特殊字符(如空格、非ASCII字符),可尝试写入桶根目录测试,排除路径格式问题。
  • 调整依赖版本:spark 3.2.1与hadoop-aws 3.3.1可能存在兼容性问题,可尝试升级spark到3.3.x搭配hadoop-aws 3.3.x,或降级hadoop-aws到3.2.x。
  • 排查代理干扰:若环境存在代理,需正确配置spark.hadoop.fs.s3a.proxy.host和spark.hadoop.fs.s3a.proxy.port;无代理则确保无相关配置影响请求。

内容的提问来源于stack exchange,提问作者Jun Q

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 15:42:05