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

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次

可能原因

  1. CSV格式解析歧义:大集群并行读取时,文件中存在未正确处理的换行符、分隔符,导致不同分区解析记录时出现偏差
  2. 并行分区拆分问题:大集群分区数更多,Spark拆分文件时可能将单条记录拆分到两个分区,造成统计丢失或重复
  3. 文件IO偶发异常:ADLS文件存在块损坏或异步写入的情况,大集群多节点读取时偶尔触发异常
  4. 重复读取的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 00:30:43