jason810496 commented on code in PR #74042:
URL: https://github.com/apache/airflow/pull/74042#discussion_r4163179706
##########
task-sdk/src/airflow/sdk/execution_time/coordinator.py:
##########
@@ -262,8 +283,43 @@ def from_config(cls) -> Self:
for key in queue_to_coordinator.values():
if key not in coordinator_specs:
raise ValueError(f"[sdk] queue_to_coordinator references
invalid coordinator key: {key!r}")
+ cls._check_dag_file_claims(coordinator_specs)
Review Comment:
Done in 266131caf2 and e610aa7648. The claim check now runs in `for_bundle`,
so `for_queue` never runs it. A coordinator that fails to import or build, with
any exception, is logged and skipped.
##########
task-sdk/src/airflow/sdk/execution_time/coordinator.py:
##########
@@ -107,6 +101,33 @@ def execute_task(
"""
raise NotImplementedError
+ @classmethod
+ def get_dag_importer_class(cls) -> type[AbstractDagImporter] | None:
Review Comment:
Done in 13d82cc57e and e610aa7648. `get_dag_importer()` is now the only
hook: it returns a coordinator's Dag importer, or `None` by default.
`serves_bundle`, `get_parsed_bundles` and `get_dag_importer_class` are gone,
and `for_bundle` reads `dag_bundle_name` from the coordinator's config.
ADR-0010 is updated to match.
##########
task-sdk/src/airflow/sdk/importers/base.py:
##########
@@ -429,6 +437,31 @@ def from_config(cls, bundle_name: str | None = None) ->
Self:
return registry
+ def _register_coordinator_importers(self, bundle_name: str) -> None:
+ """
+ Register the Dag importers of the coordinators that parse Dag files in
*bundle_name*.
+
+ A coordinator configuration that cannot be loaded registers no
coordinator importers, so
+ the bundle's other importers keep working. Two coordinators that claim
the same extension
+ are a configuration error, which is raised.
+ """
+ from airflow.sdk.execution_time.coordinator import
InvalidCoordinatorError, get_coordinator_manager
Review Comment:
Done in 2cb5c904df. The imports in `importers/base.py` are now at the top.
The one in `coordinator.py` has to stay inline, so it now has a one-line reason.
##########
task-sdk/src/airflow/sdk/execution_time/coordinator.py:
##########
@@ -293,6 +349,24 @@ def for_queue(self, queue: str) -> BaseCoordinator:
log.debug("Coordinator found for queue", coordinator=coordinator,
queue=queue)
return coordinator
+ def for_bundle(self, bundle_name: str) -> dict[str, BaseCoordinator]:
Review Comment:
Done in e610aa7648. `for_bundle` picks coordinators by their
`dag_bundle_name` kwarg and builds only those.
##########
task-sdk/tests/task_sdk/execution_time/test_coordinator.py:
##########
@@ -121,6 +159,57 @@ def test_from_config_empty(self, monkeypatch):
assert manager._coordinator_specs == {}
assert manager._queue_to_coordinator == {}
+ @pytest.mark.parametrize(
Review Comment:
Done in 266131caf2 and e610aa7648. The tests cover every-bundle-then-named
and the reverse, `.JAR` vs `jar` conflicting, and `for_queue` still working
when two coordinators conflict.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]