含约10分钟Web Post的Airflow Operator卡顿问题排查
Airflow长耗时POST请求算子卡顿/失败问题
当Airflow Operator中包含耗时约10分钟的Web Post请求时,会出现卡顿现象,以下算子执行失败:
- 包含使用requests库发送HTTP Post请求代码的PythonOperator
- 与上述逻辑相同,但使用PythonVirtualenvOperator(参考相关帖子测试)
- 使用paramiko通过SSH连接主机并执行本地cURL的PythonOperator
以下算子执行成功:
- 调用cURL请求端点的BashOperator,但该方式仅支持单个bash调用,优先希望使用Python实现
各算子代码示例
代码1:使用requests库的PythonOperator
def post_using_requests(): s = requests.session() s.auth = (<auth>) s.headers.update(<headers>) resp = None resp = s.post(<endpoint>, data = <data>) # ----- Operator gets stuck here ----- # Check and process response request_via_requests_lib_operator = PythonOperator( task_id = <task_id>, python_callable = post_using_requests )
代码2:使用requests的PythonVirtualenvOperator
def post_using_requests_virtualenv(): # Additional imports # Everything else the same as #1 request_via_requests_lib_virtualenv_operator = PythonVirtualenvOperator( task_id = <task_id>, python_callable = post_using_requests_virtualenv )
代码3:通过paramiko SSH执行cURL的PythonOperator
def post_using_ssh(): # get ssh_client ssh_client = paramiko.SSHClient() # credentials = <get_creds> ssh_client.set_missing_host_key_policy(paramiko.AutoAddPolicy()) ssh_client.connect(<credentials>) curl_command = "curl -X POST -sS --location <endpoint> --header <headers> --data <data>" _stdin, _stdout, _stderr = ssh_client.exec_command(curl_command) # ----- Operator gets stuck here ----- stderr_str = _stderr.read() stdout_str = _stdout.read() request_via_ssh_operator = PythonOperator( task_id = <task_id>, python_callable = post_using_ssh )
代码4:调用cURL的BashOperator
request_via_bash_cURL= BashOperator( task_id=<task_id>, bash_command="""curl -X POST -sS --location <endpoint> --header <headers> --data <data>""" )
已尝试的解决方法
在/aws-mwaa-local-runner/docker/config/.env.localrunner文件中添加NO_PROXY="*",但问题未解决。当前Airflow环境基于aws-mwaa-local-runner仓库v2.2.2分支搭建。
内容的提问来源于stack exchange,提问作者Maile Cupo
相关产品推荐
相关产品推荐

