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

Dataflow管道Worker磁盘空间耗尽崩溃问题排查求助

Dataflow Worker磁盘耗尽问题排查与解决

问题场景

我们使用Dataflow管道合并GCS存储桶中数千个2-4MB的Parquet文件,转换后的最终输出仅约1GB,但当待处理文件数量翻倍时,Worker频繁因磁盘空间耗尽崩溃;将Worker磁盘扩容至500GB后,管道可正常运行。目前无法明确磁盘占用的具体来源,怀疑是文件未正确关闭导致堆积,但尚未确认。

代码中的潜在问题与磁盘占用分析

结合提供的代码,磁盘耗尽的核心原因并非文件未关闭,而是以下两点:

1. GroupByKey触发的Shuffle临时文件堆积

GroupByKey操作会触发Dataflow的批处理shuffle流程,该流程会将中间数据写入Worker本地磁盘作为临时存储。当文件数量翻倍后,每个分组key对应的DataFrame列表规模随之增大,shuffle阶段的临时数据量会远超20GB标准磁盘的容纳上限,直接导致磁盘耗尽。

2. Parquet读取时的临时文件缓存

pd.read_parquet()依赖PyArrow或FASTParquet引擎,这些引擎在并发读取大量小文件时,会在本地磁盘生成临时缓存文件(默认存储在/tmp目录)。当Worker并发处理大量Parquet文件时,这些临时文件会持续堆积,进一步占用磁盘空间。

代码中的小错误

  • DownloadFile DoFn中未使用element作为文件名,file_name变量未定义,会导致运行时错误;
  • concat_and_export函数中file_name未定义,无法生成输出文件。

具体解决方案

1. 启用流式Shuffle替代本地磁盘Shuffle

流式Shuffle将中间数据存储到GCS而非Worker本地磁盘,彻底规避本地磁盘容量限制。在Pipeline配置中添加:

from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions, GoogleCloudOptions

options = PipelineOptions()
# 启用流式引擎(适用于批处理作业)
google_options = options.view_as(GoogleCloudOptions)
google_options.enable_streaming_engine = True
# 明确指定为批处理模式
standard_options = options.view_as(StandardOptions)
standard_options.streaming = False

with beam.Pipeline(options=options) as pipeline:
    # ... 后续管道逻辑 ...

2. 优化Parquet读取的临时缓存配置

调整PyArrow的读取参数,减少或避免磁盘缓存:

# 在DownloadFile的process方法中修改读取逻辑
import pyarrow.parquet as pq

with GcsIO().open(file_name, "rb") as file:
    # 使用内存映射读取,避免生成磁盘临时文件
    table = pq.read_table(file, memory_map=True, use_threads=True)
    df = table.to_pandas()
    yield get_grouping_key(file_name), do_transformations(df)

也可以通过设置环境变量指定临时目录(如指向更大的磁盘分区):

import os
os.environ["ARROW_TMP_DIR"] = "/var/dataflow/tmp"  # Dataflow默认的临时目录,空间更大

3. 优化分组逻辑,减少Shuffle数据量

如果分组key的分布不均匀,会导致单个Worker负载过高。可以在DownloadFile之后先进行局部合并,再执行GroupByKey,减少Shuffle阶段的数据量:

def partial_concat(dfs):
    # 对同一key的小批量DataFrame先合并
    return pd.concat(dfs)

# 在DownloadFile之后添加局部合并步骤
(pipeline
 # ... 前面的步骤 ...
 | "Download files" >> beam.ParDo(DownloadFile())
 | "Partial Combine Per Key" >> beam.CombinePerKey(partial_concat)  # 局部合并
 | "Group by" >> beam.GroupByKey()
 # ... 后续步骤 ...
)

4. 修复代码中的变量未定义问题

修正DownloadFile和concat_and_export中的文件名变量:

class DownloadFile(beam.DoFn):
     def process(self, element, ) -> Iterator[tuple[str, pd.DataFrame]]:
         file_name = element  # 使用输入的blob名称作为文件名
         with GcsIO().open(file_name, "rb") as file:
              yield get_grouping_key(file_name), do_transformations(pd.read_parquet(file))

def concat_and_export(element: tuple[str, list[pd.DataFrame]]):
   group_key, dfs = element
   df = pd.concat(dfs)
   output_file_name = f"gs://bucket_name/output/{group_key}.parquet"  # 根据key生成输出路径
   with GcsIO().open(output_file_name, "w") as file:
       df.to_parquet(file)

5. 监控磁盘使用情况

在DoFn中添加磁盘日志,明确占用来源:

import logging
import shutil

def log_disk_usage():
    disk_info = shutil.disk_usage('/')
    logging.info(
        f"磁盘使用情况: 总容量{disk_info.total//(1024**3)}GB, "
        f"已用{disk_info.used//(1024**3)}GB, 剩余{disk_info.free//(1024**3)}GB"
    )

# 在DownloadFile的process方法中调用
class DownloadFile(beam.DoFn):
     def process(self, element, ) -> Iterator[tuple[str, pd.DataFrame]]:
         log_disk_usage()
         # ... 后续逻辑 ...

内容的提问来源于stack exchange,提问作者unitrium

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 00:52:32