在EMR集群用PySpark读取Delta Parquet遇spark_catalog错误求助
在普通EC2实例的SageMaker Notebook(镜像:DS 3.0,内核:Python3)中,手动安装PySpark 3.2.0并配置delta-core等依赖后,可正常读取S3上的Delta格式数据;但在EMR 6.6集群(Spark版本为3.2.0-amzn-0)中运行完全相同的代码时,出现spark_catalog相关错误,怀疑亚马逊定制版Spark与delta-core不兼容。
可正常运行场景:SageMaker EC2笔记本本地PySpark
先通过pip安装指定版本的PySpark:
%pip install pyspark==3.2.0
随后配置delta-core及S3依赖,实例化Spark会话:
from pyspark.sql import SparkSession from pyspark.sql.functions import col pkg_list = [ "io.delta:delta-core_2.12:1.1.0", "org.apache.hadoop:hadoop-aws:3.2.0", "com.amazonaws:aws-java-sdk-bundle:1.12.180" ] packages = ",".join(pkg_list) spark = (SparkSession.builder.appName("EDA") .config("spark.jars.packages", packages) .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.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") .getOrCreate() ) # 配置S3访问参数 spark.sparkContext._jsc.hadoopConfiguration().set("com.amazonaws.services.s3.enableV4", "true") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") spark.sparkContext._jsc.hadoopConfiguration().set("fs.AbstractFileSystem.s3a.impl", "org.apache.hadoop.fs.s3a.S3A") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.aws.credentials.provider", "com.amazonaws.auth.DefaultAWSCredentialsProviderChain") spark.sparkContext._jsc.hadoopConfiguration().set('fs.s3a.access.key', '<key>') spark.sparkContext._jsc.hadoopConfiguration().set('fs.s3a.secret.key', '<secret>') print(f"Spark version: {spark.sparkContext.version}")
输出:Spark version: 3.2.0
执行读取Delta表代码可正常运行:
df = spark.read.format("delta").load("s3a://[my-bucket]/tmp/tmp/lake1/").filter(col('language')=='English') df.show()
无法运行场景:EMR 6.6集群上的PySpark
使用EMR 6.6版本(对应Spark 3.2.0-amzn-0),运行上述完全相同的代码,输出Spark版本为Spark version: 3.2.0-amzn-0,但执行读取Delta表代码时出现spark_catalog相关错误。
解决方案
1. 移除冲突依赖
EMR集群已预装与自身版本兼容的hadoop-aws和aws-java-sdk-bundle,手动指定外部版本会导致依赖冲突,进而引发catalog初始化错误。需从依赖列表中移除这两个包。
2. 调整Spark Catalog配置
EMR默认使用Hive作为元数据catalog,直接替换为DeltaCatalog会与EMR的定制化catalog实现冲突。如果仅需读取Delta表而非使用Delta的catalog管理功能,可暂时移除spark.sql.catalog.spark_catalog配置,仅保留Delta扩展配置。
3. 简化S3访问配置
EMR实例通过IAM角色获取S3访问权限,无需硬编码密钥,移除手动设置的fs.s3a.access.key和fs.s3a.secret.key,依赖DefaultAWSCredentialsProviderChain自动获取凭证即可。
修改后的EMR代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import col # 仅保留delta-core依赖,移除EMR自带的冲突包 pkg_list = [ "io.delta:delta-core_2.12:1.1.0" ] packages = ",".join(pkg_list) spark = (SparkSession.builder.appName("EDA") .config("spark.jars.packages", packages) .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") # 移除spark_catalog配置,使用EMR默认HiveCatalog .getOrCreate() ) # 仅保留必要的S3配置 spark.sparkContext._jsc.hadoopConfiguration().set("com.amazonaws.services.s3.enableV4", "true") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") spark.sparkContext._jsc.hadoopConfiguration().set("fs.AbstractFileSystem.s3a.impl", "org.apache.hadoop.fs.s3a.S3A") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.aws.credentials.provider", "com.amazonaws.auth.DefaultAWSCredentialsProviderChain") print(f"Spark version: {spark.sparkContext.version}") # 读取Delta表 df = spark.read.format("delta").load("s3a://[my-bucket]/tmp/tmp/lake1/").filter(col('language')=='English') df.show()
更稳妥的方案:通过EMR应用安装Delta Lake
创建EMR集群时,在“应用程序”选项中添加Delta Lake,EMR会自动配置好所有兼容的依赖和参数,无需手动设置spark.jars.packages和扩展配置,彻底避免版本冲突问题。
内容的提问来源于stack exchange,提问作者geominded

