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

Spark SQL中SELECT *与指定列查询返回行数不一致问题排查

行数差异原因分析

问题背景

在Azure Blob存储数据湖上通过PySpark执行SQL查询,数据湖采用多层嵌套分区。执行两个过滤逻辑一致的查询,结果行数出现明显差异:

  • 查询1指定col1、col2、col3列查询,返回2000万行
  • 查询2使用SELECT *查询,返回2800万行

数据湖容器结构

.
|-- company=first
|   |-- year=2021
|   |-- year=2022
|   |-- year=2023
|   `-- year=2024
|-- company=other
|   |-- year=2021
|   |-- year=2022
|   |-- year=2023
|   `-- year=2024
`-- company=second
    |-- year=2021
    |-- year=2022
    |-- year=2023
    `-- year=2024

查询1 SQL

-- File : query_1.sql
SELECT
    col1,
    col2,
    col3
FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/company=first/year=2023`
WHERE col1 IN ('A','B','C')
UNION
SELECT
    col1,
    col2,
    col3
FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/company=second/year=2023`
WHERE col1 IN ('A','B','C')
UNION
SELECT
    col1,
    col2,
    col3
FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/company=other/year=2023`
WHERE col1 IN ('A','B','C')

查询2 SQL

-- File : query_2.sql
SELECT *
FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/company=first/year=2023`
WHERE col1 IN ('A','B','C')
UNION
SELECT *
FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/company=second/year=2023`
WHERE col1 IN ('A','B','C')
UNION
SELECT *
FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/company=other/year=2023`
WHERE col1 IN ('A','B','C')

PySpark执行代码

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

query1 = open("query_1.sql", mode="r").read()
df1 = spark.sql(query1)
df1.count()   # Gives 20 million rows

query2 = open("query_2.sql", mode="r").read()
df2 = spark.sql(query2)
df2.count()   # Gives 28 million rows

原因分析

核心原因:UNION的去重逻辑差异

UNION关键字的核心特性是对结果集自动去重,但去重的判定依据是查询返回的所有列:

  • 查询1仅返回col1、col2、col3三列,只要这三列的值完全相同,无论其他列是什么,都会被判定为重复行并合并。
  • 查询2返回所有列(包括隐式的分区列company、year,以及数据中的其他未指定列),只有当所有列的值完全一致时,才会被判定为重复行并去重。

具体场景拆解

  1. 分区列的影响:
    虽然查询路径已经指定了company和year的具体值,但SELECT *会把这两个分区列自动包含到结果集中。不同company分区下的行,即使col1、col2、col3完全相同,因为company列值不同(比如first和second),在查询2中会被视为不同行,不会被UNION去重;但在查询1中,因为未包含company列,这些行的三列值相同,会被合并成一行。

  2. 其他列的重复可能性:
    如果数据中存在其他列(比如col4、col5等),即使col1、col2、col3的值相同,只要这些额外列的值不同,查询2就会保留这些行;而查询1因为只取三列,会把这些行当作重复行去掉。

验证方法

可以通过以下步骤验证上述结论:

  • 给查询1的结果添加company列,执行count:

    SELECT col1, col2, col3, company
    FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/company=first/year=2023`
    WHERE col1 IN ('A','B','C')
    UNION
    SELECT col1, col2, col3, company
    FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/company=second/year=2023`
    WHERE col1 IN ('A','B','C')
    UNION
    SELECT col1, col2, col3, company
    FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/company=other/year=2023`
    WHERE col1 IN ('A','B','C')
    

    该查询的结果行数会接近查询2的2800万行,因为加入company列后,去重条件变严格,更多行被保留。

  • 检查数据中是否存在col1、col2、col3相同但其他列不同的行:

    SELECT col1, col2, col3, COUNT(DISTINCT col4)
    FROM parquet.`abfs://container@storageaccount.dfs.core.windows.net/year=2023`
    WHERE col1 IN ('A','B','C')
    GROUP BY col1, col2, col3
    HAVING COUNT(DISTINCT col4) > 1
    

    如果有返回结果,说明这类行是导致行数差异的原因之一。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 20:03:13