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

如何实现带多业务键的多源目标视图批量数据验证脚本优化

Spark 数据验证统一脚本解决方案

核心实现思路

把重复的行计数校验、数据差异对比逻辑封装成可复用函数,通过配置化的视图映射列表传入源视图、目标视图及对应业务键,实现批量自动化验证,最终输出统一的校验结果。

完整Python Spark脚本

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, concat_ws

def init_spark():
    return SparkSession.builder.appName("DataValidation").getOrCreate()

def check_row_count(spark, src_table, tgt_table):
    """校验源表与目标表的行数差异"""
    src_count = spark.sql(f"SELECT COUNT(*) AS row_count FROM {src_table}").collect()[0]["row_count"]
    tgt_count = spark.sql(f"SELECT COUNT(*) AS row_count FROM {tgt_table}").collect()[0]["row_count"]
    return {
        "table_pair": f"{src_table} <-> {tgt_table}",
        "src_rows": src_count,
        "tgt_rows": tgt_count,
        "row_diff": abs(src_count - tgt_count)
    }

def compare_data_diff(spark, src_table, tgt_table, business_keys):
    """对比源表与目标表的A-B、B-A差异,按业务键排序结果"""
    # 获取源表和目标表的公共字段(需确保字段名一致,不一致需提前做映射)
    src_cols = spark.sql(f"DESCRIBE {src_table}").select("col_name").rdd.flatMap(lambda x: x).collect()
    tgt_cols = spark.sql(f"DESCRIBE {tgt_table}").select("col_name").rdd.flatMap(lambda x: x).collect()
    common_cols = list(set(src_cols) & set(tgt_cols))
    
    # 生成业务键拼接字段用于关联匹配
    key_concat = concat_ws("|", *[col(k) for k in business_keys])
    
    src_df = spark.table(src_table).withColumn("key_concat", key_concat)
    tgt_df = spark.table(tgt_table).withColumn("key_concat", key_concat)
    
    # A-B:源表存在但目标表缺失的记录
    a_minus_b = src_df.join(tgt_df, on="key_concat", how="left_anti") \
                     .select(*common_cols) \
                     .orderBy(*business_keys)
    
    # B-A:目标表存在但源表缺失的记录
    b_minus_a = tgt_df.join(src_df, on="key_concat", how="left_anti") \
                     .select(*common_cols) \
                     .orderBy(*business_keys)
    
    return {
        "table_pair": f"{src_table} <-> {tgt_table}",
        "a_minus_b_count": a_minus_b.count(),
        "b_minus_a_count": b_minus_a.count(),
        "a_minus_b_records": a_minus_b.collect(),
        "b_minus_a_records": b_minus_a.collect()
    }

def run_all_validations(spark, validation_config):
    """执行所有验证任务并汇总输出结果"""
    all_results = []
    for config in validation_config:
        src_table = config["src_table"]
        tgt_table = config["tgt_table"]
        business_keys = config["business_keys"]
        
        # 执行行数校验
        row_result = check_row_count(spark, src_table, tgt_table)
        # 执行数据差异对比
        diff_result = compare_data_diff(spark, src_table, tgt_table, business_keys)
        
        # 合并结果
        combined_result = {**row_result, **diff_result}
        all_results.append(combined_result)
        
        # 打印当前表对的验证详情
        print(f"\n=== 验证结果:{src_table} <-> {tgt_table} ===")
        print(f"源表行数: {row_result['src_rows']} | 目标表行数: {row_result['tgt_rows']} | 行数差异: {row_result['row_diff']}")
        print(f"A-B差异行数: {diff_result['a_minus_b_count']} | B-A差异行数: {diff_result['b_minus_a_count']}")
        
        if diff_result['a_minus_b_count'] > 0:
            print("\nA-B差异记录(按业务键排序):")
            for rec in diff_result['a_minus_b_records']:
                print(rec)
        if diff_result['b_minus_a_count'] > 0:
            print("\nB-A差异记录(按业务键排序):")
            for rec in diff_result['b_minus_a_records']:
                print(rec)
        
        if row_result['row_diff'] == 0 and diff_result['a_minus_b_count'] == 0 and diff_result['b_minus_a_count'] == 0:
            print("✅ 数据无差异")
    
    return all_results

if __name__ == "__main__":
    spark = init_spark()
    
    # --------------------------
    # 配置区:按需修改视图对及业务键
    # --------------------------
    validation_config = [
        {
            "src_table": "TEST_SCH.VIEWNAME1",
            "tgt_table": "schema.deltaview1",
            "business_keys": ["user_id", "order_date"]
        },
        {
            "src_table": "TEST_SCH.VIEWNAME2",
            "tgt_table": "schema.deltaview2",
            "business_keys": ["product_id"]
        }
        # 可添加更多视图对配置
    ]
    
    # 启动验证
    run_all_validations(spark, validation_config)
    
    spark.stop()

使用说明

  1. 配置修改:在validation_config列表中添加需要验证的视图对,每个配置项包含:
    • src_table:SQL源视图全名称(如TEST_SCH.VIEWNAME)
    • tgt_table:Delta目标视图全名称(如schema.deltaview)
    • business_keys:当前视图的业务主键列表(支持多个字段)
  2. 字段一致性:确保源视图与目标视图的字段名一致,若存在字段名差异,需提前在脚本中添加字段映射逻辑
  3. 运行方式:直接提交Spark作业执行,脚本会逐表输出验证结果,包括行数差异、A/B方向的差异记录数及具体差异数据(按业务键排序)
  4. 结果判定:若行数差异、A-B差异数、B-A差异数均为0,则判定数据无差异

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 06:54:57