如何用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)
代码说明
- 数据初始化:创建SparkSession并加载原始数据为DataFrame。
- 生成有效As-Is集合:筛选
Flag=0且modelData、Coding非空的记录,按Vehicle和ECU分组后,用collect_set收集该分组内符合条件的唯一As-Is值,形成判断用的集合。 - 关联与标记:将原始DataFrame与有效集合DataFrame按
Vehicle和ECU关联,通过when-otherwise判断每条记录的As-Is是否在对应分组的有效集合中,生成目标列后删除临时集合列。
内容的提问来源于stack exchange,提问作者karthik kk
相关产品推荐
相关产品推荐

