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

使用AWS Glue for循环同步S3表到Redshift如何借助staging表避免数据重复

适配批量循环场景的staging表方案实现

你原来的脚本重复运行重复写入的核心原因是默认write_dynamic_frame.from_jdbc_conf的写入模式是追加(APPEND),我们可以把staging表逻辑套入循环,每个表独立处理自己的临时 staging 表,具体实现如下:


前置准备

  • 所有迁移表需要提前在Redshift目标schema中创建完成,且定义好主键(增量去重场景需要,全量覆盖场景可跳过)
  • 确保你的Glue角色有Redshift的SQL执行权限,以及对应schema的读写权限

核心修改点

  1. 循环内每个表生成唯一的transformation_ctx,避免Glue Job Bookmark生效时出现上下文冲突
  2. 每个表先把全量/增量数据写入Redshift会话级临时staging表
  3. 通过Redshift事务完成「去重-写入」的原子操作,避免中间状态出错
  4. 支持两种常用场景:全量覆盖迁移、增量主键去重迁移

修改后的完整脚本

import boto3
import sys
from awsglue.context import GlueContext
from pyspark.context import SparkContext
from awsglue.job import Job
from awsglue.utils import getResolvedOptions

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
args = getResolvedOptions(sys.argv, ['TempDir', 'JOB_NAME'])
job.init(args['JOB_NAME'], args)

# 初始化Redshift Data API客户端,用于执行SQL
redshift_data = boto3.client('redshift-data', region_name='us-east-1')
REDSHIFT_CLUSTER_ID = '替换为你的Redshift集群ID'
REDSHIFT_USER = '替换为你的数据库用户名'
REDSHIFT_DB = 'db1'
TARGET_SCHEMA = 'schema1'

client = boto3.client("glue", region_name="us-east-1")
databaseName = "db1_g"
Tables = client.get_tables(DatabaseName=databaseName)
tableList = Tables["TableList"]

# 可提前维护表名和主键的映射字典,适配多表主键不同的场景
table_pk_map = {
    "表名1": "主键列名1",
    "表名2": "主键列名2"
}

for table in tableList:
    tableName = table["Name"]
    # 1. 读取S3源数据,每个表用唯一的transformation_ctx
    datasource0 = glueContext.create_dynamic_frame.from_catalog(
        database="db1_g", 
        table_name=tableName, 
        transformation_ctx=f"datasource_{tableName}"
    )

    # 2. 写入临时staging表(#开头为会话级临时表,会话结束自动删除,无命名冲突)
    staging_table = f"#staging_{tableName}"
    glueContext.write_dynamic_frame.from_jdbc_conf(
        frame=datasource0,
        catalog_connection="redshift",
        connection_options={
            "dbtable": staging_table,
            "database": REDSHIFT_DB,
        },
        redshift_tmp_dir=args["TempDir"],
        transformation_ctx=f"datasink_staging_{tableName}",
    )

    # --------------------------
    # 场景1:全量覆盖迁移(每次跑清空目标表再写入,适合全量同步场景)
    # --------------------------
    # full_load_sql = f"""
    # BEGIN TRANSACTION;
    # TRUNCATE TABLE {TARGET_SCHEMA}.{tableName};
    # INSERT INTO {TARGET_SCHEMA}.{tableName} SELECT * FROM {staging_table};
    # COMMIT;
    # """
    # redshift_data.execute_statement(
    #     ClusterIdentifier=REDSHIFT_CLUSTER_ID,
    #     DbUser=REDSHIFT_USER,
    #     Database=REDSHIFT_DB,
    #     Sql=full_load_sql
    # )

    # --------------------------
    # 场景2:增量去重迁移(根据主键去重,适合增量同步场景)
    # --------------------------
    pk_col = table_pk_map[tableName]
    upsert_sql = f"""
    BEGIN TRANSACTION;
    DELETE FROM {TARGET_SCHEMA}.{tableName} t
    USING {staging_table} s
    WHERE t.{pk_col} = s.{pk_col};
    INSERT INTO {TARGET_SCHEMA}.{tableName} SELECT * FROM {staging_table};
    COMMIT;
    """
    redshift_data.execute_statement(
        ClusterIdentifier=REDSHIFT_CLUSTER_ID,
        DbUser=REDSHIFT_USER,
        Database=REDSHIFT_DB,
        Sql=upsert_sql
    )

job.commit()

注意事项

  • 会话级临时表仅对当前连接有效,不同表的处理不会互相干扰
  • 如果你不想用Redshift Data API,也可以通过Spark JDBC的方式执行上述SQL语句

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 16:36:04