Synapse Analytics中PySpark不同池大小下DataFrame count结果不一致问题
解决Spark读取CSV时count()结果不一致的问题
问题场景
构建湖仓历史层时,读取ADLS上的CSV文件并计划转为Delta表,多次调用count()验证记录数时出现异常:
- 大集群(16 vCores / 128 GB,3至15节点):count结果波动(如50270、45935、41640等)
- 小集群(4 vCores / 32 GB,3至15节点):count结果始终为50270,完全一致
测试代码:
datDf = spark.read.option("quote", "").load(f"abfss://dls@xxxx.dfs.core.windows.net/xx/xx/xx/2022/10/20/xx.dat", format='csv', sep =';', header = False) print(datDf.count()) print(datDf.count()) # 重复调用count()共11次
可能原因
- CSV格式解析歧义:大集群并行读取时,文件中存在未正确处理的换行符、分隔符,导致不同分区解析记录时出现偏差
- 并行分区拆分问题:大集群分区数更多,Spark拆分文件时可能将单条记录拆分到两个分区,造成统计丢失或重复
- 文件IO偶发异常:ADLS文件存在块损坏或异步写入的情况,大集群多节点读取时偶尔触发异常
- 重复读取的IO差异:每次count()都会重新读取文件,大集群的分布式IO环境更容易出现偶发的读取异常
解决方案
1. 验证CSV文件格式完整性
先确认文件本身没有格式问题:
- 用Spark读取整个文件的文本内容,检查是否有异常换行、未转义的分隔符:
wholeTextDf = spark.read.text("abfss://dls@xxxx.dfs.core.windows.net/xx/xx/xx/2022/10/20/xx.dat") wholeTextDf.show(truncate=False)
- 通过本地工具(如
wc -l)或Azure CLI下载文件后统计行数,对比正确的记录数
2. 调整CSV读取参数,消除解析歧义
根据文件实际格式调整读取参数:
- 如果存在多行记录,开启
multiLine确保跨多行的记录被正确识别:
datDf = spark.read.option("quote", "") \ .option("multiLine", True) \ .load("abfss://dls@xxxx.dfs.core.windows.net/xx/xx/xx/2022/10/20/xx.dat", format='csv', sep=';', header=False)
- 如果确认没有多行记录,明确指定换行符
lineSep,避免Spark自动推断错误:
datDf = spark.read.option("quote", "") \ .option("lineSep", "\n") \ .load("abfss://dls@xxxx.dfs.core.windows.net/xx/xx/xx/2022/10/20/xx.dat", format='csv', sep=';', header=False)
3. 缓存DataFrame,避免重复读取文件
每次调用count()都会重新扫描文件,缓存后只读取一次,消除IO波动影响:
datDf = spark.read.option("quote", "") \ .load("abfss://dls@xxxx.dfs.core.windows.net/xx/xx/xx/2022/10/20/xx.dat", format='csv', sep=';', header=False) datDf.cache() # 或使用persist(StorageLevel.MEMORY_AND_DISK),适合大文件 # 后续所有count()都会使用缓存数据,结果一致 print(datDf.count()) print(datDf.count())
4. 检查ADLS文件一致性
- 使用Azure Storage Explorer查看文件的ETag和完整性,确认文件没有被异步修改或损坏
- 尝试重新上传文件到ADLS,排除文件本身的损坏问题
5. 手动指定Spark分区数
避免Spark自动拆分文件时的不合理分区,手动设置适合大集群的分区数:
datDf = spark.read.option("quote", "") \ .load("abfss://dls@xxxx.dfs.core.windows.net/xx/xx/xx/2022/10/20/xx.dat", format='csv', sep=';', header=False) \ .repartition(10) # 根据集群节点数和文件大小调整,比如每节点2-4个分区
内容的提问来源于stack exchange,提问作者justsander
相关产品推荐
相关产品推荐

