PySpark过滤行为不符合预期的问题求助
先明确你的核心需求
你想要过滤(移除)的是**满足CID = '-' 并且 trait_diff 为 null**的行,但你尝试的多个过滤语句都没达到预期,反而都只移除了CID 为 null 且 trait_diff 为 null的3行,我们来逐个拆解背后的原因:
1. 为什么前3段代码都只移除了CID和trait_diff均为null的行?
我们逐个分析你的代码逻辑:
第一段代码:
df = df.filter( ~((col('CID') == "-") & (col('CID').isNull()) & col('trait_diff').isNull()) )这里有个明显的逻辑矛盾:一个字段不可能同时等于字符串
'-'和null,所以(col('CID') == "-") & (col('CID').isNull())这个子条件永远是False,整个括号里的判断就变成了False & col('trait_diff').isNull(),结果也永远是False。取反~之后就是True,理论上这个过滤应该保留所有行,不会移除任何数据。但你说它和其他代码一样移除了3行,大概率是Databricks的DataFrame缓存问题——你之前运行过移除那3行的代码,df被自动缓存了,后续修改代码后缓存没刷新,导致你看到的还是旧结果。建议每次测试新代码前先运行
df.unpersist()清除缓存。第二段代码:
df = df.filter( ~((col('CID').isNull()) & col('trait_diff').isNull()) )这段逻辑就是明确移除
CID为null且trait_diff为null的行,所以它移除那3行是完全符合逻辑的,但这和你要移除CID='-'且trait_diff为null的需求完全无关。第三段代码:
df = df.filter( ~((col('CID') == "-") & col('trait_diff').isNull()) )这段代码的逻辑本应该是移除
CID='-'且trait_diff为null的行,但你说它还是只移除了那3行,只有两种可能:- 你的数据里根本不存在
CID='-'且trait_diff为null的行,所以这个过滤条件没有匹配到任何行,自然不会移除额外数据; - 还是缓存的锅,你看到的是之前的旧结果。
- 你的数据里根本不存在
2. 为什么(col('CID') != "-") | col('trait_diff').isNotNull()也没生效?
根据德摩根定律,这个语句和你的第三段代码~((col('CID') == "-") & col('trait_diff').isNull())是完全等价的!所以它的行为和第三段代码一致很正常——如果数据里没有CID='-'且trait_diff为null的行,那这个条件会保留所有行;如果有缓存,结果还是旧的。
3. 为什么最后一段代码看起来生效了?
你最后用的这段代码:
df = df.filter( ~((col('CID') == "-") & col('CID').isNotNull() & col('trait_diff').isNull()) )
其实和第三段代码逻辑完全等价!因为如果CID == "-",那它必然是isNotNull()的(字符串'-'不是null),所以加不加col('CID').isNotNull()这个条件,整个括号里的判断结果是一样的。
你觉得它生效,可能是因为这次你刚好清除了缓存,看到了真实结果;或者你的“预期行数”其实包含了移除那3行的情况,而这次结果刚好符合预期,但本质上它和第三段代码没有区别。
4. 快速排查建议
你现在需要先搞清楚数据的真实情况,才能解决问题:
- 先检查数据里到底有没有目标行:直接运行这段代码查看你要移除的行是否存在:
如果# 查看符合你要移除条件的行 target_rows = df.filter((col('CID') == "-") & col('trait_diff').isNull()) print(f"符合条件的行数:{target_rows.count()}") target_rows.display() # 在Databricks里直接展示数据count()返回0,说明你的数据里根本没有CID='-'且trait_diff为null的行,这就是所有过滤语句都没达到预期的根源; - 强制清除缓存:每次测试新代码前,先运行
df.unpersist(),避免缓存干扰; - 如果需要同时移除两种行:如果你其实想同时移除
CID='-'且trait_diff为null以及CID为null且trait_diff为null的行,那正确的过滤语句应该是:# 写法1:直接取反要移除的组合条件 df = df.filter( ~( ((col('CID') == "-") | col('CID').isNull()) & col('trait_diff').isNull() ) ) # 写法2:用德摩根定律改写,更易读 df = df.filter( (col('CID') != "-") & col('CID').isNotNull() | col('trait_diff').isNotNull() )
备注:内容来源于stack exchange,提问作者S. Nasir

