使用withColumn添加列后DataFrame的Source_count值变更原因排查
DB2读取DataFrame后withColumn添加列导致行数减少的可能原因
- Spark惰性执行引发的数据源变更:Spark采用惰性计算模型,创建df1时仅生成读取DB2的执行计划,并未实际拉取数据。当创建df2并触发计数等动作时,才会真正执行读取操作。如果在df1创建到df2执行的间隙,DB2源表有数据被删除,就会导致实际读取到的行数变为12600。
- 隐式脏数据过滤:即便用
withColumn+lit硬编码加列,若原df1的Schema中存在字段类型不匹配的情况,Spark在实际执行读取时会自动过滤无法解析的脏数据。比如某数值型字段混入非数值内容,这类行会被标记为无效并排除,最终导致行数减少。 - DB2分区读取异常:如果DB2表是分区表,Spark读取时可能因分区元数据过期、分片策略异常,导致遗漏部分分区的数据。比如分区统计信息未更新,Spark只扫描了部分分区,从而少读了150行。
- 未缓存df1导致重复读取:如果没有对df1调用
cache()或persist()方法,df2的计算会重新发起DB2读取请求。若两次读取期间源数据发生变化,就会出现行数差异。 - DB2事务隔离级别影响:读取操作的事务隔离级别可能导致数据不一致。比如在READ COMMITTED级别下,若其他事务在读取期间提交了删除操作,当前读取会获取到已删除后的数据集,导致行数减少。
内容的提问来源于stack exchange,提问作者phani437
相关产品推荐
相关产品推荐

