Airflow中Polars读取Parquet文件无报错却一直运行的问题排查
以下是可能导致Airflow中Polars读取Parquet任务持续运行无日志、但Pandas和本地运行正常的几个原因及对应解决方法:
1. Polars内存管理与资源限制不匹配
Airflow的worker进程(尤其是SequentialExecutor下的单进程环境)可能存在内存限制,Polars默认的内存映射或高内存占用模式可能在受限环境下触发死锁或无限等待,且未抛出可被Airflow捕获的异常。而Pandas的内存处理逻辑更偏向于渐进式加载,适配性更强。
解决方法:
- 禁用内存映射读取:
df = pl.read_parquet("file.parquet", memory_map=False) - 手动设置Polars内存上限:
import polars as pl pl.Config.set_memory_limit(1024 * 1024 * 1024) # 限制为1GB,根据实际环境调整 df = pl.read_parquet("file.parquet")
2. SQLite锁机制与Polars线程冲突
搭配SequentialExecutor的SQLite数据库是单文件锁模式,Polars读取Parquet时默认会启用多线程加速(如PyArrow的线程池),这可能与SQLite的文件锁产生冲突,导致Airflow任务进程被卡住,无法写入日志或继续执行。而Pandas读取Parquet默认使用单线程,不会触发该冲突。
解决方法:
- 强制Polars使用单线程:
import polars as pl pl.Config.set_threads(1) df = pl.read_parquet("file.parquet") - 调整Airflow的SQLite连接配置,在
airflow.cfg中修改sql_alchemy_conn,添加线程安全参数:sql_alchemy_conn = sqlite:////path/to/airflow.db?check_same_thread=False
3. Airflow日志捕获失效
Polars的内部操作日志可能未被Airflow的日志系统正确捕获,导致任务实际卡住但无日志输出;或者任务已抛出异常,但进程因锁或内存问题无法将日志写入SQLite数据库。
解决方法:
- 在任务函数中手动添加日志打点并强制捕获异常:
通过查看日志是否打印"开始执行"但无后续内容,确认任务是否卡在读取步骤。import logging from airflow.decorators import task import polars as pl @task def read_parquet_task(): try: logging.info("开始执行Polars读取Parquet操作") df = pl.read_parquet("/absolute/path/to/file.parquet") logging.info(f"读取完成,数据行数:{len(df)}") return df.shape except Exception as e: logging.error(f"读取失败:{str(e)}", exc_info=True) raise
4. 依赖版本不兼容
Airflow环境中的Polars、PyArrow版本与本地运行环境不一致,可能导致Parquet读取时出现隐性死锁(如PyArrow版本不兼容引发的线程阻塞)。Pandas对PyArrow版本的兼容性范围更广,因此未触发问题。
解决方法:
- 对比本地与Airflow环境的依赖版本:
本地执行:pip freeze | grep -E "polars|pyarrow"
Airflow worker中执行相同命令,确保版本一致。 - 升级/降级到兼容版本,例如Polars 0.20.x系列适配PyArrow 14.x版本。
5. 文件路径/权限问题
Airflow worker进程的运行用户与本地用户不同,可能导致Parquet文件无读取权限;或使用相对路径时,Airflow的工作目录与本地不一致,实际读取的文件不存在,但Polars未抛出明确异常而是无限重试。
解决方法:
- 使用绝对路径读取Parquet文件:
df = pl.read_parquet("/opt/airflow/data/file.parquet") - 检查Airflow运行用户对文件及所在目录的读取权限,确保权限一致。
内容的提问来源于stack exchange,提问作者Ekaterina Lubyankina

