如何基于条件向Delta Table插入行(PySpark实现)
PySpark实现Delta Table条件插入逻辑
现有Delta Table数据
| id | vehicle | production | asIs | EU | EU_variant | status |
|---|---|---|---|---|---|---|
| 1 | A3345 | PQ1298 | FV1 | FV1_variant | OK | |
| 2 | A3346 | B3346 | PQ1287 | FV2 | FV2_variant | NOT_OK |
| 3 | A3346 | B3346 | PQ1207 | FV2 | FV2_variant | NOT_OK |
| 4 | A3347 | QP9 | QP9_variant | OK | ||
| 5 | A3347 | QP9 | QP9_variant | NOT_OK | ||
| 6 | A3347 | QP3 | QP3_variant | OK | ||
| 7 | A3348 | MP6553 | YR34 | YR34_variant | NOT_OK | |
| 8 | A3348 | MP6554 | YR35 | YR35_variant | NOT_OK | |
| 9 | A3348 | MP6554 | YR35 | YR35_variant | NOT_OK |
核心插入规则
- 待插入行
vehicle非空时:检查同vehicle的现有记录中是否存在status=NOT_OK的条目,存在则不插入,否则插入 - 待插入行
vehicle为空时:检查同production的现有记录中是否存在status=NOT_OK的条目,存在则不插入,否则插入
实现步骤与代码
1. 初始化Spark环境与读取Delta表
from pyspark.sql import SparkSession from pyspark.sql.functions import expr, max from delta.tables import DeltaTable # 初始化带Delta Lake支持的Spark Session spark = SparkSession.builder \ .appName("DeltaConditionalInsert") \ .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.forPath(spark, "/path/to/your/delta/table")
2. 预处理目标表:标记分组状态
先对现有表按vehicle(非空时)或production(vehicle空时)分组,统计每个分组是否存在NOT_OK记录:
# 生成分组状态汇总表:group_key为vehicle或production,has_not_ok标记是否存在NOT_OK target_summary = delta_table.toDF().groupBy( expr("CASE WHEN vehicle IS NOT NULL THEN vehicle ELSE production END AS group_key") ).agg( max(expr("CASE WHEN status = 'NOT_OK' THEN 1 ELSE 0 END")).alias("has_not_ok") )
3. 筛选符合插入条件的待插入行
将待插入数据与分组汇总表关联,筛选出允许插入的行:
# 构造待插入数据(示例包含题目中的3条测试用例) insert_data = spark.createDataFrame([ (1, "A3345", None, "PQ1298", "FV1", "FV1_variant", "OK"), (2, "A3346", "B3346", "PQ1287", "FV2", "FV2_variant", "NOT_OK"), (9, None, "A3348", "MP6554", "YR35", "YR35_variant", "NOT_OK") ], schema="id int, vehicle string, production string, asIs string, EU string, EU_variant string, status string") # 筛选可插入的行:分组不存在 或 分组无NOT_OK记录 insert_candidates = insert_data.alias("source") \ .join( target_summary.alias("summary"), expr("CASE WHEN source.vehicle IS NOT NULL THEN source.vehicle ELSE source.production END = summary.group_key"), how="left_outer" ) \ .filter("summary.group_key IS NULL OR summary.has_not_ok = 0") \ .select("source.*")
4. 执行插入
将筛选后的行写入Delta表:
# 以append模式插入符合条件的行 insert_candidates.write.mode("append").format("delta").save("/path/to/your/delta/table")
逻辑验证
对应题目中的示例:
- 示例1:A3345分组的
has_not_ok=0(仅存在OK记录),符合插入条件,执行插入 - 示例2:A3346分组的
has_not_ok=1(存在NOT_OK记录),被过滤,不插入 - 示例3:A3348分组的
has_not_ok=1(存在NOT_OK记录),被过滤,不插入
内容的提问来源于stack exchange,提问作者karthik
相关产品推荐
相关产品推荐

