Skip to content

Let a coordinator parse the Dag files of the bundles it serves - #74042

Open
jason810496 wants to merge 7 commits into
jason/lang-sdk-e2e/03b-serialized-dag-validationfrom
jason/lang-sdk-e2e/03c-coordinator-dag-parsing
Open

jason810496 wants to merge 7 commits into
jason/lang-sdk-e2e/03b-serialized-dag-validationfrom
jason/lang-sdk-e2e/03c-coordinator-dag-parsing

Conversation

@jason810496

@jason810496 jason810496 commented Oct 1, 2026 •

Copy link
Copy Markdown
Member

Stack (bottom to top): #74040, #74041, #74042, #74035, #74043, #74036, #74037, #73841, #73845, #73846, #73847

Why

To parse native Dags, a coordinator has to say which Dag bundles it parses, and the Dag importer registry has to use its importer for them.

A coordinator's artifacts now always come from a Dag bundle, so the explicit filesystem roots go away.

What changes

  • BaseCoordinator.get_dag_importer() returns the Dag importer that parses a coordinator's native Dag files. It returns None by default, so a coordinator opts in.
  • A coordinator parses the Dag bundle its dag_bundle_name kwarg names, or every bundle when that is unset.
  • DagImporterRegistry.from_config(bundle_name) registers the importers of the coordinators that parse the bundle.
  • CoordinatorManager.for_bundle returns the Dag importers of the coordinators that parse the bundle, and rejects two that claim one extension there. Give each its own dag_bundle_name. Routing a task to a coordinator's queue does not run this check.
  • A Lang-SDK runtime gets the resolved [api] base_url, [operators] default_deferrable and [triggerer] queues_enabled in its environment, like the log levels.

A coordinator opts in by handing out its Dag importer:

class MyCoordinator(SubprocessCoordinator):
    def get_dag_importer(self):
        return MyDagImporter(coordinator=self)  # supported_extensions = [".jar"]
[dag_processor]
dag_bundle_config_list = [
    {"name": "dags-folder", "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", "kwargs": {}},
    {"name": "jar-dags", "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle", "kwargs": {"path": "/opt/airflow/jars"}}
  ]

[sdk]
# Before: "kwargs": {"jars_root": ["/opt/airflow/jars"]}
coordinators = {"mine": {"classpath": "my_pkg.MyCoordinator", "kwargs": {"dag_bundle_name": "jar-dags"}}}

Was generative AI tooling used to co-author this PR?

Comment thread task-sdk/src/airflow/sdk/execution_time/coordinator.py Outdated
Comment thread task-sdk/src/airflow/sdk/execution_time/coordinator.py Outdated
Comment thread task-sdk/src/airflow/sdk/execution_time/coordinator.py Outdated
Comment thread task-sdk/src/airflow/sdk/importers/base.py Outdated
Comment thread task-sdk/tests/task_sdk/execution_time/test_coordinator.py Outdated
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03c-coordinator-dag-parsing branch 2 times, most recently from 107f1eb to e3fc440 Compare October 2, 2026 04:07
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03c-coordinator-dag-parsing branch from e3fc440 to 3018722 Compare October 2, 2026 05:44
@pierrejeambrun
pierrejeambrun force-pushed the jason/lang-sdk-e2e/03c-coordinator-dag-parsing branch from 3018722 to d21aea9 Compare October 2, 2026 09:00
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03c-coordinator-dag-parsing branch from d21aea9 to 42c00d4 Compare October 2, 2026 11:41
@jason810496
jason810496 marked this pull request as ready for review October 2, 2026 12:28

@jason810496 jason810496 left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the reivew.

Comment thread task-sdk/src/airflow/sdk/execution_time/coordinator.py Outdated
Comment thread task-sdk/src/airflow/sdk/execution_time/coordinator.py Outdated
Comment thread task-sdk/src/airflow/sdk/execution_time/coordinator.py Outdated
Comment thread task-sdk/src/airflow/sdk/importers/base.py Outdated
Comment thread task-sdk/tests/task_sdk/execution_time/test_coordinator.py Outdated
jason810496 and others added 7 commits October 2, 2026 23:09
A coordinator opts in to parsing native Dags by returning a Dag importer
from get_dag_importer(). It serves the Dag bundle its dag_bundle_name
names, or every bundle when that is unset. Two coordinators claiming one
extension in a bundle is a config error.

The jars_root, bundles_root and executables_root kwargs are removed. A
coordinator now always finds its artifacts in a Dag bundle: the one
dag_bundle_name names, or the task's own bundle.

SubprocessCoordinator.parse_dag execs the runtime that parses a file.
The runtime inherits only the standard streams and connects back to the
given addresses.

Like the log levels, a runtime now gets the resolved [api] base_url,
[operators] default_deferrable and [triggerer] queues_enabled in its
environment, with the fallbacks Python uses.
importers/base.py and execution_time/coordinator.py do not import
each other at load time, so base.py now imports the coordinator
module at the top. The imports back in coordinator.py stay lazy, with
a comment naming the cycle.

Rename _normalize_extensions to normalize_extensions, since
coordinator.py uses it too, and let reset_coordinator_manager reuse
reset_importer_registry.

Co-Authored-By: Claude <noreply@anthropic.com>
The claim check ran in CoordinatorManager.from_config, which every
task hits through for_queue, so one conflicting pair of coordinators
stopped all tasks from starting. It now runs in for_bundle, before
anything is built, and only for the bundle being parsed.

for_bundle selects the coordinators to build from their specs alone,
so a coordinator that cannot parse the bundle is never built. A class
that fails to import is logged and skipped, and an importer whose
supported_extensions is not a class-level list is rejected.

Co-Authored-By: Claude <noreply@anthropic.com>
for_bundle selects coordinators from get_parsed_bundles, so
serves_bundle was a second answer to the same question that could
disagree with it. Remove it.

get_dag_importer is called only on coordinators whose
get_dag_importer_class is set, so it now always returns an importer
and the base raises NotImplementedError. ADR-0010 describes the
class-level lookups.

Co-Authored-By: Claude <noreply@anthropic.com>
Reword the docs so dag_bundle_name reads as optional with one coordinator
per class, and needed when one class is configured more than once.

Co-Authored-By: Claude <noreply@anthropic.com>
Make get_dag_importer the only hook a coordinator implements to parse
native Dags: it returns the Dag importer, or None, the default, to
parse none. get_dag_importer_class and get_parsed_bundles are removed.

for_bundle selects coordinators by the dag_bundle_name kwarg in their
config, unset meaning every bundle, builds them, and returns their
importers keyed by coordinator. One that cannot be built, or returns
None, is skipped. Two importers claiming one extension in the bundle
still raise.

Rename normalize_extensions back to _normalize_extensions, since the
coordinator module no longer uses it.

Co-Authored-By: Claude <noreply@anthropic.com>
dag_bundle_name lets Dag processing pick which of several coordinators
of one class starts the runtime that parses an artifact. With a single
coordinator it can be omitted. Reword the Java, TypeScript and Go docs
and READMEs to say so, without tying it to the Python stub Dag.

Co-Authored-By: Claude <noreply@anthropic.com>
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03c-coordinator-dag-parsing branch from 63a5cd3 to b830155 Compare October 2, 2026 15:09

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:coordinator Coordinator: The interface to spawn Lang-SDK subprocesses area:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants