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]