jason810496 commented on code in PR #73970:
URL: https://github.com/apache/airflow/pull/73970#discussion_r4165331797
##########
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:
Agreed, I consolidate both loop and validate them by getting the unique
coordinators from `queue_to_coordinator` first in c705540.
--
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]