使用Pyspark与Glue Job实现Redshift到S3增量数据加载方案咨询
关于Redshift到S3增量加载方案的解答
Glue Job Bookmarks是否适用
Glue Job Bookmarks仅对Glue原生数据源连接器的读取操作生效,你当前使用的是原生Spark JDBC接口读取Redshift,没有通过Glue内置的Redshift连接器,所以即使开启Job Bookmarks,也无法追踪你的数据读取进度,无法实现自动增量加载。
如果要适配Glue Job Bookmarks能力,可以将读取逻辑替换为Glue Redshift连接器的读取方式,示例代码如下:
from awsglue.context import GlueContext glueContext = GlueContext(spark.sparkContext) dynamic_frame = glueContext.create_dynamic_frame.from_options( connection_type="redshift", connection_options={ "url": url, "dbtable": f"{schema}.{tablename}", "user": Properties["user"], "password": Properties["password"], "redshiftTmpDir": "s3://你的临时存储桶路径/" }, transformation_ctx="redshift_read" ) df = dynamic_frame.toDF()
配置transformation_ctx参数并开启Job Bookmarks后,Glue会自动追踪该读取操作的进度,首次运行加载全量数据,后续重跑任务时只会拉取新增的记录。
其他可行的增量实现方案
如果不想依赖Glue特性,也可以选择以下通用方案:
- 水印列增量方案(最通用)
前提:Redshift目标表存在可标识数据新旧的字段,比如自增主键ID、创建时间create_time、更新时间update_time
实现逻辑:- 每次任务启动时,先读取S3中已写入的历史数据,取出最大的增量标识值(比如
max_update_time) - 读取Redshift时拼接查询条件,首次运行加载全量,非首次运行只查询
update_time > ${max_update_time}的新增/变更数据 - 将过滤后的增量数据写入S3对应路径
优势:完全可控,不需要调整现有核心读取逻辑,适配所有Spark运行环境。
- 每次任务启动时,先读取S3中已写入的历史数据,取出最大的增量标识值(比如
- Redshift CDC方案
如果需要同步数据的更新、删除操作,可以开启Redshift CDC功能,通过读取解析CDC日志获取全量数据变更记录,适合对数据一致性要求高的场景。
注意事项
- 若使用时间列作为增量标识,建议每次查询的时间窗口结束值设为当前时间往前推3~5分钟,避免Redshift事务提交延迟导致的增量数据丢失
- 定期做全量重刷校验,避免增量逻辑异常导致的数据不一致
- 基于水印列的方案无法捕获表中的硬删除记录,需要配合CDC或者定期全量比对补全数据
内容的提问来源于stack exchange,提问作者whatsinthename
相关产品推荐
相关产品推荐

