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
相关产品推荐
相关产品推荐

