通过Glue ETL作业写入RDS PostgreSQL表时如何过滤异常记录
Glue ETL异常记录过滤及分流方案
核心思路
提前基于目标表元数据做记录级校验,拆分合法、异常数据集分别写入,从根源避免写入阶段触发类型/长度不匹配异常导致作业终止。
具体实现步骤
- 拉取目标表元数据约束:从Glue Catalog获取目标表的字段类型、字符串最大长度等校验规则,也可直接连接RDS PostgreSQL查询
information_schema.columns表获取更精准的表结构定义。 - 编写校验逻辑自定义UDF:对每条记录逐字段校验,返回两个附加字段:
is_valid(布尔值,标记是否符合规则)、error_msg(字符串,标注异常原因,比如“字段a长度超出定义的20位”、“字段b类型应为int实际为字符串”)。 - 拆分数据集:根据
is_valid字段把转换后的数据集拆分为合法记录集、异常记录集两个独立数据集。 - 分别写入目标位置:合法记录写入RDS PostgreSQL目标表,异常记录写入S3路径或者专门的异常归档表,写入时配置容错参数避免偶发异常终止作业。
代码实现示例
import sys from pyspark.sql.functions import udf, struct, lit from pyspark.sql.types import BooleanType, StringType, StructType, StructField from awsglue.utils import getResolvedOptions from awsglue.context import GlueContext from awsglue.job import Job from awsglue.dynamicframe import DynamicFrame from pyspark.context import SparkContext # 初始化参数 args = getResolvedOptions(sys.argv, ['JOB_NAME', 'SRC_DB', 'SRC_TABLE', 'TGT_DB', 'TGT_TABLE', 'ERROR_S3_PATH']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) # 1. 读取源数据并做基础转换 DataSource0 = glueContext.create_dynamic_frame.from_catalog(database = args['SRC_DB'], table_name = args['SRC_TABLE'], transformation_ctx = "DataSource0") Transform0 = sparkSqlQuery(glueContext, query = SqlQuery0, mapping = {"sparkDataSource": DataSource0}, transformation_ctx = "Transform0") df = Transform0.toDF() # 2. 获取目标表元数据规则(示例从Glue Catalog获取,也可替换为查询RDS元数据的逻辑) tgt_table = glueContext.get_table(database_name=args['TGT_DB'], name=args['TGT_TABLE']) tgt_columns = {col['Name']: {'type': col['Type'], 'max_length': int(col['Type'].split('(')[1].split(')')[0]) if col['Type'].startswith('varchar') else None} for col in tgt_table['Table']['StorageDescriptor']['Columns']} # 3. 定义校验UDF def validate_record(row): error_list = [] for col_name, rule in tgt_columns.items(): val = row[col_name] if val is None: continue # 校验字符串长度 if rule['max_length'] is not None: if len(str(val)) > rule['max_length']: error_list.append(f"字段{col_name}长度超出上限{rule['max_length']},实际长度{len(str(val))}") # 校验数据类型(示例仅校验int、timestamp类型,可按需扩展) if rule['type'] == 'int': if not isinstance(val, int): error_list.append(f"字段{col_name}类型应为int,实际为{type(val).__name__}") elif rule['type'] == 'timestamp': if not hasattr(val, 'minute'): error_list.append(f"字段{col_name}类型应为timestamp,实际为{type(val).__name__}") if len(error_list) == 0: return (True, "") else: return (False, ";".join(error_list)) validate_udf = udf(validate_record, StructType([ StructField("is_valid", BooleanType(), nullable=False), StructField("error_msg", StringType(), nullable=True) ])) # 4. 执行校验并拆分数据集 df_with_check = df.withColumn("check_result", validate_udf(struct([df[x] for x in df.columns]))) valid_df = df_with_check.filter(df_with_check.check_result.is_valid == True).select(df.columns) error_df = df_with_check.filter(df_with_check.check_result.is_valid == False).withColumn("error_msg", df_with_check.check_result.error_msg) # 5. 写入合法数据到RDS,配置容错参数 valid_dyf = DynamicFrame.fromDF(valid_df, glueContext, "valid_dyf") DataSink0 = glueContext.write_dynamic_frame.from_catalog( frame = valid_dyf, database = args['TGT_DB'], table_name = args['TGT_TABLE'], transformation_ctx = "DataSink0", additional_options = { "rewriteBatchedStatements": "true", "batchsize": "1000", "error_threshold": "0" # 合法数据集仍可配置容错阈值,0表示不允许任何错误 } ) # 6. 写入异常数据到S3(也可配置写入其他Glue表) error_dyf = DynamicFrame.fromDF(error_df, glueContext, "error_dyf") glueContext.write_dynamic_frame.from_options( frame = error_dyf, connection_type = "s3", connection_options = {"path": args['ERROR_S3_PATH']}, format = "json", # 也可选择parquet、csv格式 transformation_ctx = "ErrorSink" ) job.commit()
可选优化点
- 可调整
error_threshold参数设置允许的最大错误记录比例,超过该比例才终止作业,避免极少量漏校验的异常数据导致作业中断。 - 异常数据可按日期分区存储,方便后续回溯排查。
- 可添加监控告警,当异常记录占比超过预设阈值时自动触发通知给数据源团队。
内容的提问来源于stack exchange,提问作者Vishal Mishra
相关产品推荐
相关产品推荐

