Airflow PythonOperator执行HTTP请求时调度器停止问题排查
在M1 Mac本地Airflow 2.6.3中调用HTTP请求时任务卡住、调度器崩溃的问题
环境与现象
- M1芯片Mac,本地虚拟环境运行Airflow 2.6.3
- 通过PythonOperator调用
utils.get_rate()执行HTTP请求获取汇率时,手动触发DAG后:- 任务一直处于运行状态,日志卡在发送HTTP请求步骤
- Airflow调度器停止工作
- 单独运行
utils.get_rate()脚本正常,PythonOperator仅执行打印操作也正常 - 重启调度器时提示端口8793被占用
DAG代码
from Exchange_Rate import utils from airflow import DAG from airflow.operators.python import PythonOperator import datetime default_args = { 'owner': 'admin', 'email':"****@gmail.com" } dag = DAG( dag_id="exchange-rate-to-line", schedule="0 18 * * *", start_date=datetime.datetime(2023, 7, 15), catchup=False, default_args=default_args, ) task = PythonOperator( task_id=f'exchange_rate_to_line', python_callable= utils.get_rate, dag=dag )
DAG执行日志
AIRFLOW_CTX_DAG_ID='exchange-rate-to-line' AIRFLOW_CTX_TASK_ID='exchange_rate_to_line' AIRFLOW_CTX_EXECUTION_DATE='2023-07-29T12:21:56.333954+00:00' AIRFLOW_CTX_TRY_NUMBER='1' AIRFLOW_CTX_DAG_RUN_ID='manual__2023-07-29T12:21:56.333954+00:00' [2023-07-29T13:22:00.084+0100] {utils.py:23} INFO - Getting exchange rate API key... [2023-07-29T13:22:00.089+0100] {utils.py:42} INFO - Sending HTTP request to exchange rate API...
调度器重启日志
[2023-07-29 13:25:45 +0100] [73796] [INFO] Starting gunicorn 20.1.0 [2023-07-29 13:25:45 +0100] [73796] [ERROR] Connection in use: ('::', 8793) [2023-07-29 13:25:45 +0100] [73796] [ERROR] Retrying in 1 second. [2023-07-29T13:25:46.326+0100] {settings.py:60} INFO - Configured default timezone Timezone('Europe/London') [2023-07-29T13:25:46.345+0100] {manager.py:411} WARNING - Because we cannot use more than 1 thread (parsing_processes = 2) when using sqlite. So we set parallelism to 1. [2023-07-29 13:25:46 +0100] [73796] [ERROR] Connection in use: ('::', 8793) [2023-07-29 13:25:46 +0100] [73796] [ERROR] Retrying in 1 second. [2023-07-29 13:25:47 +0100] [73796] [ERROR] Connection in use: ('::', 8793) [2023-07-29 13:25:47 +0100] [73796] [ERROR] Retrying in 1 second. [2023-07-29 13:25:48 +0100] [73796] [ERROR] Connection in use: ('::', 8793) [2023-07-29 13:25:48 +0100] [73796] [ERROR] Retrying in 1 second. [2023-07-29 13:25:49 +0100] [73796] [ERROR] Connection in use: ('::', 8793) [2023-07-29 13:25:49 +0100] [73796] [ERROR] Retrying in 1 second. [2023-07-29T13:25:50.338+0100] {dagrun.py:609} ERROR - Marking run <DagRun exchange-rate-to-line @ 2023-07-29 12:21:56.333954+00:00: manual__2023-07-29T12:21:56.333954+00:00, state:running, queued_at: 2023-07-29 12:21:56.352181+00:00. externally triggered: True> failed [2023-07-29T13:25:50.339+0100] {dagrun.py:681} INFO - DagRun Finished: dag_id=exchange-rate-to-line, execution_date=2023-07-29 12:21:56.333954+00:00, run_id=manual__2023-07-29T12:21:56.333954+00:00, run_start_date=2023-07-29 12:21:57.108894+00:00, run_end_date=2023-07-29 12:25:50.339046+00:00, run_duration=233.230152, state=failed, external_trigger=True, run_type=manual, data_interval_start=2023-07-27 17:00:00+00:00, data_interval_end=2023-07-28 17:00:00+00:00, dag_hash=5321608c331baa96d84ddc8b0af4ee7d [2023-07-29T13:25:50.349+0100] {dag.py:3504} INFO - Setting next_dagrun for exchange-rate-to-line to 2023-07-28T17:00:00+00:00, run_after=2023-07-29T17:00:00+00:00 [2023-07-29 13:25:50 +0100] [73796] [ERROR] Can't connect to ('::', 8793)
utils.py代码
import requests import os import json from dotenv import load_dotenv from datetime import date import logging # Logging setting logger = logging.getLogger() logger.setLevel(20) fhandler = logging.FileHandler(filename=r'main.log') formatter = logging.Formatter('%(asctime)s - %(message)s') load_dotenv() def get_api_key(): """ Get api key of ExchangeRate-API """ try: logging.info('Getting exchange rate API key...') api_key = os.environ['EXCHANGE_RATE_API_KEY'] return api_key except KeyError: logging.error('Key Error in get_api_key(). Check the environmental variable EXCHANGE_RATE_API_KEY.') raise except Exception as e: logging.error(f'Error in get_api_key(). {type(e)} : {e}') raise def get_rate(): """ Get JPY exchange rate of 1 GBP """ api_key = get_api_key() req = f'https://v6.exchangerate-api.com/v6/{api_key}/latest/GBP' try: logging.info('Sending HTTP request to exchange rate API...') response = requests.get(req, timeout = 10) jpy = response.json()['conversion_rates']['JPY'] return jpy except Exception as e: logging.error(f'Error in get_rate(). {type(e)} : {e}') raise
排查信息
- 测试
get_api_key()在DAG中可正常返回正确值,问题集中在HTTP请求执行环节
解决方案
1. 修复日志配置冲突问题
你的utils.py中直接修改了根日志器的配置并添加文件处理器,这会和Airflow自身的日志系统产生冲突,导致任务进程阻塞。修改日志配置,改为获取当前模块的日志器:
# 替换原日志配置 logger = logging.getLogger(__name__) # 获取当前模块的日志器,而非根日志器 logger.setLevel(logging.INFO) # 用明确的常量替代数字20 # 移除自定义的FileHandler,Airflow会自动将日志输出到任务日志文件
2. 清理占用端口的进程
任务卡住后进程未正常退出,导致端口8793被占用。执行以下命令强制清理:
# 查找占用8793端口的进程ID lsof -i :8793 # 杀死对应进程,替换<PID>为查到的进程ID kill -9 <PID>
3. 验证HTTP请求的运行环境
虽然单独运行脚本正常,但Airflow任务的运行环境可能存在代理或网络配置差异。可以在get_rate()中添加代理参数(如需要),或增加请求日志排查:
# 若环境需要代理,添加proxies参数 proxies = { 'http': 'http://your-proxy:port', 'https': 'http://your-proxy:port', } response = requests.get(req, timeout=10, proxies=proxies) # 增加日志输出请求详情 logger.info(f"Request URL: {req}") logger.info(f"Response status code: {response.status_code}")
4. 切换元数据库提升稳定性
当前使用的SQLite不支持多线程/多进程,任务执行时容易导致调度器阻塞。建议切换到PostgreSQL或MySQL,从根源上解决Airflow的稳定性问题。
内容的提问来源于stack exchange,提问作者yngstgy
相关产品推荐
相关产品推荐

