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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 10:32:06