【问题标题】:Why does a pythonoperator callable not need to accept parameters in airflow?为什么可调用的pythonoperator不需要接受气流中的参数?
【发布时间】:2022-07-22 11:58:57
【问题描述】:

我不明白可调用对象(由 PythonOperator 指定的函数调用)n Airflow 应该如何设置其参数列表。我见过没有参数或命名参数或 **kwargs。我总是可以添加“ti”或 **allargs 作为参数,并且 ti 似乎用于任务实例信息,或 ds 用于执行日期。但是我的可调用对象显然不需要参数。它们可以简单地是“def function():”。如果我编写了一个常规的 python 函数 func() 而不是 func(**kwargs),除非没有传递任何参数,否则它会在运行时失败。气流似乎总是通过 t1,那么可调用函数签名怎么可能不需要呢?下面的示例来自培训站点,其中 _process_data 函数获取 ti,但 _extract_bitcoin_price() 没有。我在想那是因为 xcom 推送,但似乎 ti 总是可用的,那么“def somefunc()”怎么能工作呢?我尝试查看 pythonoperator 源代码,但我不清楚它是如何工作的或将参数包含在可调用对象中的最佳实践。谢谢!!

from airflow import DAG
from airflow.operators.python_operator 
import PythonOperator

from datetime import datetime
import json
from typing import Dict
import requests
import logging

API = "https://api.coingecko.com/api/v3/simple/price?ids=bitcoin&vs_currencies=usd&include_market_cap=true&include_24hr_vol=true&include_24hr_change=true&include_last_updated_at=true"

def \_extract_bitcoin_price():
    return requests.get(API).json()\['bitcoin'\]

def \_process_data(ti):
    response = ti.xcom_pull(task_ids='extract_bitcoin_price')
    logging.info(response)
    processed_data = {'usd': response\['usd'\], 'change': response\['usd_24h_change'\]}
    ti.xcom_push(key='processed_data', value=processed_data)

def \_store_data(ti):
    data = ti.xcom_pull(task_ids='process_data', key='processed_data')
    logging.info(f"Store: {data\['usd'\]} with change {data\['change'\]}")

with DAG('classic_dag', schedule_interval='@daily', start_date=datetime(2021, 12, 1), catchup=False) as dag:

extract_bitcoin_price = PythonOperator(
    task_id='extract_bitcoin_price',
    python_callable=_extract_bitcoin_price
)

process_data = PythonOperator(
    task_id='process_data',
    python_callable=_process_data
)

store_data = PythonOperator(
    task_id='store_data',
    python_callable=_store_data
)

extract_bitcoin_price >> process_data >> store_data

尝试了没有参数的可调用对象 somefunc() 期望得到错误,说传递了太多参数,但它成功了。添加 somefunc(ti) 也可以!两者如何工作?

【问题讨论】:

    标签: airflow


    【解决方案1】:

    我认为您缺少的是 Airflow 允许将任务的 context 传递给 python 可调用对象(您可以看到其中之一是 ti)。这些是 Airflow 提供的其他有用参数,您可以在任务中使用它们。

    在较旧的 Airflow 版本中,用户必须设置 provide_context=True 才能使用:

     process_data = PythonOperator(
        ...,
        provide_context=True
    )
    

    由于 Airflow>=2.0,因此无需使用 provide_context。气流在引擎盖下处理它。

    当你在 Python Callable 中看到类似的签名时:

    def func(ti, **kwargs):
        ...
    

    这意味着 ti 已从 kwargs 中“解包”。你也可以这样做:

    def func(**kwargs):
        ti = kwargs['ti']
    

    编辑: 我认为您在写作时缺少的是:

    def func()
            ...
    
    store_data = PythonOperator(
        task_id='task',
        python_callable=func
    )
    

    Airflow 不仅仅是调用 func。正在执行的代码是PythonOperator 的execute() 函数,该函数调用您提供的带有args 和kwargs 的python callable

    【讨论】:

    • “由于 Airflow>=2.0,因此无需使用 provide_context。气流句柄在引擎盖下处理它。” Tnanks,这回答了一个我什至没有问但想知道的问题:为什么在没有提供上下文的情况下首先传递 ti。我仍然不明白如果上下文是没有参数的 def func() 这样的签名如何成功总是过去。气流是否将函数签名检查为某种 python 元数据并进行调整?或者也许无参数签名总是失败......将不得不重试并回复。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-09-23
    • 2021-04-29
    • 2014-07-25
    • 2020-06-29
    • 2016-07-27
    相关资源
    最近更新 更多