uranusjr commented on code in PR #73970:
URL: https://github.com/apache/airflow/pull/73970#discussion_r4164119039


##########
task-sdk/src/airflow/sdk/execution_time/coordinator.py:
##########
@@ -262,6 +269,15 @@ 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}")
+        for key, spec in coordinator_specs.items():
+            bundle_name = spec.kwargs.get("task_handler_bundle_name")
+            if bundle_name is None:
+                continue
+            if not isinstance(bundle_name, str) or not 
DagBundlesManager.is_bundle_configured(bundle_name):
+                raise InvalidCoordinatorError(
+                    f"[sdk] coordinators {key!r} sets 
task_handler_bundle_name={bundle_name!r}, "
+                    f"which is not a bundle in [dag_processor] 
dag_bundle_config_list"
+                )

Review Comment:
   This loop validates `task_handler_bundle_name` for every configured 
coordinator, whereas the `queue_to_coordinator` check just above only validates 
keys that are actually routed to. Should configured coordinators that are not 
used by `queue_to_coordinator` simply be ignored? (If not, I’d suggest moving 
this check _before_ checking `queue_to_coordinator`.)



-- 
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