如何在AWS Glue中使用zstd压缩编码写入Delta Lake?
在AWS Glue 4.0中配置Delta Lake使用Zstandard(zstd)压缩
问题描述
使用AWS Glue 4.0(基于Spark 3.3)任务读取Parquet文件并写入Delta Lake时,尝试通过"write.parquet.compression-codec": "zstd"配置切换压缩编码,但生成的Parquet文件仍使用Snappy压缩,文件名也保留snappy标识。
原因分析
"write.parquet.compression-codec"是AWS Glue专属的Parquet Sink参数,仅在使用Glue原生的Parquet写入方式(如glueContext.write_dynamic_frame.from_options指定format="parquet")时生效。当使用format("delta")写入Delta Lake时,该参数不会被Delta Lake引擎识别,Delta Lake遵循Spark原生的Parquet压缩配置规则。
解决方案
以下两种方法均可正确配置Delta Lake使用zstd压缩:
方法1:设置SparkSession全局压缩配置
在初始化SparkSession后添加全局配置,所有后续的Parquet/Delta写入都会默认使用zstd压缩:
# 初始化SparkSession后添加这行配置 spark.conf.set("spark.sql.parquet.compression.codec", "zstd")
完整示例代码片段:
import sys from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job args = getResolvedOptions(sys.argv, ["JOB_NAME"]) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session # 添加全局压缩配置 spark.conf.set("spark.sql.parquet.compression.codec", "zstd") job = Job(glueContext) job.init(args["JOB_NAME"], args) S3bucket_node1 = glueContext.create_dynamic_frame.from_options( format_options={}, connection_type="s3", format="parquet", connection_options={ "paths": ["s3://my-bucket/data/raw-parquet/motor/"], "recurse": True, }, transformation_ctx="S3bucket_node1", ) additional_options = { "path": "s3://my-bucket/data/new-zstd-delta-tables/motor/", "mergeSchema": "true", } sink_to_delta_lake_node3_df = S3bucket_node1.toDF() sink_to_delta_lake_node3_df.write.format("delta").options(**additional_options).mode("overwrite").save() job.commit()
方法2:在Delta写入的options中直接指定compression参数
Delta Lake支持通过compression参数直接指定Parquet压缩编码,优先级高于全局配置:
additional_options = { "path": "s3://my-bucket/data/new-zstd-delta-tables/motor/", "compression": "zstd", # 替换原有的write.parquet.compression-codec "mergeSchema": "true", } sink_to_delta_lake_node3_df = S3bucket_node1.toDF() sink_to_delta_lake_node3_df.write.format("delta").options(**additional_options).mode("overwrite").save()
验证方法
写入完成后,可通过以下方式确认压缩类型:
- 查看生成的Parquet文件名,应包含
zstd标识(如part-xxxx-xxxx.c000.zstd.parquet) - 使用pyarrow验证:
import pyarrow.parquet as pq parquet_file_path = "part-xxxx-xxxx.c000.zstd.parquet" print(pq.ParquetFile(parquet_file_path).metadata.row_group(0).column(0).compression) # 应输出 "ZSTD"
内容的提问来源于stack exchange,提问作者Hongbo Miao
相关产品推荐
相关产品推荐

