kaxil commented on code in PR #74042:
URL: https://github.com/apache/airflow/pull/74042#discussion_r4159728602
##########
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:
`from_config` is also what `supervise()` hits for every task
(`get_coordinator_manager().for_queue(...)` at supervisor.py:2973, which
re-raises), so a parse-claim conflict now stops every task on every queue from
starting, Python tasks on `default` included, before the task log even exists.
Two task-bundle `JavaCoordinator`s routed by queue to different JDKs are enough
once #74036 opts Java in. The same move turns a coordinator module that raises
anything other than `ImportError` on import, or an importer with a callable
`supported_extensions`, into a failure for every queue instead of its own.
Could the check run on the parse path instead, at the top of `for_bundle` and
still before anything is built, so a conflict surfaces as a Dag processor error
for the bundles involved?
##########
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:
None of the new function-body imports between `importers/base.py` and
`execution_time/coordinator.py` says why it is inline. Neither module imports
the other at load time, so the two here can go to the top; the pair in
coordinator.py would then have to stay lazy, with a comment naming the cycle.
`_normalize_extensions` is also private to this module, so a public name would
be better if coordinator.py keeps using it.
##########
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:
ADR-0010 (lines 41-50) has `get_dag_importer()` always return an importer,
and says a second `None` answer would let it disagree with the registry's own
lookup. This adds that `None` path plus `get_parsed_bundles(kwargs)` and
`serves_bundle()`, held together by a docstring line saying they agree. Could
`for_bundle` filter on `get_parsed_bundles(spec.kwargs)` and drop
`serves_bundle`, so there is one answer to "which bundles does this coordinator
parse"? Either way the ADR and the code should say the same thing.
##########
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:
`for_bundle` builds every configured coordinator to ask `serves_bundle`,
including ones whose class has no importer, and it now runs for every
`DagBag(bundle_name=...)`: each parse child and each Python task's `parse()`. A
coordinator whose package is only installed on its own worker image then logs
`Cannot load coordinator` with a traceback in every other task's log. Filtering
on `get_dag_importer_class()` and `get_parsed_bundles(spec.kwargs)` first, as
`_check_dag_file_claims` does, would build only what can parse the bundle. The
class docstring's "never looked up incurs no startup cost" needs updating
either way.
##########
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:
Two branches of `_check_dag_file_claims` have no test that fails without
them. The `bundles is not None and not bundles` skip only matters when the
other coordinator parses every bundle, but `no-bundle` pairs it with a named
one, so deleting the skip stays green. Every-bundle first and named second is
never tested either; dropping `other_bundles is None` stays green and turns
that config into a `TypeError`. Could the first spec be parametrized too, and
one pair use two different importer classes claiming `.JAR` and `jar`? Every
conflicting pair here uses one classpath twice.
--
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]