From 8dbc76e5ee769189aefdc28ad407fd36f6f7f02c Mon Sep 17 00:00:00 2001 From: Florian Salfenmoser Date: Wed, 27 May 2026 13:29:59 +0200 Subject: [PATCH 1/6] =?UTF-8?q?Implemented=20Dag-level=20outles=20in=20Tas?= =?UTF-8?q?k=20SDK=20based=20on=20issue=20#39105,=20and=20kept=20the=20API?= =?UTF-8?q?=20consistent=20with=20how=20BaseOperator=20handles=20outlets.?= =?UTF-8?q?=20What=20changed=20In=20task-sdk/src/airflow/sdk/definitions/d?= =?UTF-8?q?ag.py:=201.=20Added=20outlets=20to=20DAG=20init=20=E2=97=A6=20N?= =?UTF-8?q?ew=20field:=20outlets:=20list[Any]=20=E2=97=A6=20Handles=20both?= =?UTF-8?q?=20a=20single=20outlet=20and=20a=20colection.=202.=20Added=20op?= =?UTF-8?q?erator-style=20outlet=20composition=20on=20Dag=20=E2=97=A6=20DA?= =?UTF-8?q?G.=5F=5Fgt=5F=5F=20now=20supports=20dag=20>=20outlet=20?= =?UTF-8?q?=E2=97=A6=20DAG.add=5Foutlets(...)=20appends=20validated=20outl?= =?UTF-8?q?ets=203.=20Wired=20Dag=20success=20to=20outlet=20emission=20?= =?UTF-8?q?=E2=97=A6=20In=20=5F=5Fattrs=5Fpost=5Finit=5F=5F,=20if=20Dag=20?= =?UTF-8?q?outlets=20are=20set,=20an=20outlet-emission=20callback=20is=20a?= =?UTF-8?q?ppended=20to=20on=5Fsuccess=5Fcallback=20=E2=97=A6=20Existing?= =?UTF-8?q?=20user=20callbacks=20are=20preserved=20=E2=97=A6=20has=5Fon=5F?= =?UTF-8?q?success=5Fcallback=20is=20updated=204.=20Added=20the=20callback?= =?UTF-8?q?=20implementation=20=E2=97=A6=20Emits=20asset=20events=20when?= =?UTF-8?q?=20a=20Dag=20run=20succeeds=20=E2=97=A6=20Only=20emits=20concre?= =?UTF-8?q?te=20Asset=20outlets=20=E2=97=A6=20Includes=20Dag/run=20source?= =?UTF-8?q?=20metadata=20=E2=97=A6=20Creates=20missing=20asset=20models=20?= =?UTF-8?q?before=20retrying=20event=20registration=20(same=20pattern=20as?= =?UTF-8?q?=20task=20outlet=20handling)=20In=20airflow-core/src/airflow/as?= =?UTF-8?q?sets/manager.py:=205.=20Extended=20register=5Fasset=5Fchange(..?= =?UTF-8?q?.)=20so=20Dag-level=20emitters=20can=20provide=20source=20field?= =?UTF-8?q?s=20directly=20=E2=97=A6=20Added=20optional=20args:=20=E2=96=AA?= =?UTF-8?q?=20source=5Fdag=5Fid=20=E2=96=AA=20source=5Ftask=5Fid=20?= =?UTF-8?q?=E2=96=AA=20source=5Frun=5Fid=20=E2=96=AA=20source=5Fmap=5Finde?= =?UTF-8?q?x=20=E2=97=A6=20Task-instance=20flow=20is=20unchanged=20?= =?UTF-8?q?=E2=97=A6=20This=20lets=20Dag-level=20events=20carry=20proper?= =?UTF-8?q?=20provenance=20without=20pretending=20to=20come=20from=20a=20t?= =?UTF-8?q?ask=20In=20task-sdk/tests/task=5Fsdk/definitions/test=5Fdag.py:?= =?UTF-8?q?=206.=20Added=20tests=20for:=20=E2=97=A6=20outlets=20init=20fro?= =?UTF-8?q?m=20single=20value=20=E2=97=A6=20outlets=20init=20from=20collec?= =?UTF-8?q?tion=20=E2=97=A6=20dag=20>=20outlet=20=E2=97=A6=20rejecting=20i?= =?UTF-8?q?nvalid=20outlet=20input=20=E2=97=A6=20auto-registration=20of=20?= =?UTF-8?q?Dag=20success=20outlet=20callback=20=E2=97=A6=20preserving=20ex?= =?UTF-8?q?isting=20success=20callback=20when=20outlets=20are=20enabled=20?= =?UTF-8?q?Why=20The=20issue=20asks=20for=20Dag-level=20outlets=20with=20t?= =?UTF-8?q?he=20same=20practical=20effect=20as=20operator=20outlets,=20but?= =?UTF-8?q?=20triggered=20when=20the=20Dag=20run=20succeeds.=20This=20chan?= =?UTF-8?q?ge=20gives=20you=20that=20behavior=20while=20keeping=20the=20Da?= =?UTF-8?q?g=20API=20familiar=20(outlets=3D...,=20dag=20>=20outlet)=20and?= =?UTF-8?q?=20making=20sure=20emitted=20events=20are=20traceable=20back=20?= =?UTF-8?q?to=20the=20Dag=20run.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- airflow-core/src/airflow/assets/manager.py | 11 +++ task-sdk/src/airflow/sdk/definitions/dag.py | 83 +++++++++++++++++++ .../tests/task_sdk/definitions/test_dag.py | 50 +++++++++++ 3 files changed, 144 insertions(+) diff --git a/airflow-core/src/airflow/assets/manager.py b/airflow-core/src/airflow/assets/manager.py index 7c12c31f979d6..d6e324c238a16 100644 --- a/airflow-core/src/airflow/assets/manager.py +++ b/airflow-core/src/airflow/assets/manager.py @@ -251,6 +251,10 @@ def register_asset_change( asset: SerializedAsset | AssetModel | SerializedAssetUniqueKey, extra=None, source_alias_names: Collection[str] = (), + source_dag_id: str | None = None, + source_task_id: str | None = None, + source_run_id: str | None = None, + source_map_index: int | None = None, session: Session, partition_key: str | None = None, source_is_api: bool = False, @@ -309,6 +313,13 @@ def register_asset_change( source_run_id=task_instance.run_id, source_map_index=task_instance.map_index, ) + elif source_dag_id is not None or source_run_id is not None or source_task_id is not None: + event_kwargs.update( + source_task_id=source_task_id, + source_dag_id=source_dag_id, + source_run_id=source_run_id, + source_map_index=source_map_index, + ) asset_event = AssetEvent(**event_kwargs) session.add(asset_event) diff --git a/task-sdk/src/airflow/sdk/definitions/dag.py b/task-sdk/src/airflow/sdk/definitions/dag.py index fd47de467fea4..126a876510c24 100644 --- a/task-sdk/src/airflow/sdk/definitions/dag.py +++ b/task-sdk/src/airflow/sdk/definitions/dag.py @@ -233,6 +233,60 @@ def _convert_deadline(deadline: list[DeadlineAlert] | DeadlineAlert | None) -> l return list(deadline) +def _collect_from_input(value_or_values: Any | Collection[Any] | None) -> list[Any]: + if not value_or_values: + return [] + if isinstance(value_or_values, Collection) and not isinstance(value_or_values, str): + return list(value_or_values) + return [value_or_values] + + +def _emit_dag_asset_events(*, context: Context) -> None: + """Emit asset events for Dag-level outlets on successful Dag run.""" + dag = context.get("dag") + if dag is None: + return + + from airflow.assets.manager import asset_manager + from airflow.models.asset import AssetModel + from airflow.serialization.definitions.assets import SerializedAsset + from airflow.utils.session import create_session + + run_id = context.get("run_id") + with create_session() as session: + for outlet in dag.outlets: + if not isinstance(outlet, BaseAsset): + continue + if not isinstance(outlet, Asset): + # Dag-level outlets only emit concrete assets as events. + continue + + serialized_asset = SerializedAsset( + name=outlet.name, + uri=outlet.uri, + group=outlet.group, + extra=outlet.extra, + watchers=[], + ) + event = asset_manager.register_asset_change( + asset=serialized_asset, + source_dag_id=dag.dag_id, + source_run_id=run_id, + partition_key=getattr(context.get("dag_run"), "partition_key", None), + session=session, + ) + if event is None: + session.add(AssetModel.from_serialized(serialized_asset)) + session.flush() + asset_manager.register_asset_change( + asset=serialized_asset, + source_dag_id=dag.dag_id, + source_run_id=run_id, + partition_key=getattr(context.get("dag_run"), "partition_key", None), + session=session, + ) + + def _convert_doc_md(doc_md: str | None) -> str | None: if doc_md is None: return doc_md @@ -410,6 +464,7 @@ class DAG: :param owner_links: Dict of owners and their links, that will be clickable on the Dags view UI. Can be used as an HTTP link (for example the link to your Slack channel), or a mailto link. e.g: ``{"dag_owner": "https://airflow.apache.org/"}`` + :param outlets: List of outlets that the Dag should emit when the Dag run is successful. :param auto_register: Automatically register this DAG when it is used in a ``with`` block :param fail_fast: Fails currently running tasks when task in Dag fails. **Warning**: A fail stop dag can only have tasks with the default trigger rule ("all_success"). @@ -524,6 +579,7 @@ def __rich_repr__(self): render_template_as_native_obj: bool = attrs.field(default=False, converter=bool) tags: MutableSet[str] = attrs.field(factory=set, converter=_convert_tags) owner_links: dict[str, str] = attrs.field(factory=dict) + outlets: list[Any] = attrs.field(factory=list, converter=_collect_from_input) auto_register: bool = attrs.field(default=True, converter=bool) fail_fast: bool = attrs.field(default=False, converter=bool) allowed_run_types: DagRunType | Collection[DagRunType] | None = attrs.field( @@ -590,6 +646,12 @@ def __attrs_post_init__(self): f"requires max_active_runs <= {active_runs_limit}" ) + if self.outlets: + callbacks = _collect_from_input(self.on_success_callback) + callbacks.append(_emit_dag_asset_events) + self.on_success_callback = callbacks + self.has_on_success_callback = True + @params.validator def _validate_params(self, _, params: ParamsDict): """ @@ -735,6 +797,26 @@ def __hash__(self): hash_components.append(repr(val)) return hash(tuple(hash_components)) + def __gt__(self, other): + """ + Return [Dag] > [Outlet]. + + If other is an attr annotated object it is set as an outlet of this Dag. + """ + if not isinstance(other, Iterable): + other = [other] + + for obj in other: + if not attrs.has(obj): + raise TypeError(f"Left hand side ({obj}) is not an outlet") + self.add_outlets(other) + + return self + + def add_outlets(self, outlets: Iterable[Any]) -> None: + """Define the outlets of this Dag.""" + self.outlets.extend(outlets) + def __enter__(self) -> Self: from airflow.sdk.definitions._internal.contextmanager import DagContext @@ -1610,6 +1692,7 @@ def dag( render_template_as_native_obj: bool = False, tags: Collection[str] | None = None, owner_links: dict[str, str] | None = None, + outlets: Any | None = None, auto_register: bool = True, fail_fast: bool = False, allowed_run_types: DagRunType | Collection[DagRunType] | None = None, diff --git a/task-sdk/tests/task_sdk/definitions/test_dag.py b/task-sdk/tests/task_sdk/definitions/test_dag.py index c42eb8dfc80a1..787bd9df4c1f7 100644 --- a/task-sdk/tests/task_sdk/definitions/test_dag.py +++ b/task-sdk/tests/task_sdk/definitions/test_dag.py @@ -29,6 +29,7 @@ from airflow.sdk.bases.timetable import BaseTimetable from airflow.sdk.definitions.dag import DAG, dag as dag_decorator from airflow.sdk.definitions.param import DagParam, Param, ParamsDict +from airflow.sdk.definitions.asset import Asset from airflow.sdk.definitions.timetables import assets, events, interval, simple, trigger # noqa: F401 from airflow.sdk.exceptions import AirflowDagCycleException, DuplicateTaskIdFound, RemovedInAirflow4Warning from airflow.utils.types import DagRunType @@ -37,6 +38,55 @@ class TestDag: + def test_dag_outlets_init_from_single_value(self): + outlet = Asset("asset://dag_outlet") + dag = DAG("dag-with-outlet", schedule=None, outlets=outlet) + + assert dag.outlets == [outlet] + + def test_dag_outlets_init_from_collection(self): + outlet_1 = Asset("asset://dag_outlet_1") + outlet_2 = Asset("asset://dag_outlet_2") + dag = DAG("dag-with-outlets", schedule=None, outlets=[outlet_1, outlet_2]) + + assert dag.outlets == [outlet_1, outlet_2] + + def test_dag_gt_sets_outlets(self): + dag = DAG("dag-gt-outlet", schedule=None) + outlet = Asset("asset://dag_outlet") + + result = dag > outlet + + assert result is dag + assert dag.outlets == [outlet] + + def test_dag_gt_rejects_non_outlets(self): + dag = DAG("dag-invalid-outlet", schedule=None) + + with pytest.raises(TypeError, match=r"Left hand side \(not-an-outlet\) is not an outlet"): + dag > "not-an-outlet" + + def test_dag_outlets_register_success_callback(self): + dag = DAG("dag-callback-outlet", schedule=None, outlets=Asset("asset://dag_outlet")) + + assert dag.has_on_success_callback + callbacks = dag.on_success_callback if isinstance(dag.on_success_callback, list) else [dag.on_success_callback] + assert any(getattr(cb, "__name__", "") == "_emit_dag_asset_events" for cb in callbacks) + + def test_dag_outlets_preserve_existing_success_callback(self): + def callback(context): + pass + + dag = DAG( + "dag-callback-merge", + schedule=None, + outlets=Asset("asset://dag_outlet"), + on_success_callback=callback, + ) + callbacks = dag.on_success_callback if isinstance(dag.on_success_callback, list) else [dag.on_success_callback] + assert callback in callbacks + assert any(getattr(cb, "__name__", "") == "_emit_dag_asset_events" for cb in callbacks) + @pytest.mark.parametrize( ("dag_id", "exc_type", "exc_value"), [ From eb2d18bff2e1ecbf6c1673956d4da5a7164fc8e8 Mon Sep 17 00:00:00 2001 From: "F.S" Date: Wed, 27 May 2026 14:59:08 +0200 Subject: [PATCH 2/6] =?UTF-8?q?Fix=20Dag=20outlet=20validation=20for=20str?= =?UTF-8?q?ing=20operands=20Treat=20string=20values=20passed=20through=20d?= =?UTF-8?q?ag=20>=20outlet=20as=20a=20single=20invalid=20outlet=20instead?= =?UTF-8?q?=20of=20iterating=20over=20characters.=20Import=20Asset=20where?= =?UTF-8?q?=20Dag-level=20outlet=20event=20emission=20checks=20for=20concr?= =?UTF-8?q?ete=20assets,=20and=20update=20the=20invalid-outlet=20test=20to?= =?UTF-8?q?=20exercise=20DAG.=5F=5Fgt=5F=5F=20without=20triggering=20ruff?= =?UTF-8?q?=E2=80=99s=20pointless-comparison=20rule.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- task-sdk/src/airflow/sdk/definitions/dag.py | 6 ++++-- task-sdk/tests/task_sdk/definitions/test_dag.py | 17 +++++++++++++---- 2 files changed, 17 insertions(+), 6 deletions(-) diff --git a/task-sdk/src/airflow/sdk/definitions/dag.py b/task-sdk/src/airflow/sdk/definitions/dag.py index 126a876510c24..a3dbe05f9c8e9 100644 --- a/task-sdk/src/airflow/sdk/definitions/dag.py +++ b/task-sdk/src/airflow/sdk/definitions/dag.py @@ -45,7 +45,7 @@ from airflow.sdk.bases.timetable import BaseTimetable from airflow.sdk.definitions._internal.node import DAGNode, validate_key from airflow.sdk.definitions._internal.types import NOTSET, ArgNotSet, is_arg_set -from airflow.sdk.definitions.asset import AssetAll, BaseAsset +from airflow.sdk.definitions.asset import Asset, AssetAll, BaseAsset from airflow.sdk.definitions.context import Context from airflow.sdk.definitions.deadline import DeadlineAlert from airflow.sdk.definitions.param import DagParam, ParamsDict @@ -803,8 +803,10 @@ def __gt__(self, other): If other is an attr annotated object it is set as an outlet of this Dag. """ - if not isinstance(other, Iterable): + if isinstance(other, str) or not isinstance(other, Iterable): other = [other] + else: + other = list(other) for obj in other: if not attrs.has(obj): diff --git a/task-sdk/tests/task_sdk/definitions/test_dag.py b/task-sdk/tests/task_sdk/definitions/test_dag.py index 787bd9df4c1f7..edd5af5f00f28 100644 --- a/task-sdk/tests/task_sdk/definitions/test_dag.py +++ b/task-sdk/tests/task_sdk/definitions/test_dag.py @@ -16,6 +16,7 @@ # under the License. from __future__ import annotations +import operator import re import warnings import weakref @@ -27,9 +28,9 @@ from airflow.sdk import Context, Label, PartitionAtRuntime, TaskGroup from airflow.sdk.bases.operator import BaseOperator from airflow.sdk.bases.timetable import BaseTimetable +from airflow.sdk.definitions.asset import Asset from airflow.sdk.definitions.dag import DAG, dag as dag_decorator from airflow.sdk.definitions.param import DagParam, Param, ParamsDict -from airflow.sdk.definitions.asset import Asset from airflow.sdk.definitions.timetables import assets, events, interval, simple, trigger # noqa: F401 from airflow.sdk.exceptions import AirflowDagCycleException, DuplicateTaskIdFound, RemovedInAirflow4Warning from airflow.utils.types import DagRunType @@ -64,13 +65,17 @@ def test_dag_gt_rejects_non_outlets(self): dag = DAG("dag-invalid-outlet", schedule=None) with pytest.raises(TypeError, match=r"Left hand side \(not-an-outlet\) is not an outlet"): - dag > "not-an-outlet" + operator.gt(dag, "not-an-outlet") def test_dag_outlets_register_success_callback(self): dag = DAG("dag-callback-outlet", schedule=None, outlets=Asset("asset://dag_outlet")) assert dag.has_on_success_callback - callbacks = dag.on_success_callback if isinstance(dag.on_success_callback, list) else [dag.on_success_callback] + callbacks = ( + dag.on_success_callback + if isinstance(dag.on_success_callback, list) + else [dag.on_success_callback] + ) assert any(getattr(cb, "__name__", "") == "_emit_dag_asset_events" for cb in callbacks) def test_dag_outlets_preserve_existing_success_callback(self): @@ -83,7 +88,11 @@ def callback(context): outlets=Asset("asset://dag_outlet"), on_success_callback=callback, ) - callbacks = dag.on_success_callback if isinstance(dag.on_success_callback, list) else [dag.on_success_callback] + callbacks = ( + dag.on_success_callback + if isinstance(dag.on_success_callback, list) + else [dag.on_success_callback] + ) assert callback in callbacks assert any(getattr(cb, "__name__", "") == "_emit_dag_asset_events" for cb in callbacks) From db8bf037ed23f9d631a5a1ebffd80af88ab0bd1d Mon Sep 17 00:00:00 2001 From: "F.S" Date: Wed, 10 Jun 2026 16:29:19 +0200 Subject: [PATCH 3/6] Fix type errors in DAG.test() in Task SDK --- task-sdk/src/airflow/sdk/definitions/dag.py | 32 ++++++++++++--------- 1 file changed, 19 insertions(+), 13 deletions(-) diff --git a/task-sdk/src/airflow/sdk/definitions/dag.py b/task-sdk/src/airflow/sdk/definitions/dag.py index a3dbe05f9c8e9..829b9c96df80b 100644 --- a/task-sdk/src/airflow/sdk/definitions/dag.py +++ b/task-sdk/src/airflow/sdk/definitions/dag.py @@ -1360,25 +1360,31 @@ def test( scheduler_dag = DagSerialization.deserialize_dag(DagSerialization.serialize_dag(self)) # Allow users to explicitly pass None. If it isn't set, we default to current time. - logical_date = logical_date if is_arg_set(logical_date) else timezone.utcnow() + logical_date_val: datetime | None = ( + logical_date if is_arg_set(logical_date) else timezone.utcnow() + ) - log.debug("Clearing existing task instances for logical date %s", logical_date) + log.debug("Clearing existing task instances for logical date %s", logical_date_val) # TODO: Replace with calling client.dag_run.clear in Execution API at some point SerializedDAG.clear_dags( dags=[scheduler_dag], - start_date=logical_date, - end_date=logical_date, + start_date=logical_date_val, + end_date=logical_date_val, dag_run_state=False, ) log.debug("Getting dagrun for dag %s", self.dag_id) - logical_date = timezone.coerce_datetime(logical_date) - run_after = timezone.coerce_datetime(run_after) or timezone.coerce_datetime(timezone.utcnow()) - if logical_date is None: + logical_date_val = timezone.coerce_datetime(logical_date_val) + run_after_val: datetime = timezone.coerce_datetime(run_after) or timezone.coerce_datetime( + timezone.utcnow() + ) + if logical_date_val is None: data_interval: DataInterval | None = None else: timetable = coerce_to_core_timetable(self.timetable) - data_interval = timetable.infer_manual_data_interval(run_after=logical_date) + # logical_date_val is not None here, but mypy might still be unsure about its type + # We cast to Any because the core Timetable expects pendulum.DateTime which is a subclass of datetime + data_interval = timetable.infer_manual_data_interval(run_after=cast(Any, logical_date_val)) # These imports are intentionally lazy: this Task SDK module must not # pull in airflow-core at import time (worker isolation). from airflow.dag_processing.bundles.manager import DagBundlesManager @@ -1441,14 +1447,14 @@ def test( dr: DagRun = get_or_create_dagrun( dag=scheduler_dag, - start_date=logical_date or run_after, - logical_date=logical_date, + start_date=logical_date_val or run_after_val, + logical_date=logical_date_val, data_interval=data_interval, - run_after=run_after, + run_after=run_after_val, run_id=DagRun.generate_run_id( run_type=DagRunType.MANUAL, - logical_date=logical_date, - run_after=run_after, + logical_date=logical_date_val, + run_after=run_after_val, ), session=session, conf=run_conf, From 724784bc716748d041ed83ca579be8a6b5376ecb Mon Sep 17 00:00:00 2001 From: "F.S" Date: Mon, 13 Jul 2026 10:59:43 +0200 Subject: [PATCH 4/6] Fix serialized Dag outlets compatibility --- airflow-core/src/airflow/serialization/definitions/dag.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/airflow-core/src/airflow/serialization/definitions/dag.py b/airflow-core/src/airflow/serialization/definitions/dag.py index 33e48aa72bdd1..a28d89c819540 100644 --- a/airflow-core/src/airflow/serialization/definitions/dag.py +++ b/airflow-core/src/airflow/serialization/definitions/dag.py @@ -116,7 +116,9 @@ class SerializedDAG: max_active_runs: int = 16 max_active_tasks: int = 16 max_consecutive_failed_dag_runs: int = 0 + inlets: Sequence[Any] = () owner_links: dict[str, str] = attrs.field(factory=dict) + outlets: Sequence[Any] = () params: SerializedParamsDict = attrs.field(factory=SerializedParamsDict) partial: bool = False render_template_as_native_obj: bool = False From c6ffe3b07cb90c58d8302c673a71076e94cfa122 Mon Sep 17 00:00:00 2001 From: "F.S" Date: Mon, 13 Jul 2026 13:36:42 +0200 Subject: [PATCH 5/6] Fix serialized Dag outlet defaults --- airflow-core/src/airflow/serialization/definitions/dag.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/airflow-core/src/airflow/serialization/definitions/dag.py b/airflow-core/src/airflow/serialization/definitions/dag.py index a28d89c819540..93850c3abceba 100644 --- a/airflow-core/src/airflow/serialization/definitions/dag.py +++ b/airflow-core/src/airflow/serialization/definitions/dag.py @@ -116,9 +116,9 @@ class SerializedDAG: max_active_runs: int = 16 max_active_tasks: int = 16 max_consecutive_failed_dag_runs: int = 0 - inlets: Sequence[Any] = () + inlets: Sequence[Any] = attrs.field(factory=list) owner_links: dict[str, str] = attrs.field(factory=dict) - outlets: Sequence[Any] = () + outlets: Sequence[Any] = attrs.field(factory=list) params: SerializedParamsDict = attrs.field(factory=SerializedParamsDict) partial: bool = False render_template_as_native_obj: bool = False From 7ad33e191df1cc74b48e5e4d4e0e7303faf1c6ae Mon Sep 17 00:00:00 2001 From: "F.S" Date: Mon, 13 Jul 2026 14:46:38 +0200 Subject: [PATCH 6/6] Keep serialized Dag schema in sync --- airflow-core/src/airflow/serialization/definitions/dag.py | 1 + airflow-core/src/airflow/serialization/schema.json | 1 + 2 files changed, 2 insertions(+) diff --git a/airflow-core/src/airflow/serialization/definitions/dag.py b/airflow-core/src/airflow/serialization/definitions/dag.py index 93850c3abceba..e39f62c19bb9d 100644 --- a/airflow-core/src/airflow/serialization/definitions/dag.py +++ b/airflow-core/src/airflow/serialization/definitions/dag.py @@ -167,6 +167,7 @@ def get_serialized_fields(cls) -> frozenset[str]: "max_active_runs", "max_active_tasks", "max_consecutive_failed_dag_runs", + "outlets", "owner_links", "relative_fileloc", "render_template_as_native_obj", diff --git a/airflow-core/src/airflow/serialization/schema.json b/airflow-core/src/airflow/serialization/schema.json index b9efe88448039..955b1162abe35 100644 --- a/airflow-core/src/airflow/serialization/schema.json +++ b/airflow-core/src/airflow/serialization/schema.json @@ -168,6 +168,7 @@ "tasks": { "$ref": "#/definitions/tasks" }, "timezone": { "$ref": "#/definitions/timezone" }, "owner_links": { "type": "object" }, + "outlets": { "type": "array", "default": [] }, "timetable": { "type": "object", "properties": {