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

