Repository navigation
Validate a serialized Dag and fill its unset settings from the Airflow config - #74041
Conversation
3570480 to
ad14011
Compare
There was a problem hiding this comment.
Do we need prek run sync-go-sdk-schemas --hook-stage manual ? (I checked the code and not automatically syncing Java and Go schemas seems deliberate because it need a maintainer decision, but the CI doesn't warn about it and its really easy to miss the step).
I'll open a PR so at least the drift is caught so the author is asked to take a look instead of silently failing. Maybe something like #74091
LGTM beside kaxil comment and 1 suggestion as well.
Actually, I just noted down the ideal solution to prevent the schema drift in a proper way in #74047 (comment). TL;DR; having bi-weekly job to notify us there's a drift between the Python source and the vendored one. Yes, I will run the |
0afbeb1 to
c930dd6
Compare
jason810496
left a comment
There was a problem hiding this comment.
Thanks for the review
c930dd6 to
a07d7ad
Compare
DagSerialization.validate_serialized_dag checks the JSON schema, deserializes a copy and looks for a cycle in the task graph. It raises DeserializationError naming the Dag, and for a cycle a task on it. It is meant for a serialized Dag that no Python code built, such as one a Lang-SDK runtime returns, so the Python SDK's cycle check never ran on it. The cycle check runs the shared detect_cycle over the serialized downstream edges.
A Lang-SDK runtime cannot read the Airflow config, so it leaves out the Dag settings that a Python Dag reads from it when unset. The new DagSerialization.fill_config_defaults fills max_active_tasks, max_active_runs, max_consecutive_failed_dag_runs, catchup and disable_bundle_versioning in a serialized Dag from the config. A value the Dag sets is kept.
The schema put additionalProperties on the tasks array, which JSON
Schema ignores, so a task entry like {} passed validation. Each entry
must now be an operator with a string task_id, and
validate_serialized_dag rejects a repeated task id. A schema error
now names the failing field's path.
Also pin both boolean config defaults to different values in the
fill_config_defaults tests, and cover the unwrapped cause of a
deserialization error.
receive() now calls DagSerialization.fill_config_defaults and validate_serialized_dag, so the conformance job runs the same checks as the Dag processor and the config-backed field list lives in one place.
These tests passed the dag_maker proxy of the serialized Dag to LazyDeserializedDAG.from_dag, which serialized each task as its repr. The tasks schema now rejects that, so they serialize the SDK Dag.
validate_serialized_dag now applies the max_active_runs and catchup rules the SDK Dag runs when it is built, to the Dag it deserializes. The cycle check walks that Dag's tasks, so edges from the legacy _downstream_task_ids key and from client_defaults count too. It returns the deserialized Dag.
a07d7ad to
962acf2
Compare
jsonschema 4.26 writes $.dag.tasks[2]['__type'] where 4.23 writes $.dag.tasks[2].__type, so the import error and three tests changed with the installed version. Build the path from the error's absolute path.
a729bf8 to
a7efc7c
Compare
jason810496
left a comment
There was a problem hiding this comment.
Thanks for the review!
Stack (bottom to top): #74041, #74042, #74035, #74043, #74036, #74037, #73845, #73846, #73847
Why
A Lang-SDK runtime returns its Dags already serialized, so Airflow has to check them and fill what the runtime cannot know before storing them.
What changes
DagSerialization.validate_serialized_dag(data)checks a serialized Dag against the JSON schema, rejects duplicate task ids, checksmax_active_runs,catchupand cycles as the SDK does, and returns the loadedSerializedDAG.DagSerialization.fill_config_defaults(data)fills an unsetmax_active_tasks,max_active_runs,max_consecutive_failed_dag_runs,catchupanddisable_bundle_versioningfrom[core] max_active_tasks_per_dag,[core] max_active_runs_per_dag,[core] max_consecutive_failed_dag_runs_per_dag,[scheduler] catchup_by_defaultand[dag_processor] disable_bundle_versioning, as a Python Dag does. A value the Dag sets is kept.airflow-core/adr/lang-sdk/0004-dag-parsing.mdsays Airflow fills an unset field from its config.Was generative AI tooling used to co-author this PR?