DataFrame代际黑发信息传递问题及代码优化需求
系谱树DataFrame黑发信息向上传递优化方案
需求说明
现有呈系谱树结构的DataFrame,包含列:Generation(世代,0为最老世代,最大值为最年轻世代)、Child_name(子代姓名)、child_hair_color(子代发色)、Parent_name(父代姓名)、parent_hair_color(父代发色)、parent_eye_color(父代眼色)。需实现黑发信息从子代向上传递至父代、祖母代等,仅当父代眼色为Hazel时停止传递。
输入数据集
Schema定义
from pyspark.sql.types import StructType, StructField, IntegerType, StringType schema = StructType([ StructField('Generation', IntegerType()), StructField('Child_name', StringType()), StructField('child_hair_color', StringType()), StructField('Parent_name', StringType()), StructField('parent_hair_color', StringType()), StructField('parent_eye_color', StringType()) ])
初始数据行
from pyspark.sql import Row rows = [ Row(0, "Susan", "Brown", None, None, None ), Row(1, "Mary", "Blond", "Susan", "Brown", "Green"), Row(2, "Alexandra", "Grey", "Mary", "Blond", "Blue"), Row(1, "Robin", "Blond", "Susan", "Brown", "Hazel"), Row(2, "Rachel", "Brown", "Robin", "Blond", "Blue"), Row(2, "Mika", "Grey", "Robin", "Blond", "Blue"), Row(3, "Patricia", "Blond", "Mika", "Grey", "Green"), Row(4, "Molly", "Black", "Patricia", "Blond", "Green") ]
用户原代码(无效果)
以下代码运行后输入输出DataFrame无差异,且无法高效处理30-70万行大数据集:
selected_column = "Generation" # Use function agg() to find min and max from "Generation" column min_max_values = df.agg(F.min(selected_column).alias("min_value"), F.max(selected_column).alias("max_value")).first() max_value = min_max_values["max_value"] # Loop through generations for current_level in range(max_value , 0, -1): df = df.withColumn("parent_hair_color", when((df.child_hair_color == "Black") & (df.Generation == current_level) & (df.parent_eye_color != "Hazel"), ("Black")) .otherwise(F.col("parent_hair_color"))) df1 = df.filter(df.parent_hair_color == "Black").select(F.col("Parent_name").alias("Parent")) df = df.join(df1, df.Child_name == df1.Parent, "left").select(df["*"], df1["Parent"]) df = df.withColumn("child_hair_color", when(df.Child_name == df.Parent, ("Black")) .otherwise(F.col("child_hair_color"))) df = df.drop("Parent")
期望输出数据
rows = [ Row(0, "Susan", "Brown", None, None, None ), Row(1, "Mary", "Blond", "Susan", "Brown", "Green"), Row(2, "Alexandra", "Grey", "Mary", "Blond", "Blue"), Row(1, "Robin", "Black", "Susan", "Brown", "Hazel"), Row(2, "Rachel", "Brown", "Robin", "Blond", "Blue"), Row(2, "Mika", "Black", "Robin", "Black", "Blue"), Row(3, "Patricia", "Black", "Mika", "Black", "Green"), Row(4, "Molly", "Black", "Patricia", "Black", "Green") ]
问题分析与优化方案
原代码问题
- 逻辑错位:修改
parent_hair_color时,仅修改当前行的父代发色字段,未关联到父代自身的行进行更新;后续join逻辑也未正确定位父代行完成发色传递。 - 效率低下:循环中反复对全量DataFrame操作,几十万行数据会产生大量重复计算,性能极差。
- 传递中断逻辑缺失:未处理父代被染黑后继续向上传递的连锁反应,仅单步传递,未完成完整的系谱传递。
优化方案(递归CTE实现)
利用Spark递归CTE高效处理树状结构传递,避免循环损耗,适配大数据量场景:
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 初始化SparkSession spark = SparkSession.builder.appName("HairColorPropagation").getOrCreate() # 构建初始DataFrame df = spark.createDataFrame(rows, schema) # 创建临时视图供SQL使用 df.createOrReplaceTempView("original_df") # 递归CTE实现发色传递 result_df = spark.sql(""" WITH RECURSIVE hair_propagation AS ( -- 基础层:筛选初始黑发节点 SELECT Generation, Child_name, child_hair_color, Parent_name, parent_hair_color, parent_eye_color, 1 AS propagate_flag -- 标记需要传递的节点 FROM original_df WHERE child_hair_color = 'Black' UNION ALL -- 递归层:向上传递黑发标记,遇到Hazel眼色停止 SELECT p.Generation, p.Child_name, 'Black' AS child_hair_color, p.Parent_name, 'Black' AS parent_hair_color, p.parent_eye_color, CASE WHEN p.parent_eye_color != 'Hazel' THEN 1 ELSE 0 END AS propagate_flag FROM hair_propagation hp JOIN original_df p ON hp.Parent_name = p.Child_name WHERE hp.propagate_flag = 1 AND p.parent_eye_color != 'Hazel' ) -- 合并原始数据与递归结果,取最终发色 SELECT o.Generation, o.Child_name, COALESCE(hp.child_hair_color, o.child_hair_color) AS child_hair_color, o.Parent_name, COALESCE(hp.parent_hair_color, o.parent_hair_color) AS parent_hair_color, o.parent_eye_color FROM original_df o LEFT JOIN ( SELECT Child_name, MAX(child_hair_color) AS child_hair_color, MAX(parent_hair_color) AS parent_hair_color FROM hair_propagation GROUP BY Child_name ) hp ON o.Child_name = hp.Child_name ORDER BY Generation, Child_name """) # 查看结果 result_df.show()
方案优势
- 高效处理大数据:递归CTE是Spark原生优化的树状结构处理方式,避免循环中的重复计算,适配30-70万行数据量。
- 逻辑精准:从初始黑发节点逐层向上传递,遇到Hazel眼色父代立即停止,完全符合需求。
- 结果完整:通过合并原始数据与递归结果,确保所有节点的发色都被正确更新。
内容的提问来源于stack exchange,提问作者user25190795
相关产品推荐
相关产品推荐

