【发布时间】:2021-09-30 06:43:27
【问题描述】:
我需要创建一个动态 DAG,默认情况下,该 DAG 变量或 (next_execution_date - 1 天) 期间的每个日期都有一个单独的任务(必须使用 dag 的执行日期)。 我的 dag 的一个例子:
dag_vars = Variable.get("dag_dates", deserialize_json=True) # dag_dates = {"dag_start_dt": "NULL", "dag_end_dt": "NULL"}, but can be different dates
DAG_NAME = "dag_test"
def get_params(vars):
if vars["dag_start_dt"] == "NULL":
start_dt = "{{(next_execution_date - macros.timedelta(days=1))}}"
else:
start_dt = vars["dag_start_dt"]
if vars["dag_end_dt"] == "NULL":
end_dt = "{{ next_execution_date }}"
else:
end_dt = vars["dag_end_dt"]
return start_dt, end_dt
start_dag_params, end_dag_params = get_params(dag_vars)
def get_dag_daterange(start_date, end_date):
for n in range(int((end_date - start_date).days)):
yield start_date + timedelta(n)
dag = DAG(
dag_id=DAG_NAME,
default_args=default_args,
schedule_interval= None,
concurrency=1,
max_active_runs=1,
)
with dag:
start_date, end_date = start_dag_params, end_dag_params
for one_date in get_dag_daterange(start_date, end_date):
task_1 = PostgresOperator(
sql = """CALL test_procedure({l_one_date})""".format(l_one_date=one_date),
task_id = "test_procedure_{l_one_date}".format(l_one_date=str(one_date)),
postgres_conn_id = "xxx",
pool = "pool_test",
dag = dag,
autocommit = True,
)
但我有一个错误“-: 'str' 和 'str' 的操作数类型不受支持”。
我知道原因在宏({{next_execution_date}})中,它是在运行时通过Jinja 解析的,但我不知道如何解决这个问题以及如何在气流DAG 中使用宏作为变量.
我很乐意提供任何帮助。谢谢!
【问题讨论】: