Airflow: pass execution_date as param to DataflowPythonOperator

Viewed 345

I'd like to use execution date from Airflow as a parameter of my Dataflow job through using DataflowPythonOperator. Specifically, this job read data from Google BigQuery so that I need to provide the execution_date as a part of the query.

I tried op_kwarg and provide_context but it is seemed only applied to PythonOperator.

It's kind of look like this. In the DAG:

run_dataflow = DataFlowPythonOperator(
        task_id='run_dataflow',
        py_file="/path/to/main.py",
        options=dataflowoptions,
        params = execution_date

In main.py:

query = ('select * from `project_id.dataset.table`'
             'where date = {}')
query = query.format(params)
0 Answers
Related