【发布时间】:2022-10-15 01:06:46
【问题描述】:
我有一些表单的自定义运算符
class DataPreparationOperator(BaseOperator):
template_fields = ['file']
def __init__(self, file, **kwargs):
super().__init__(**kwargs)
self.file = file
...
def execute(self, context):
filename = f'prepared_data_{str(time()*1000).replace(".","_")}.json'
DataDownloader(filename, self.filters()).dataframe_downloader()
return filename
(...)
class DataPreparationOperatorArrivals(DataPreparationOperator):
template_fields = ['file']
def __init__(self, file, **kwargs):
super(DataPreparationOperatorArrivals, self).__init__(file=file, **kwargs)
...
def execute(self, context):
filename = f'prepared_data_{str(time()*1000).replace(".","_")}.json'
DataDownloader(filename, self.charge_change()).dataframe_downloader()
return filename
(...)
运算符基于 BranchPythonOperator 执行,在我的 dag 中如下所示
def choose_data_preparation_operator(**kwargs):
if float(kwargs.get("arrival_factor")) != 1.0:
return ['data_preparation_arrivals_change', 'parameters_constructor']
else:
return ['data_preparation_normal', 'parameters_constructor']
opr_data_preparation_path = BranchPythonOperator(
provide_context=True,
task_id='choose_data_preparation_path',
python_callable=choose_data_preparation_operator,
op_kwargs = {'arrival_factor': '{{ dag_run.conf["arrival_factor"] }}'}
)
opr_data_prep = DataPreparationOperator(
task_id ='data_preparation_normal',
file = 'data.json'
)
opr_data_prep_arr = DataPreparationOperatorArrivals(
task_id ='data_preparation_arrivals_change',
file = 'data.json'
)
如您所见,两个运算符都返回一个文件名,现在我想使用另一个自定义运算符并使用各自的文件名在另一个步骤中调用此文件,例如
opr_parameters_constructor = ParametersConstructor(
task_id ='parameters_constructor',
file = '{{ ti.xcom_pull(task_ids="CHOOSE_THE_CORRECT_TASK_ID") }}',
initial_time = '{{ dag_run.conf.get("initial_time") }}',
final_time = '{{ dag_run.conf.get("final_time") }}',
)
我的问题是,我怎样才能把正确的task_id在 BranchPythonOperator? 中选择,即最后一段代码中的 CHOOSE_THE_CORRECT_TASK_ID 变量。
非常感谢您的帮助:D
【问题讨论】:
标签: python airflow directed-acyclic-graphs