【发布时间】: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