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]

Reply via email to