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
相关产品推荐
相关产品推荐

