PySpark写入RedShift失败问题求助
尝试通过PySpark将数据写入RedShift,Spark版本3.2.0,Scala版本2.12.15,按照官方指南操作(包括aws_iam_role写入方式)均出现相同错误,所有依赖已匹配Scala 2.12版本。
环境信息
- Spark 3.2
- Scala 2.12.15
- PySpark 3.2.3
- Java 11
- Ubuntu 22.04 LTS
- Python 3.8
代码片段
from pyspark.sql import SparkSession spark = SparkSession.builder.appName('abc')\ .config("spark.jars.packages","com.eclipsesource.minimal-json:minimal-json:0.9.5,com.amazon.redshift:redshift-jdbc42:2.1.0.12,com.google.guava:guava:31.1-jre,com.amazonaws:aws-java-sdk-s3:1.12.437,org.apache.spark:spark-avro_2.12:3.3.2,io.github.spark-redshift-community:spark-redshift_2.12:5.1.0,org.apache.hadoop:hadoop-aws:3.2.2,com.google.guava:failureaccess:1.0")\ .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") \ .config("spark.hadoop.fs.s3a.access.key", "etc") \ .config("spark.hadoop.fs.s3a.secret.key", "etc") \ .config('spark.hadoop.fs.s3a.aws.credentials.provider', 'org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider')\ .getOrCreate() df=spark.read.option("header",True) \ .csv("demo.csv") df.write \ .format("io.github.spark_redshift_community.spark.redshift") \ .option("url", "jdbc:redshift:iam://host:5439/dev?user=user&password=pass") \ .option("dbtable", "demo") \ .option("forward_spark_s3_credentials","True") \ .option("tempdir", "s3a://mubucket/folder") \ .mode("append") \ .save()
错误日志
23/03/30 18:51:47 WARN MetricsConfig: Cannot locate configuration: tried hadoop-metrics2-s3a-file-system.properties,hadoop-metrics2.properties 23/03/30 18:51:50 WARN Utils$: The S3 bucket demo does not have an object lifecycle configuration to ensure cleanup of temporary files. Consider configuring `tempdir` to point to a bucket with an object lifecycle policy that automatically deletes files after an expiration period. For more information, see https://docs.aws.amazon.com/AmazonS3/latest/dev/object-lifecycle-mgmt.html 23/03/30 18:51:51 WARN AbstractS3ACommitterFactory: Using standard FileOutputCommitter to commit work. This is slow and potentially unsafe. 23/03/30 18:51:53 WARN AbstractS3ACommitterFactory: Using standard FileOutputCommitter to commit work. This is slow and potentially unsafe. 23/03/30 18:51:53 WARN AbstractS3ACommitterFactory: Using standard FileOutputCommitter to commit work. This is slow and potentially unsafe. 23/03/30 18:51:54 ERROR Utils: Aborting task java.lang.NoSuchMethodError: 'scala.Function1 org.apache.spark.sql.execution.datasources.DataSourceUtils$.createDateRebaseFuncInWrite(scala.Enumeration$Value, java.lang.String)'
注:已移除敏感凭证,相同凭证可正常创建数据库/表,且拥有S3完全权限,多次尝试不同方法均报错。
解决方案
核心问题
错误根源是Spark依赖版本不兼容:你使用的Spark版本是3.2.x,但引入的spark-avro版本是3.3.2。createDateRebaseFuncInWrite方法是Spark 3.3才新增的API,Spark 3.2的DataSourceUtils类中没有这个方法,导致运行时抛出找不到方法的错误。
修复步骤
匹配spark-avro版本与Spark版本
将spark.jars.packages中的org.apache.spark:spark-avro_2.12:3.3.2替换为与你的Spark/PySpark版本一致的3.2.x版本,比如org.apache.spark:spark-avro_2.12:3.2.3(与PySpark 3.2.3对应)。调整AWS SDK版本以匹配hadoop-aws
hadoop-aws:3.2.2对应的AWS SDK推荐版本是1.11.x系列,将com.amazonaws:aws-java-sdk-s3:1.12.437替换为com.amazonaws:aws-java-sdk-s3:1.11.901,避免版本冲突。
修改后的spark.jars.packages配置如下:
.config("spark.jars.packages","com.eclipsesource.minimal-json:minimal-json:0.9.5,com.amazon.redshift:redshift-jdbc42:2.1.0.12,com.google.guava:guava:31.1-jre,com.amazonaws:aws-java-sdk-s3:1.11.901,org.apache.spark:spark-avro_2.12:3.2.3,io.github.spark-redshift-community:spark-redshift_2.12:5.1.0,org.apache.hadoop:hadoop-aws:3.2.2,com.google.guava:failureaccess:1.0")
- 可选:优化S3提交器(解决警告)
添加以下配置替换默认的FileOutputCommitter,提升S3写入性能与安全性:.config("spark.hadoop.fs.s3a.committer.name", "directory") \ .config("spark.sql.sources.commitProtocolClass", "org.apache.spark.sql.execution.datasources.SQLHadoopMapReduceCommitProtocol")
内容的提问来源于stack exchange,提问作者digital_monk

