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,以及数据中的其他未指定列),只有当所有列的值完全一致时,才会被判定为重复行并去重。
具体场景拆解
分区列的影响:
虽然查询路径已经指定了company和year的具体值,但SELECT *会把这两个分区列自动包含到结果集中。不同company分区下的行,即使col1、col2、col3完全相同,因为company列值不同(比如first和second),在查询2中会被视为不同行,不会被UNION去重;但在查询1中,因为未包含company列,这些行的三列值相同,会被合并成一行。其他列的重复可能性:
如果数据中存在其他列(比如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
相关产品推荐
相关产品推荐

