向DigitalOcean Spaces写入PySpark DataFrame时出现403禁止访问错误
PySpark写入DigitalOcean Spaces报403 Forbidden错误排查
问题概述
尝试将PySpark DataFrame写入DigitalOcean Spaces时触发Forbidden (403)错误,堆栈跟踪显示权限验证失败,但使用Python boto客户端通过相同密钥可正常访问该Spaces。当前使用PySpark 3.5,jar包配置遵循Hadoop-AWS版本兼容规则。
代码实现
from pyspark.sql import SparkSession def get_spark() -> SparkSession: """Provides a well configured spark session""" return ( SparkSession.builder.master("local[*]") .appName("test") .config("spark.jars.packages", "io.delta:delta-spark_2.12:3.0.0") .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config( "spark.jars.packages", "io.delta:delta-spark_2.12:3.0.0,org.apache.hadoop:hadoop-aws:3.3.4,com.amazonaws:aws-java-sdk-bundle:1.12.262", ) .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") .config("spark.hadoop.fs.s3a.access.key", "KEY") .config("spark.hadoop.fs.s3a.secret.key", "SECRET") .config("spark.haddop.fs.s3a.endpoint", "https://REGION.digitaloceanspaces.com") .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") .config("spark.sql.warehouse.dir", WAREHOUSE_LOCATION) .enableHiveSupport() .getOrCreate() ) def test_spaces(): """Creates the NHL roster table in the bronze layer""" spark = get_spark() # Create a simple DataFrame data = [("John", 25), ("Alice", 30), ("Bob", 28)] columns = ["Name", "Age"] df = spark.createDataFrame(data, columns) # Show the DataFrame df.show() # Write DataFrame to DigitalOcean Spaces df.write.json(f"s3a://bucket_name/test")
错误堆栈跟踪
py4j.protocol.Py4JJavaError: An error occurred while calling o65.json. : java.nio.file.AccessDeniedException: s3a://PATH: getFileStatus on s3a://PATH: com.amazonaws.services.s3.model.AmazonS3Exception: Forbidden (Service: Amazon S3; Status Code: 403; Error Code: 403 Forbidden; Request ID: 069ZFJ7PEE4SDT1B; S3 Extended Request ID: iExUnbCSQn8Tued6vyOcmZvu7BMm/6NVRWtopsmAHk572kPJxY5lV8C4BSkalexpg/18EgWnpAkpf2bTUElcxQ==; Proxy: null), S3 Extended Request ID: iExUnbCSQn8Tued6vyOcmZvu7BMm/6NVRWtopsmAHk572kPJxY5lV8C4BSkalexpg/18EgWnpAkpf2bTUElcxQ==:403 Forbidden at org.apache.hadoop.fs.s3a.S3AUtils.translateException(S3AUtils.java:255) at org.apache.hadoop.fs.s3a.S3AUtils.translateException(S3AUtils.java:175) at org.apache.hadoop.fs.s3a.S3AFileSystem.s3GetFileStatus(S3AFileSystem.java:3796) at org.apache.hadoop.fs.s3a.S3AFileSystem.innerGetFileStatus(S3AFileSystem.java:3688) at org.apache.hadoop.fs.s3a.S3AFileSystem.lambda$exists$34(S3AFileSystem.java:4703) at org.apache.hadoop.fs.statistics.impl.IOStatisticsBinding.lambda$trackDurationOfOperation$5(IOStatisticsBinding.java:499) at org.apache.hadoop.fs.statistics.impl.IOStatisticsBinding.trackDuration(IOStatisticsBinding.java:444) at org.apache.hadoop.fs.s3a.S3AFileSystem.trackDurationAndSpan(S3AFileSystem.java:2337) at org.apache.hadoop.fs.s3a.S3AFileSystem.trackDurationAndSpan(S3AFileSystem.java:2356) at org.apache.hadoop.fs.s3a.S3AFileSystem.exists(S3AFileSystem.java:4701) at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand.run(InsertIntoHadoopFsRelationCommand.scala:120) 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:107) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:125) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:201) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:108) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:900) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:66) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:107) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:98) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:461) at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(origin.scala:76) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:461) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDownWithPruning(LogicalPlan.scala:32) 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:32) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:32) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:437) at org.apache.spark.sql.execution.QueryExecution.eagerlyExecuteCommands(QueryExecution.scala:98) at org.apache.spark.sql.execution.QueryExecution.commandExecuted$lzycompute(QueryExecution.scala:85) at org.apache.spark.sql.execution.QueryExecution.commandExecuted(QueryExecution.scala:83) at org.apache.spark.sql.execution.QueryExecution.assertCommandExecuted(QueryExecution.scala:142) at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:859) at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:388) at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:361) at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:240) at org.apache.spark.sql.DataFrameWriter.json(DataFrameWriter.scala:774)
解决方案
1. 修复配置拼写错误
代码中spark.haddop.fs.s3a.endpoint存在拼写错误,正确配置应为spark.hadoop.fs.s3a.endpoint。拼写错误会导致Spark无法识别DigitalOcean的存储端点,默认使用AWS S3地址引发权限验证失败。
2. 添加DigitalOcean专属配置
针对DigitalOcean的S3兼容存储,需额外配置路径样式访问和签名算法:
# 添加到SparkSession.builder配置中 .config("spark.hadoop.fs.s3a.path.style.access", "true") .config("spark.hadoop.fs.s3a.signing-algorithm", "S3SignerType")
3. 合并重复配置
代码中重复设置了spark.jars.packages和spark.sql.extensions,建议合并为单次配置避免冲突:
.config( "spark.jars.packages", "io.delta:delta-spark_2.12:3.0.0,org.apache.hadoop:hadoop-aws:3.3.4,com.amazonaws:aws-java-sdk-bundle:1.12.262", ) .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
修正后的get_spark函数
def get_spark() -> SparkSession: """Provides a well configured spark session""" return ( SparkSession.builder.master("local[*]") .appName("test") .config( "spark.jars.packages", "io.delta:delta-spark_2.12:3.0.0,org.apache.hadoop:hadoop-aws:3.3.4,com.amazonaws:aws-java-sdk-bundle:1.12.262", ) .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") .config("spark.hadoop.fs.s3a.access.key", "KEY") .config("spark.hadoop.fs.s3a.secret.key", "SECRET") .config("spark.hadoop.fs.s3a.endpoint", "https://REGION.digitaloceanspaces.com") .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") .config("spark.hadoop.fs.s3a.path.style.access", "true") .config("spark.hadoop.fs.s3a.signing-algorithm", "S3SignerType") .config("spark.sql.warehouse.dir", WAREHOUSE_LOCATION) .enableHiveSupport() .getOrCreate() )
内容的提问来源于stack exchange,提问作者nerdizzle
相关产品推荐
相关产品推荐

