Apache Airflow任务报错No module named 'airflow'及日志查看求助
场景说明
在Windows+WSL+Ubuntu环境下通过Docker部署Apache Airflow,编写my_dag.py后运行,read_api_data任务报错"No module name 'airflow'",且无法查看DAG详细日志回溯信息。
一、解决"No module named 'airflow'"报错
1. 确认DAG文件路径正确性
Docker部署的Airflow要求DAG文件必须放在挂载到容器的dags目录下。检查你的my_dag.py是否存放在WSL中对应挂载的本地目录(比如部署时指定的./dags:/opt/airflow/dags),确保容器能正常读取该文件。
2. 验证容器内Airflow环境
进入Airflow Worker容器(PythonOperator任务在Worker中执行),检查Airflow模块是否存在:
# 列出所有Airflow容器 docker ps # 替换<worker-container-name>为实际容器名,进入容器 docker exec -it <worker-container-name> bash # 检查Python环境中是否安装Airflow pip list | grep airflow
如果未找到Airflow模块,说明容器部署存在问题,重新执行Airflow的Docker部署步骤(参考官方docker-compose.yaml配置)。
3. 修正DAG代码中的依赖与逻辑错误
你的代码存在多处缺失依赖和逻辑问题,这也会导致任务执行异常,修正后的完整代码如下:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime import pandas as pd import requests # 新增缺失的requests导入 import json # 新增缺失的json导入 def _read_api_data(): x1 = requests.get('http://api.openweathermap.org/data/2.5/weather?q=London&appid=APIKEY') x2 = requests.get('http://api.openweathermap.org/data/2.5/weather?q=Moscow&appid=APIKEY') responses = [x1, x2] output_object = {"responses": []} for x in responses: json_object = x.json() output_object["responses"].append(json_object) json_str = json.dumps(output_object, indent=4) with open("/tmp/weather_data.json", "w") as outfile: outfile.write(json_str) def _download_data(): with open('/tmp/weather_data.json') as f: d = json.load(f) responses = d["responses"] temps_K = [round(r["main"]["temp"], 2) for r in responses] # 原代码重复转换温度,此处改为直接取开尔文温度 names = [r["name"] for r in responses] df = pd.DataFrame({"names": names, "temps_K": temps_K}) return df def _process_data(ti): df = ti.xcom_pull(task_ids="download_data") df["temps_C"] = round(df["temps_K"] - 273.15, 2) # 修正未定义的temps变量为df["temps_K"] df.to_csv('/tmp/processed_weather_data.csv') return df def _save_data(ti): df = ti.xcom_pull(task_ids="process_data") df.to_parquet('/tmp/weather.parquet') with DAG("weather_data_pipeline_dag", schedule_interval="@once", start_date=datetime(2024,7,14), catchup=False) as dag: # 添加catchup=False避免重复执行历史任务 taskPython0 = PythonOperator( task_id = "read_api_data", python_callable = _read_api_data ) taskPython1 = PythonOperator( task_id = "download_data", python_callable = _download_data ) taskPython2 = PythonOperator( task_id = "process_data", python_callable = _process_data ) taskPython3 = PythonOperator( task_id = "save_data", python_callable = _save_data ) taskPython0 >> taskPython1 >> taskPython2 >> taskPython3
4. 安装任务所需依赖
任务依赖requests、pandas、pyarrow(保存parquet格式需要),可通过以下两种方式安装:
- 自定义镜像:在Airflow Dockerfile中添加
RUN pip install requests pandas pyarrow,重新构建镜像 - 使用requirements.txt:在DAG目录下创建
requirements.txt,写入依赖项,然后在docker-compose.yaml中配置_PIP_ADDITIONAL_REQUIREMENTS环境变量,让Airflow自动安装
二、查看Airflow任务日志的方法
1. 通过Airflow Web UI查看
登录Airflow Web UI(默认端口8080),找到目标DAG,点击任务实例的Log按钮即可查看详细日志。如果无法查看,检查:
- 容器内
/opt/airflow/logs目录的权限,确保Airflow运行用户有读写权限 - 部署时是否正确挂载本地logs目录到容器(比如
./logs:/opt/airflow/logs)
2. 直接进入容器查看日志文件
Airflow日志默认存放在容器内/opt/airflow/logs/<dag-id>/<task-id>/<execution-date>目录下,执行以下命令查看:
docker exec -it <worker-container-name> bash cd /opt/airflow/logs/weather_data_pipeline_dag/read_api_data/<YYYY-MM-DDTHH:MM:SS> cat 1.log
3. 使用Docker命令查看容器日志
如果是容器本身运行异常,可直接查看Worker容器的系统日志:
docker logs <worker-container-name>
内容的提问来源于stack exchange,提问作者Excentricitet

