使用AWS Glue for循环同步S3表到Redshift如何借助staging表避免数据重复
适配批量循环场景的staging表方案实现
你原来的脚本重复运行重复写入的核心原因是默认write_dynamic_frame.from_jdbc_conf的写入模式是追加(APPEND),我们可以把staging表逻辑套入循环,每个表独立处理自己的临时 staging 表,具体实现如下:
前置准备
- 所有迁移表需要提前在Redshift目标schema中创建完成,且定义好主键(增量去重场景需要,全量覆盖场景可跳过)
- 确保你的Glue角色有Redshift的SQL执行权限,以及对应schema的读写权限
核心修改点
- 循环内每个表生成唯一的
transformation_ctx,避免Glue Job Bookmark生效时出现上下文冲突 - 每个表先把全量/增量数据写入Redshift会话级临时staging表
- 通过Redshift事务完成「去重-写入」的原子操作,避免中间状态出错
- 支持两种常用场景:全量覆盖迁移、增量主键去重迁移
修改后的完整脚本
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
相关产品推荐
相关产品推荐

