jason810496 commented on code in PR #74042:
URL: https://github.com/apache/airflow/pull/74042#discussion_r4181288566
##########
task-sdk/src/airflow/sdk/execution_time/coordinator.py:
##########
@@ -297,17 +294,121 @@ def for_queue(self, queue: str) -> BaseCoordinator:
except KeyError:
log.debug("Queue not configured to a coordinator; defaulting to
Python", queue=queue)
return _build_python_coordinator()
+ if key not in self._coordinator_specs:
+ raise InvalidCoordinatorError(f"Queue {queue!r} configured to
nonexistent coordinator")
+ coordinator = self.get_coordinator(key)
+ log.debug("Coordinator found for queue", coordinator=coordinator,
queue=queue)
+ return coordinator
+
+ def get_coordinator(self, key: str) -> BaseCoordinator:
+ """
+ Return the coordinator configured under *key* in ``[sdk]
coordinators``, building it on first use.
+
+ :raises InvalidCoordinatorError: when *key* is not configured, or its
class cannot be imported or
+ called with its kwargs. Other errors from building the coordinator
propagate.
+ """
+ with contextlib.suppress(KeyError):
+ return self._created_coordinators[key]
try:
- coordinator = self._find_queue(key)
+ spec = self._coordinator_specs[key]
except KeyError:
- raise InvalidCoordinatorError(f"Queue {queue!r} configured to
nonexistent coordinator")
+ raise InvalidCoordinatorError(f"No coordinator {key!r} in [sdk]
coordinators")
+ try:
+ coordinator = import_string(spec.classpath)(**spec.kwargs)
except ImportError:
raise InvalidCoordinatorError(f"Cannot import coordinator {key!r}")
except TypeError:
raise InvalidCoordinatorError(f"Cannot instantiate coordinator
{key!r}")
- log.debug("Coordinator found for queue", coordinator=coordinator,
queue=queue)
+ self._created_coordinators[key] = coordinator
return coordinator
+ def _get_coordinator_class(self, key: str) -> type | None:
+ """
+ Return the class of the coordinator under *key* without building it.
+
+ ``None`` means no class can be found: *key* is not configured, its
classpath cannot be
+ imported or fails while importing, or it is not a class. The reason is
logged once.
+ """
+ with contextlib.suppress(KeyError):
+ return self._coordinator_classes[key]
+ coordinator_class: type | None = None
+ if (spec := self._coordinator_specs.get(key)) is None:
+ log.error("No coordinator in [sdk] coordinators", coordinator=key)
+ else:
+ try:
+ resolved = import_string(spec.classpath)
+ except Exception:
+ log.exception("Cannot import coordinator", coordinator=key,
classpath=spec.classpath)
+ else:
+ if isinstance(resolved, type):
+ coordinator_class = resolved
+ else:
+ log.error(
+ "Coordinator classpath is not a class",
coordinator=key, classpath=spec.classpath
+ )
+ self._coordinator_classes[key] = coordinator_class
+ return coordinator_class
+
+ def get_coordinator_keys_for_class(self, coordinator_classpath: str) ->
list[str]:
+ """
+ Return, in config order, the keys of the coordinators of the class at
*coordinator_classpath*.
+
+ A coordinator is of the class when its class is that class or a
subclass of it. Nothing is built.
+ """
+ target = import_string(coordinator_classpath)
Review Comment:
I think we should let it fail loudly if there's import error since it should
be considered as the bug of Airflow Core itself.
--
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]