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

SQL Server与PySpark中NTILE函数输出差异的原因探究

PySpark与SQL Server中NTILE函数结果差异的原因

问题背景

数据集

原始数据集(SQL表中Value为整数类型,Spark DataFrame中显示为浮点型):

+--+-----------+-----+
|Id|Date       |Value|
+---------------------|
|1| 2021-08-11|  4.0 |
|1| 2021-08-12|  1.0 |
|1| 2021-08-13|  4.0 |
|1| 2021-08-14|  2.0 |
|1| 2021-08-15|  4.0 |
|1| 2021-08-16|  1.0 |
|1| 2021-08-17|  2.0 |
|1| 2021-08-18|  0.0 |
|1| 2021-08-19|  0.0 |
|1| 2021-08-20|  2.0 |
|1| 2021-08-21|  2.0 |
|1| 2021-08-22|  4   |
|1| 2021-08-23|  1.0 |
|1| 2021-08-24|  3.0 |
+-+-----------+------+

执行的查询

  • SQL Server查询(注:原文存在拼写错误SELECET,正确应为SELECT):
SELECT ntile(4) over (PARTITION BY Id ORDER BY Value) AS [NTILE(SQL)]
  • PySpark 3.2.1代码:
from pyspark.sql import functions as f
from pyspark.sql.window import Window

w = Window.partitionBy("Id").orderBy("Value")
df = df.withColumn("ntile", f.ntile(4).over(w))

结果差异

  • SQL Server输出:
+--------------+
|NTILE(SQL)    |
---------------+
|             4|
|             2|
|             4|
|             2|
|             4|
|             1|
|             2|
|             1|
|             1|
|             2|
|             3|
|             3|
|             1|
|             3|
+--------------+
  • PySpark输出:
+--------------+
|NTILE(Spark)  |
---------------+
|             4|
|             1|
|             4|
|             2|
|             4|
|             1|
|             2|
|             1|
|             1|
|             2|
|             3|
|             3|
|             2|
|             3|
+--------------+

差异原因分析

1. 排序的稳定性不同

NTILE的分组结果完全依赖ORDER BY定义的行顺序:

  • SQL Server:当ORDER BY列存在重复值时,会自动使用行的物理存储顺序(或聚集索引顺序)作为隐式第二排序键,保证排序的稳定性。相同Value的行顺序固定,因此NTILE分组结果一致。
  • PySpark:默认Window排序是不稳定的。当ORDER BY列有重复值时,相同值的行在分区内的顺序不确定(分布式计算的shuffle阶段行顺序可能随机变化),直接导致NTILE对这些行的分组分配出现差异。

2. 数据类型的影响

你观察到修改Value列类型会改变PySpark输出,原因在于:

  • 不同数据类型的排序比较逻辑存在细微差异,SQL表中Value为整数、Spark DataFrame中为浮点型时,底层排序规则不一致。
  • 数据类型变化会影响Spark shuffle过程中的行哈希值,改变相同值行的顺序,最终导致NTILE结果变化。

3. NTILE分组的细节差异

当分区行数无法被桶数整除时,两者的多余行分配策略存在细微差异:

  • SQL Server受稳定排序影响,会将多余行优先分配给前面的桶,结果固定。
  • PySpark在不稳定排序的前提下,多余行的分配会随相同值行的顺序变化而改变,进一步放大结果差异。

解决方案

若要让PySpark的NTILE结果与SQL Server对齐,可执行以下操作:

  • 在Window的orderBy中添加唯一排序键(如Date),保证排序稳定性:
    w = Window.partitionBy("Id").orderBy("Value", "Date")
    
  • 确保Spark DataFrame中Value列类型与SQL表一致(转换为整数类型):
    df = df.withColumn("Value", f.col("Value").cast("integer"))
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 11:10:38