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

在AWS Glue环境中如何将PySpark DataFrame转换为Delta Table以执行Merge操作?

Absolutely, you can convert your transformed PySpark DataFrame into a Delta Table and perform upsert (merge) operations in AWS Glue—here's a step-by-step guide tailored to your use case:

1. First, Configure Delta Lake Dependencies in AWS Glue

AWS Glue doesn't include Delta Lake out of the box, so you need to add the required dependencies to your Glue Job:

  • Navigate to your Glue Job's Job parameters section
  • Add these key-value pairs:
    • --additional-python-modules: delta-spark==2.4.0 (match the version to your Glue Spark version; e.g., Glue 4.0 uses Spark 3.3, so Delta 2.4.0 works seamlessly)
    • --conf: spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension --conf spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog

2. Convert Your PySpark DataFrame to a Delta Table

You have two practical approaches here—either save to an S3 path or register directly as a Glue Catalog table:

Option 1: Save to an S3 Path (reference by path)

from delta.tables import DeltaTable

# Assume you already have your transformed PySpark DataFrame: transformed_df
# Save the DataFrame as Delta format to your S3 bucket
transformed_df.write.format("delta") \
    .mode("overwrite")  # Use "append" to add to existing data, or "overwrite" for initial table creation
    .save("s3://your-bucket/path/to/delta-table")

# Reference this Delta Table for merge operations
delta_table = DeltaTable.forPath(spark, "s3://your-bucket/path/to/delta-table")

This lets you use the Glue Catalog name to access the Delta Table instead of the raw S3 path:

# Save the transformed DataFrame as a managed/external Delta Table in Glue Catalog
transformed_df.write.format("delta") \
    .mode("overwrite")
    .saveAsTable("your_glue_database.your_delta_table_name")

# Reference the Delta Table using its Catalog name
delta_table = DeltaTable.forName(spark, "your_glue_database.your_delta_table_name")

3. Perform Upsert (Merge) Operation

Once you have your Delta Table (either the one created from your transformed DF or an existing target table), you can run the merge operation to upsert data. For example, merging your transformed DF into an existing target Delta Table:

# Get the target Delta Table (the table you want to upsert into)
target_delta_table = DeltaTable.forPath(spark, "s3://your-bucket/path/to/target-delta-table")
# OR use Catalog name: DeltaTable.forName(spark, "your_db.target_table")

# Execute the merge/upsert logic
target_delta_table.alias("target") \
    .merge(
        transformed_df.alias("source"),
        "target.primary_key = source.primary_key"  # Replace with your actual match condition (e.g., user_id, order_id)
    ) \
    .whenMatchedUpdate(set={
        # Define columns to update when a match is found
        "column1": "source.column1",
        "column2": "source.column2",
        "last_updated": current_timestamp()  # Optional: add an update timestamp
    }) \
    .whenNotMatchedInsert(values={
        # Define all columns to insert when no match is found
        "primary_key": "source.primary_key",
        "column1": "source.column1",
        "column2": "source.column2",
        "created_at": current_timestamp()
    }) \
    .execute()

Key Notes for AWS Glue:

  • IAM Permissions: Ensure your Glue Job's IAM role has:
    • S3 read/write/list permissions for the Delta Table's S3 path
    • Glue Catalog permissions (glue:CreateTable, glue:UpdateTable, glue:GetTable) if using Catalog tables
  • Version Compatibility: Always align Delta Lake version with your Glue Spark version (check Delta's compatibility matrix for specifics)
  • Schema Evolution: If your transformed DF has new columns, add .option("mergeSchema", "true") to your write operation to automatically merge schemas

内容的提问来源于stack exchange,提问作者Harish J

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 08:57:48