You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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()

验证方法

写入完成后,可通过以下方式确认压缩类型:

  1. 查看生成的Parquet文件名,应包含zstd标识(如part-xxxx-xxxx.c000.zstd.parquet)
  2. 使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 15:29:58