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

如何用PySpark为Spark DataFrame按规则添加As_Is_Flag列

PySpark实现添加As_Is_Flag列的解决方案

原始数据

Vehicle Production Model  ECU         As-Is   modelData     Coding  Flag
22223   22233        C24      BR-11   A42222  A2223          A3442   0
22223   22233        C24      BR-11   A3442   A2223          A3442   0
22223   22233        C24      BR-11   A3443   A2223          A3443   1
22223   22233        C24      BR-12   A3442   A2223          A3443   1
22224   22234        C24      BR-12   A3442   A2223          A3443   1
22224   22234        C24      BR-13   A3443   A2223          A3446   0
22224   22234        C24      BR-13   A3447   A2223          A3444   0
22224   22234        C24      BR-13   A3448   A2225          A3443   0

需求规则

  • 按Vehicle和ECU分组处理
  • 每个分组内提取:
    • 所有As-Is值
    • Flag=0且modelData、Coding非空的唯一(modelData, Coding)组合对应的As-Is值集合
  • 对每条数据的As-Is值判断是否存在于上述集合中,存在则标记As_Is_Flag=1,否则为0

预期输出

Vehicle Production Model  ECU     As-Is      modelData     Coding  Flag  As-Is Flag
22223   22233        C24      BR-11  A42222  A2223         A3442    0       0
22223   22233        C24      BR-11  A3442    A2223        A3442    0       1
22223   22233        C24      BR-11  A3443    A2223        A3443    1       0
22223   22233        C24      BR-12   A3442   A2223        A3443    1       1
22224   22234        C24      BR-12   A3442   A2223        A3443    1       0
22224   22234        C24      BR-13   A3443   A2223        A3446    0       1
22224   22234        C24      BR-13   A3447   A2223        A3444    0       0
22224   22234        C24      BR-13   A3448   A2225        A3443    0       0

PySpark实现代码

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

# 初始化SparkSession
spark = SparkSession.builder.appName("AsIsFlagCalculation").getOrCreate()

# 构建原始DataFrame
data = [
    ("22223", "22233", "C24", "BR-11", "A42222", "A2223", "A3442", "0"),
    ("22223", "22233", "C24", "BR-11", "A3442", "A2223", "A3442", "0"),
    ("22223", "22233", "C24", "BR-11", "A3443", "A2223", "A3443", "1"),
    ("22223", "22233", "C24", "BR-12", "A3442", "A2223", "A3443", "1"),
    ("22224", "22234", "C24", "BR-12", "A3442", "A2223", "A3443", "1"),
    ("22224", "22234", "C24", "BR-13", "A3443", "A2223", "A3446", "0"),
    ("22224", "22234", "C24", "BR-13", "A3447", "A2223", "A3444", "0"),
    ("22224", "22234", "C24", "BR-13", "A3448", "A2225", "A3443", "0")
]

columns = ["Vehicle", "Production", "Model", "ECU", "As-Is", "modelData", "Coding", "Flag"]
df = spark.createDataFrame(data, columns)

# 1. 生成每个Vehicle-ECU分组下的有效As-Is集合
valid_as_is_df = df.filter(
    (F.col("Flag") == "0") &
    (F.col("modelData").isNotNull()) & (F.col("modelData") != "") &
    (F.col("Coding").isNotNull()) & (F.col("Coding") != "")
).groupBy("Vehicle", "ECU").agg(
    F.collect_set("As-Is").alias("valid_as_is_set")
)

# 2. 关联原始数据并计算As-Is Flag
result_df = df.join(valid_as_is_df, on=["Vehicle", "ECU"], how="left") \
    .withColumn(
        "As-Is Flag",
        F.when(F.col("As-Is").isin(F.col("valid_as_is_set")), 1).otherwise(0)
    ) \
    .drop("valid_as_is_set")

# 查看结果
result_df.show(truncate=False)

代码说明

  1. 数据初始化:创建SparkSession并加载原始数据为DataFrame。
  2. 生成有效As-Is集合:筛选Flag=0且modelData、Coding非空的记录,按Vehicle和ECU分组后,用collect_set收集该分组内符合条件的唯一As-Is值,形成判断用的集合。
  3. 关联与标记:将原始DataFrame与有效集合DataFrame按Vehicle和ECU关联,通过when-otherwise判断每条记录的As-Is是否在对应分组的有效集合中,生成目标列后删除临时集合列。

内容的提问来源于stack exchange,提问作者karthik kk

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 04:45:41