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)