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

Airflow僵尸任务(Zombie Job)报错排查及解决请求

问题分析与解决方案

错误原因定位

Airflow中出现"Detected zombie job"通常意味着Worker进程被外部系统强制终止,最常见的触发场景是:

  • 任务处理大文件时内存耗尽,被系统OOM Killer杀掉
  • Worker资源配置不足,无法承载任务负载
  • 代码中未捕获的异常导致进程崩溃(你的错误日志无直接异常栈,更倾向于资源问题)

结合你的代码来看,核心问题是format_imdb_files函数直接用pd.read_csv加载整个TSV文件到内存,如果文件体积较大(比如几个G),很容易触发内存溢出,导致Worker进程被终止,任务变成僵尸状态。


具体解决方案

1. 分块处理大文件(优先推荐)

修改TSV转Parquet的逻辑,使用Pandas的chunksize参数分批读取写入,避免一次性加载全量数据到内存:

def format_imdb_files():
    pm = PathManager('/opt/airflow/datalake')
    imdb_path = pm.get_local_path('raw', 'imdb')
    chunk_size = 100000  # 根据内存情况调整,比如每次读10万行
    
    for file in os.listdir(imdb_path):
        if file.endswith('.tsv'):
            file_path = os.path.join(imdb_path, file)
            parquet_file = file_path.replace(".tsv", ".parquet")
            
            # 分块读取并写入Parquet
            chunks = pd.read_csv(file_path, sep='\t', chunksize=chunk_size)
            # 写入第一个chunk时创建文件,后续追加
            first_chunk = True
            for chunk in chunks:
                if first_chunk:
                    chunk.to_parquet(parquet_file, index=False)
                    first_chunk = False
                else:
                    chunk.to_parquet(parquet_file, index=False, mode='a')
            
            os.remove(file_path)

2. 复用已实现的FileHandler类

你已经写了FileHandler.convert_tsv_to_parquet方法,建议直接复用,同时给这个方法加上分块处理的优化:

# 修改FileHandler的convert_tsv_to_parquet方法
def convert_tsv_to_parquet(self, tsv_file, chunk_size=100000):
    parquet_file = tsv_file.replace(".tsv", ".parquet")
    chunks = pd.read_csv(tsv_file, sep='\t', chunksize=chunk_size)
    first_chunk = True
    for chunk in chunks:
        if first_chunk:
            chunk.to_parquet(parquet_file, index=False)
            first_chunk = False
        else:
            chunk.to_parquet(parquet_file, index=False, mode='a')
    os.remove(tsv_file)
    return parquet_file  # 改为返回文件路径更实用

# 然后修改format_imdb_files函数
def format_imdb_files():
    pm = PathManager('/opt/airflow/datalake')
    imdb_path = pm.get_local_path('raw', 'imdb')
    # 用FileHandler的list_files_in_directory遍历文件,更可靠
    tsv_files = file_handler.list_files_in_directory(imdb_path)
    for file_path in tsv_files:
        file_handler.convert_tsv_to_parquet(file_path)

3. 调整Airflow Worker资源配置

如果你的Airflow用CeleryExecutor,需要调整Worker的内存限制:

  • 在airflow.cfg中设置worker_max_memory_per_child = 4096(单位MB,根据实际资源调整)
  • 或者启动Worker时指定--max-memory-per-child=4096参数

如果用KubernetesExecutor,需要在Pod模板中增加资源请求和限制:

resources:
  requests:
    memory: "4Gi"
  limits:
    memory: "8Gi"

4. 调试与排查建议

  • 查看Worker节点的系统日志(比如/var/log/syslog或dmesg),搜索OOM killer确认是否是内存溢出导致进程被终止
  • 用top或htop命令监控任务运行时的内存占用,确认是否超过Worker的内存上限
  • 先用小体积的TSV文件测试,验证转换逻辑本身是否正常,排除文件格式或权限问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 06:53:13