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

如何基于条件向Delta Table插入行(PySpark实现)

PySpark实现Delta Table条件插入逻辑

现有Delta Table数据

idvehicleproductionasIsEUEU_variantstatus
1A3345PQ1298FV1FV1_variantOK
2A3346B3346PQ1287FV2FV2_variantNOT_OK
3A3346B3346PQ1207FV2FV2_variantNOT_OK
4A3347QP9QP9_variantOK
5A3347QP9QP9_variantNOT_OK
6A3347QP3QP3_variantOK
7A3348MP6553YR34YR34_variantNOT_OK
8A3348MP6554YR35YR35_variantNOT_OK
9A3348MP6554YR35YR35_variantNOT_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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 10:15:46