如何在Polars反连接中将Null视为匹配值?
Polars反连接中让Null视为匹配的解决方案
问题场景
使用Polars的anti join生成数据集增量时,默认规则下Null值不被视为匹配项。比如当两张表中存在相同的行但某列均为Null时,该行会被误识别为增量数据,不符合预期需求——需要将两边均为Null的行视为匹配,仅保留真正新增的行。
可行方案
方案1:自定义匹配条件实现半连接后筛选
通过生成列级的匹配规则(列值相等 或 两边均为Null),先找到my_update中与my_dataset匹配的所有行,再从my_update中排除这些行得到增量。该方案无需修改原始数据,适配多列、多数据类型场景。
代码示例:
import polars as pl # 定义原始数据集(lazy模式) my_dataset = pl.DataFrame( { "a": [1, 2, 3, 5], "b": ["x", "y", "z", pl.Null], "c": [0.1, 0.2, 0.3, 0.5], } ).lazy() my_update = pl.DataFrame( { "a": [1, 2, 3, 4, 5], "b": ["x", "y", "z", "q", pl.Null], "c": [0.1, 0.2, 0.3, 0.4, 0.5], } ).lazy() cols = my_update.columns # 生成每一列的匹配规则:值相等 或 两边都是Null match_conditions = pl.all_horizontal( [(pl.col(c) == pl.col(f"{c}_right")) | (pl.col(c).is_null() & pl.col(f"{c}_right").is_null()) for c in cols] ) # 找到my_update中与my_dataset匹配的行 matched_rows = my_update.join( my_dataset, condition=match_conditions, how="semi" ) # 反选匹配行,得到正确增量 my_delta = my_update.join(matched_rows, on=cols, how="anti") # 执行并查看结果 print(my_delta.collect())
执行结果:
shape: (1, 3) ┌─────┬─────┬──────┐ │ a ┆ b ┆ c │ │ --- ┆ --- ┆ --- │ │ i64 ┆ str ┆ f64 │ ╞═════╪═════╪══════╡ │ 4 ┆ q ┆ 0.4 │ └─────┴─────┴──────┘
方案2:用唯一占位符填充Null后恢复
若更倾向于使用原生anti join逻辑,可先为不同数据类型设置唯一占位符,填充Null后执行连接,最后再将占位符还原为Null。需注意占位符不能与实际数据中的值冲突。
代码示例:
import polars as pl # 定义原始数据集(lazy模式) my_dataset = pl.DataFrame( { "a": [1, 2, 3, 5], "b": ["x", "y", "z", pl.Null], "c": [0.1, 0.2, 0.3, 0.5], } ).lazy() my_update = pl.DataFrame( { "a": [1, 2, 3, 4, 5], "b": ["x", "y", "z", "q", pl.Null], "c": [0.1, 0.2, 0.3, 0.4, 0.5], } ).lazy() cols = my_update.columns # 为不同数据类型定义唯一占位符(需确保不会与业务数据冲突) placeholder_map = { pl.Utf8: "__POLARS_UNIQUE_NULL__", pl.Int64: 999999999, pl.Float64: 999999999.999, # 根据实际数据类型补充 } # 填充Null为占位符 def fill_null_placeholder(df): return df.with_columns( [pl.col(c).fill_null(placeholder_map[pl.col(c).dtype()]) for c in cols] ) dataset_filled = fill_null_placeholder(my_dataset) update_filled = fill_null_placeholder(my_update) # 执行原生anti join delta_filled = update_filled.join(dataset_filled, on=cols, how="anti") # 将占位符还原为Null my_delta = delta_filled.with_columns( [pl.col(c).replace(placeholder_map[pl.col(c).dtype()], None) for c in cols] ) # 执行并查看结果 print(my_delta.collect())
执行结果与方案1一致。
方案对比
- 方案1:无需修改原始数据,性能更优,适合大数据量、多列场景;代码需动态生成匹配条件,逻辑稍复杂。
- 方案2:逻辑简单,依赖原生join;需适配所有数据类型的占位符,存在与业务数据冲突的风险,多了填充/还原的计算步骤。
内容的提问来源于stack exchange,提问作者rxFt20
相关产品推荐
相关产品推荐

