SQL Server递归存储过程转PySpark后数据计数不符排查
SQL Server递归CTE转PySpark后数据计数不一致问题排查
原SQL Server递归查询代码
IF (OBJECT_ID('TEMPDB..#SUPPLYPARTS') IS NOT NULL) BEGIN DROP TABLE #SUPPLYPARTS; END; WITH A (TYPE, QUANTITY, ID, PARENTMSF, TOPPARENT) AS ( SELECT DISTINCT CAST(CHILDMSF AS VARCHAR(MAX)) AS CHILDMSF, PARENTMSF, PARENTMSF AS TOPPARENT, CAST(CONCAT(PARENTMSF, '-->', CHILDMSF) AS VARCHAR(MAX)) AS TRANSITION, QUANTITY, TYPE FROM #BO BO(NOLOCK) WHERE TYPE IN (0, 1) UNION ALL SELECT CAST(B.CHILDMSF AS VARCHAR(MAX)), B.PARENTMSF, A.TOPPARENT, CAST(CONCAT(A.TRANSITION, '-->', B.CHILDMSF) AS VARCHAR(MAX)) AS TRANSITION, B.QUANTITY, B.TYPE FROM A JOIN #BO B (NOLOCK) ON B.PARENTMSF = A.CHILDMSF WHERE A.TYPE IN (0, 1) ) SELECT DISTINCT * INTO #SUPPLYPARTS FROM A WHERE TYPE IN (0, 1) ORDER BY CHILDMSF OPTION(MAXRECURSION 0);
当前PySpark实现代码
PARENT = root_df.select( \ col("CHILDMSF").cast("string"), \ col("ParentMsf").alias("PARENTMSF"), \ col("ParentMsf").alias("TOPPARENT"), \ col("TYPE"), \ col("QUANTITY"))\ .withColumn("TRANSITION", concat(col("ParentMsf"), F.lit(' --> '), col("ChildMsf")).cast("string")) \ .filter(col("TYPE").isin([0, 1])).distinct().alias("A") BO= root_df.select("CHILDMSF","PARENTMSF",col("PARENTMSF").alias("TOPPARENT"),"TYPE","ISSPARABLE","QUANTITY").alias("B") inter1 = PARENT.join(BO, col("P.CHILDMSF") == col("T.ParentMsf"), "inner") \ .select( col("B.CHILDMSF"), col("B.PARENTMSF"), col("B.TOPPARENT"), col("B.TYPE"), col("B.QUANTITY"), col("A.TRANSITION"), ).withColumn("TRANSITION", concat(col("A.TRANSITION"), F.lit(' --> '), col("B.CHILDMSF")).cast("string")).drop("A.BOMPATH").filter(col("B.TYPE").isin([0, 1])) BO1=PARENT.unionAll(inter1).orderBy("CHILDMSF").distinct() BO_FINAL=BO1.filter(col("TYPE").isin([0, 1])) BO_FINAL.count()
核心差异与问题点排查
1. 递归逻辑完全缺失
原SQL使用递归CTE,会持续迭代关联子节点,直到没有更深层级的关系为止(MAXRECURSION 0允许无限递归)。而当前PySpark代码仅做了单次父-子节点关联,只处理了两层数据,完全遗漏了多层级递归的逻辑,这是计数不一致的核心原因。
2. 关联条件别名错误
代码中PARENT.alias("A")、BO.alias("B"),但join条件却写了col("P.CHILDMSF") == col("T.ParentMsf"),这里的别名P、T完全不匹配,正确条件应为:
PARENT.join(BO, col("A.CHILDMSF") == col("B.PARENTMSF"), "inner")
3. TOPPARENT字段赋值错误
原SQL递归中,TOPPARENT始终继承自最顶层的父节点(即初始节点的PARENTMSF),但当前PySpark代码中:
- 初始
PARENT表的TOPPARENT赋值正确 inter1中错误地取了B.TOPPARENT(当前子节点的PARENTMSF),正确逻辑应该取A.TOPPARENT,保持顶层父节点不变。
4. 冗余字段与无效操作
BO表中选取了ISSPARABLE字段,原SQL完全未使用该字段,属于冗余操作drop("A.BOMPATH")语句中,A.BOMPATH字段从未被定义,属于无效操作
5. 去重时机差异
原SQL是在所有递归结果生成后执行SELECT DISTINCT去重,而当前PySpark代码分步对初始节点、合并后数据去重,这种逻辑可能导致与原SQL的去重结果不一致,建议调整为最终递归完成后再统一去重。
修正后的PySpark递归实现思路
from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 定义结果Schema,与原SQL字段对齐 result_schema = StructType([ StructField("CHILDMSF", StringType()), StructField("PARENTMSF", StringType()), StructField("TOPPARENT", StringType()), StructField("TRANSITION", StringType()), StructField("QUANTITY", IntegerType()), StructField("TYPE", IntegerType()) ]) # 初始节点:对应CTE的锚点查询 current_df = root_df.select( F.col("CHILDMSF").cast(StringType()), F.col("ParentMsf").alias("PARENTMSF"), F.col("ParentMsf").alias("TOPPARENT"), F.concat(F.col("ParentMsf"), F.lit(' --> '), F.col("CHILDMSF")).cast(StringType()).alias("TRANSITION"), F.col("QUANTITY"), F.col("TYPE") ).filter(F.col("TYPE").isin([0, 1])).distinct() # 存储最终结果 final_df = current_df # 递归循环:直到没有新数据加入 while True: # 关联下一层子节点 next_level_df = current_df.join( root_df.alias("B"), current_df["CHILDMSF"] == F.col("B.ParentMsf"), "inner" ).select( F.col("B.CHILDMSF").cast(StringType()), F.col("B.ParentMsf").alias("PARENTMSF"), current_df["TOPPARENT"], # 继承顶层父节点 F.concat(current_df["TRANSITION"], F.lit(' --> '), F.col("B.CHILDMSF")).cast(StringType()).alias("TRANSITION"), F.col("B.QUANTITY"), F.col("B.TYPE") ).filter(F.col("B.TYPE").isin([0, 1])).distinct() # 检查是否有新数据 new_rows = next_level_df.subtract(final_df) if new_rows.count() == 0: break # 合并新数据到结果集 final_df = final_df.union(new_rows) current_df = new_rows # 最终去重并过滤 BO_FINAL = final_df.distinct().filter(F.col("TYPE").isin([0, 1])).orderBy("CHILDMSF") BO_FINAL.count()
内容的提问来源于stack exchange,提问作者Data writer
相关产品推荐
相关产品推荐

