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

使用Pyarrow is_null函数查询Delta表疑似存在异常

Delta表转Pyarrow Dataset后is_null过滤仅返回全空文件行的问题

问题描述

将Delta表通过delta-rs转为Pyarrow Dataset后,使用Pyarrow的is_null()表达式过滤时,仅返回对应Parquet文件中过滤列所有行均为null的记录,无法匹配到文件中部分行是null的情况。例如示例中key='a'的文件包含1条null行,但被过滤逻辑跳过,最终只返回key='b'的2条全null行,与预期的3行结果不符。

复现步骤

1. 用PySpark创建Delta表

import pyspark
from delta import configure_spark_with_delta_pip
from pyspark.sql.types import StructType, StructField, StringType

builder = pyspark.sql.SparkSession.builder.appName("MyApp") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
spark = configure_spark_with_delta_pip(builder).getOrCreate()

data = [['a', None],
        ['a', 'exists'],
        ['b', None],
        ['b', None]]

schema = StructType([
    StructField("key", StringType()),
    StructField('null_column', StringType())
])

spark_df = spark.createDataFrame(data=data, schema=schema)
spark_df.repartition(2, "key").write.format("delta").mode("overwrite").save("./table")

2. Delta表元数据(日志)

日志中正确记录了各文件的null数量:

{"commitInfo":{"timestamp":1687784715118,"operation":"WRITE","operationParameters":{"mode":"Overwrite","partitionBy":"[]"},"isolationLevel":"Serializable","isBlindAppend":false,"operationMetrics":{"numFiles":"2","numOutputRows":"4","numOutputBytes":"1436"},"engineInfo":"Apache-Spark/3.4.0 Delta-Lake/2.4.0","txnId":"d718e2b3-c062-4783-987f-97cab1c5f110"}}
{"protocol":{"minReaderVersion":1,"minWriterVersion":2}}
{"metaData":{"id":"eec58790-7456-4bfb-9fd5-6a576e46851b","format":{"provider":"parquet","options":{}},"schemaString":"{\"type\":\"struct\",\"fields\":[{\"name\":\"key\",\"type\":\"string\",\"nullable\":true,\"metadata\":{}},{\"name\":\"null_column\",\"type\":\"string\",\"nullable\":true,\"metadata\":{}}]}","partitionColumns":[],"configuration":{},"createdTime":1687784710499}}
{"add":{"path":"part-00000-ada91cf1-2cd6-4928-af3b-0b1bcc987bfc-c000.snappy.parquet","partitionValues":{},"size":745,"modificationTime":1687784712864,"dataChange":true,"stats":"{\"numRecords\":2,\"minValues\":{\"key\":\"a\",\"null_column\":\"exists\"},\"maxValues\":{\"key\":\"a\",\"null_column\":\"exists\"},\"nullCount\":{\"key\":0,\"null_column\":1}}"}}
{"add":{"path":"part-00001-85bbf245-a29d-40bc-be96-86391569710c-c000.snappy.parquet","partitionValues":{},"size":691,"modificationTime":1687784712854,"dataChange":true,"stats":"{\"numRecords\":2,\"minValues\":{\"key\":\"b\"},\"maxValues\":{\"key\":\"b\"},\"nullCount\":{\"key\":0,\"null_column\":2}}"}}

3. delta-rs查询代码

from deltalake import DeltaTable
import pyarrow.compute as pc

ds = DeltaTable("./table").to_pyarrow_dataset()
pat = ds.to_table(filter=(pc.field("null_column").is_null()))
print(pat.num_rows)
print(pat.to_pandas().head())

执行后仅返回2条key='b'的行,不符合预期。

原因分析

这是Pyarrow Dataset的谓词下推优化结合Delta表统计信息的特性导致的:

  • Delta表的每个文件元数据中,minValues和maxValues默认不包含null值(null单独通过nullCount字段记录)
  • Pyarrow在处理Delta转Dataset的谓词下推时,仅依据minValues和maxValues判断文件是否匹配过滤条件:如果文件的null_column的min/max都是非null值,Pyarrow会直接跳过扫描该文件,认为其中不存在null行
  • 但实际上key='a'的文件包含1条null行,只是因为Spark默认生成stats时未将null纳入min/max,导致Pyarrow误判。

解决方案

方案1:先读取全表再过滤(临时快速解决)

跳过Pyarrow的文件级谓词下推,先读取所有文件的数据,再执行过滤操作:

from deltalake import DeltaTable
import pyarrow.compute as pc

ds = DeltaTable("./table").to_pyarrow_dataset()
# 读取全表数据
full_table = ds.to_table()
# 全表内过滤null行
filtered_table = full_table.filter(pc.field("null_column").is_null())
print(filtered_table.num_rows)  # 输出3
print(filtered_table.to_pandas().head())

方案2:修改Spark配置生成包含null的统计信息(长期解决方案)

通过Spark配置让Delta表的stats将null纳入minValues/maxValues,这样Pyarrow就能正确识别文件中是否包含null行:
在创建SparkSession时添加spark.sql.delta.stats.collectNulls=true配置:

builder = pyspark.sql.SparkSession.builder.appName("MyApp") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .config("spark.sql.delta.stats.collectNulls", "true")  # 新增配置
spark = configure_spark_with_delta_pip(builder).getOrCreate()

重新生成Delta表后,原有的过滤代码就能直接返回3行结果,无需修改查询逻辑。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 20:08:09