如何从Airflow SimpleHttpOperator GET请求访问响应

Rac*_*man 6 airflow data-pipeline apache-airflow

我正在学习Airflow并且有一个简单的问题.下面是我的DAG叫dog_retriever

import airflow
from airflow import DAG
from airflow.operators.http_operator import SimpleHttpOperator
from airflow.operators.sensors import HttpSensor
from datetime import datetime, timedelta
import json



default_args = {
    'owner': 'Loftium',
    'depends_on_past': False,
    'start_date': datetime(2017, 10, 9),
    'email': 'rachel@loftium.com',
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 3,
    'retry_delay': timedelta(minutes=3),
}

dag = DAG('dog_retriever',
    schedule_interval='@once',
    default_args=default_args)

t1 = SimpleHttpOperator(
    task_id='get_labrador',
    method='GET',
    http_conn_id='http_default',
    endpoint='api/breed/labrador/images',
    headers={"Content-Type": "application/json"},
    dag=dag)

t2 = SimpleHttpOperator(
    task_id='get_breeds',
    method='GET',
    http_conn_id='http_default',
    endpoint='api/breeds/list',
    headers={"Content-Type": "application/json"},
    dag=dag)

t2.set_upstream(t1)
Run Code Online (Sandbox Code Playgroud)

作为测试Airflow的一种方法,我只是在这个非常简单的http://dog.ceo API中向一些端点发出两个GET请求.目标是学习如何处理通过Airflow检索的一些数据

执行正在运行 - 我的代码成功调用了任务t1和t2中的enpoints,我可以看到它们按照set_upstream我写的规则以正确的顺序记录在Airflow UI中.

我无法弄清楚如何访问这两个任务的json响应.看起来很简单,但我无法弄清楚.在SimpleHtttpOperator中,我看到了response_check的参数,但没有简单的打印,存储或查看json响应.

谢谢.

Che*_*zhi 11

因此,这是SimpleHttpOperator,实际的json被推送到XCOM,你可以从那里得到它.以下是该操作的代码行:https://github.com/apache/incubator-airflow/blob/master/airflow/operators/http_operator.py#L87

您需要做的是设置xcom_push=True,所以您的第一个t1将是以下内容:

t1 = SimpleHttpOperator(
    task_id='get_labrador',
    method='GET',
    http_conn_id='http_default',
    endpoint='api/breed/labrador/images',
    headers={"Content-Type": "application/json"},
    xcom_push=True,
    dag=dag)
Run Code Online (Sandbox Code Playgroud)

您应该能够return value在XCOM中找到所有JSON,有关XCOM的更多详细信息,请访问:https://airflow.incubator.apache.org/concepts.html#xcoms

  • @成志。你好!您能否分享一下第二个 SimpleHttpOperator 任务“t2”的样子,它可能使用第一个任务中的数据。问题是,我看到无数的例子,它们说 - 只需使用 xcom 并推送数据,但它们没有显示接收器部分或其他任务,这些任务可能使用前一个任务推送的数据。 (4认同)

Dou*_*son 7

我主要为试图(或想要)从流程中调用Airflow 工作流 DAG 并接收由 DAG 活动产生的任何数据的任何人添加此答案。

重要的是要了解运行 DAG 需要 HTTP POST 并且对此 POST的响应在 Airflow 中硬编码,即如果不更改 Airflow 代码本身,Airflow 将永远不会返回任何东西,除了状态代码和消息给请求者过程。

Airflow 似乎主要用于为 ETL(提取、转换、加载)工作流创建数据管道,现有的Airflow Operators,例如 SimpleHttpOperator,可以从 RESTful web 服务中获取数据,处理它,并使用其他运算符将其写入数据库,但是不要在对运行工作流 DAG的 HTTP POST 的响应中返回它

即使操作员确实在响应中返回了此数据,查看 Airflow 源代码也确认 trigger_dag() 方法不会检查或返回它:

apache_airflow_airflow_www_api_experimental_endpoints.py

apache_airflow_airflow_api_client_json_client.py

它返回的只是这个确认消息:

编排服务中收到的气流 DagRun 消息

由于 Airflow 是开源的,我想我们可以修改 trigger_dag() 方法来返回数据,但随后我们会被困在维护分叉的代码库中,而且我们将无法使用云托管的、基于 Airflow 的服务,例如Google Cloud Platform 上的 Cloud Composer,因为它不包括我们的修改。

更糟糕的是,Apache Airflow 甚至没有正确返回其硬编码的状态消息。

当我们POST成功到 Airflow/dags/{DAG-ID}/dag_runs端点时,我们会收到“200 OK”响应,而不是我们应该收到的“201 Created”响应。并且 Airflow 使用其“已创建……”状态消息对响应的内容主体进行“硬编码”。 然而,标准是在响应头中返回新创建的资源的 Uri,而不是在正文中......这将使正文可以自由地返回在此创建期间(或由此创建)产生/聚合的任何数据。

我将这个缺陷归因于“盲目”(或我称之为“天真”)敏捷/MVP 驱动的方法,它只添加被要求的功能,而不是保持了解并为更通用的效用留出空间。由于Airflow 绝大多数用于为(和由)数据科学家(而非软件工程师)创建数据管道,因此Airflow 运营商可以使用其专有的内部 XCom 功能相互共享数据,正如@Chengzhi 的有用回答指出的那样(谢谢!)但在任何情况下都不能将数据返回给请求者启动 DAG,即 SimpleHttpOperator 可以从第三方 RESTful 服务检索数据,并可以与丰富、聚合和/或转换数据的 PythonOperator(通过 XCom)共享该数据。然后 PythonOperator 可以与 PostgresOperator 共享其数据,后者将结果直接存储在数据库中。但是结果永远无法返回到请求完成工作的流程,即我们的 Orchestration 服务,这使得 Airflow 对任何用例都无用,但由其当前用户驱动的用例。

这里的要点(至少对我而言)是 1)永远不要将过多的专业知识归功于任何人或任何组织。Apache 是一个重要的组织,在软件开发方面有着深厚而重要的根基……但它们并不完美。2) 始终提防内部的、专有的解决方案。开放的、基于标准的解决方案已经从许多不同的角度进行了检查和审查,而不仅仅是一个。

我花了将近一周的时间去寻找不同的方法来做看起来非常简单合理的事情。我希望这个答案可以为其他人节省一些时间。