我正在尝试从触发的 Airflow SimpleHttpOperator 接收 HTTP 响应代码。我已经看到使用 'lambda' 类型的示例,并且正在通过查看响应正文来这样做,但我希望能够将响应代码传递给一个函数。我当前的代码(其中 90% 来自 example_http_operator):
import json
from datetime import timedelta
from airflow import DAG
from airflow.operators.http_operator import SimpleHttpOperator
from airflow.sensors.http_sensor import HttpSensor
from airflow.utils.dates import days_ago
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'start_date': days_ago(2),
'email': ['me@company.com'],
'email_on_failure': False,
'email_on_retry': False,
'retries': 0,
}
dag = DAG(dag_id='kick_off_java_task', default_args=default_args)
kickoff_task = SimpleHttpOperator(
task_id='kick_off_c2c_java_task',
http_conn_id='test',
method='GET',
endpoint='',
data={ "command": "run" },
response_check=lambda response: True if "Ok Message" in response.text else False,
headers={},
xcom_push=False,
dag=dag
)
Run Code Online (Sandbox Code Playgroud)
根据文档和代码,似乎有一种方法可以让 response_check 指向可调用对象,但我不清楚语法,或者我是否需要朝着完全不同的方向前进,例如利用 xcom。
经过一些试验和错误,解决方案非常简单:
dag = DAG(dag_id='kick_off_java_task', default_args=default_args)
def check(response):
if response == 200:
print("Returning True")
return True
else:
print("Returning False")
return False
kickoff_task = SimpleHttpOperator(
task_id='kick_off_c2c_java_task',
http_conn_id='c2c_test',
method='GET',
endpoint='',
data={ "command": "run" },
response_check=lambda response: True if check(response.status_code) is True else False,
headers={},
xcom_push=False,
dag=dag
)
Run Code Online (Sandbox Code Playgroud)
在 lambda 中使用之前定义了 python 函数“check”,我可以将参数“response.status_code”传递给该函数。
| 归档时间: |
|
| 查看次数: |
2316 次 |
| 最近记录: |