Airflow任务发送HTTP请求时进程挂起问题求助
环境信息
- 系统:MacOS Apple M1(本地机器)
- Airflow版本:2.5.3
- 执行器:Local搭配Postgres数据库
问题描述
我正在实现一个从REST API加载数据的外部触发工作流,使用Python Operator运行代码并通过Airflow UI手动触发流程。但当执行到包含HTTP请求代码的任务时,进程会永久挂起,且笔记本电脑发热严重。


任务文件(tasks/import_logs.py)
import requests def import_logs(**context): print("[Sasha] Running log importer") context["ti"].xcom_push( key="logs", value=["log/location/1", "log/location/2"]) print("Log locations pushed to xcom") # Define the URL for the dummy endpoint url = 'https://jsonplaceholder.typicode.com/posts' # Define the payload for the JSON request payload = { "title": "foo", "body": "bar", "userId": 1 } # Define the headers for the request headers = {'Content-Type': 'application/json'} # Send the POST request to the dummy endpoint response = requests.post(url, json=payload, headers=headers) # Print the response status code and content print(f'Response status code: {response.status_code}') print(f'Response content: {response.content}')
DAG定义文件
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime, timedelta from tasks.import_logs import import_logs from tasks.import_tops import import_tops from tasks.process_input import process_input from tasks.process_log_data import process_logs from tasks.output_logs import output_logs from tasks.cleanup import cleanup from tasks.trigger_data_update import trigger_data_update default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2023, 3, 31), 'email_on_failure': False, 'email_on_retry': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } dag = DAG('process_log', default_args=default_args, schedule_interval=None) validate_input = PythonOperator( task_id='validate_input', python_callable=process_input, provide_context=True, dag=dag ) import_log = PythonOperator( task_id='import_logs', python_callable=import_logs, provide_context=True, dag=dag ) import_top = PythonOperator( task_id='import_tops', python_callable=import_tops, provide_context=True, dag=dag ) process_log = PythonOperator( task_id='process_logs', python_callable=process_logs, provide_context=True, dag=dag ) output_log = PythonOperator( task_id='write_logs', python_callable=output_logs, provide_context=True, dag=dag ) cleanup_task = PythonOperator( task_id='cleanup', python_callable=cleanup, provide_context=True, dag=dag ) update_task = PythonOperator( task_id='trigger_data_update', python_callable=trigger_data_update, provide_context=True, dag=dag ) validate_input >> [import_top, import_log] >> process_log >> output_log >> [cleanup_task, update_task] if __name__ == "__main__": import json with open('test_conf/process_log.json', 'r') as f: conf = json.load(f) dag.test( run_conf=conf )
已尝试的排查动作
- 最初使用Sequential executor时出现问题,切换到Local executor后问题依旧
- 精简代码,只保留模拟HTTP请求,问题仍存在
- 尝试使用Airflow HTTP Hook搭配Connection,结果相同
- 其他任务均为仅打印内容的模拟任务,运行DAG底部的测试代码一切正常
内容的提问来源于stack exchange,提问作者Alexandr Sarioglo
相关产品推荐
相关产品推荐

