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

PySpark数据入库校验:去空去重及registrationid的Upsert实现问询

PySpark实现精修层数据Upsert(更新插入)方案

前置清洗步骤(基础操作补充)

先完成基础清洗确保新数据符合入库要求:

  • 移除空值:优先保证主键registrationid非空,通过dropna指定关键列过滤:
    df_clean = raw_df.dropna(subset=["registrationid"])
    
  • 移除原始数据重复记录:若原始层存在重复registrationid,建议保留最新版本(假设dischargedate为更新时间):
    df_dedup = df_clean.orderBy("dischargedate", ascending=False).dropDuplicates(["registrationid"])
    

核心:精修层Upsert实现

方案1:基于Delta Lake的原生Merge操作(推荐)

Delta Lake支持ACID事务和原生Upsert,是大数据场景下的最优解。

代码实现:

from pyspark.sql import SparkSession
from delta.tables import DeltaTable

# 初始化SparkSession(需配置Delta相关参数)
spark = SparkSession.builder \
    .appName("PatientDataUpsert") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

# 读取精修层Delta表(支持表名或路径)
delta_table = DeltaTable.forName(spark, "refined_layer.patient_data")
# 或通过路径读取:DeltaTable.forPath(spark, "/warehouse/refined_layer/patient_data")

# 执行Merge(Upsert)
delta_table.alias("target") \
    .merge(
        df_dedup.alias("source"),
        "target.registrationid = source.registrationid"
    ) \
    .whenMatchedUpdateAll()  # 匹配到则全量更新目标表记录
    .whenNotMatchedInsertAll()  # 未匹配到则插入新记录
    .execute()

自定义更新/插入逻辑:

若无需全量更新,可指定具体列:

delta_table.alias("target") \
    .merge(
        df_dedup.alias("source"),
        "target.registrationid = source.registrationid"
    ) \
    .whenMatchedUpdate(
        set={
            "dischargedate": "source.dischargedate",
            "status_description": "source.status_description",
            "stage": "source.stage"
        }
    ) \
    .whenNotMatchedInsert(
        values={
            "registrationid": "source.registrationid",
            "dischargedate": "source.dischargedate",
            "status_description": "source.status_description",
            "dischargetype": "source.dischargetype",
            "stage": "source.stage"
        }
    ) \
    .execute()

方案2:普通Spark表的Upsert(无Delta场景)

若无法使用Delta Lake,可通过"过滤旧数据+合并新数据+覆盖写入"实现:

代码实现:

# 读取精修层现有数据
refined_df = spark.read.table("refined_layer.patient_data")

# 获取新数据中所有的registrationid
new_reg_ids = df_dedup.select("registrationid").distinct()

# 过滤掉精修层中与新数据重复的记录
filtered_refined_df = refined_df.join(new_reg_ids, on="registrationid", how="left_anti")

# 合并过滤后的旧数据与清洗后的新数据
final_df = filtered_refined_df.unionByName(df_dedup)

# 覆盖写入精修层(若为分区表,建议按分区处理提升效率)
final_df.write.mode("overwrite").saveAsTable("refined_layer.patient_data")

注意事项:

  • 该方式无事务保障,写入过程中出错会导致数据丢失;
  • 数据量较大时全量覆盖效率极低,建议结合分区或分桶优化。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 05:44:50