airflow.providers.celery.executors.celery_kubernetes_executor¶
属性¶
类¶
CeleryKubernetesExecutor 由 CeleryExecutor 和 KubernetesExecutor 组成。 |
模块内容¶
- class airflow.providers.celery.executors.celery_kubernetes_executor.CeleryKubernetesExecutor(celery_executor=None, kubernetes_executor=None)[source]¶
Bases:
airflow.executors.base_executor.BaseExecutorCeleryKubernetesExecutor 由 CeleryExecutor 和 KubernetesExecutor 组成。
它根据任务上定义的队列来选择要使用的 executor。当队列的值为配置文件
[celery_kubernetes_executor] 中的 kubernetes_queue(默认值:kubernetes)时,选择 KubernetesExecutor 来运行任务;否则使用 CeleryExecutor。- property queued_tasks: dict[airflow.models.taskinstancekey.TaskInstanceKey, Any][source]¶
返回来自 Celery 和 Kubernetes executor 的排队任务。
- property running: set[airflow.models.taskinstancekey.TaskInstanceKey][source]¶
返回来自 Celery 和 Kubernetes executor 的运行任务。
- job_id: int | str | None[source]¶
继承自 BaseExecutor 的属性。
由于这实际上不是一个 executor,而是多个 executor 的封装,我们将其实现为属性,以便可以自定义 setter。
- queue_command(task_instance, command, priority=1, queue=None)[source]¶
通过 Celery 或 Kubernetes executor 将命令排队。
- queue_task_instance(task_instance, mark_success=False, ignore_all_deps=False, ignore_depends_on_past=False, wait_for_past_depends_before_skipping=False, ignore_task_deps=False, ignore_ti_state=False, pool=None, cfg_path=None, *, kwargs)[source]¶
通过 Celery 或 Kubernetes executor 将任务实例排队。
- has_task(task_instance)[source]¶
检查任务是否在 Celery 或 Kubernetes executor 中排队或运行。
- 参数:
task_instance (airflow.models.taskinstance.TaskInstance) – 任务实例
- 返回:
如果此 executor 已知该任务则返回 True
- 返回类型:
- try_adopt_task_instances(tis)[source]¶
尝试收养因 SchedulerJob 失败而被遗弃的运行中任务实例。
未被收养的任务将被调度器清除(随后可重新调度)。
- 返回:
未能被收养的 TaskInstance 列表
- 返回类型:
collections.abc.Sequence[airflow.models.taskinstance.TaskInstance]
- revoke_task(*, ti)[source]¶
尝试从执行器中移除任务。
它应尝试确保该任务不再在工作节点上运行,并确保其已从内部数据结构中清除。
它 不应 更改 Airflow 中任务的状态,也不应向事件缓冲区添加任何事件。
它不应抛出任何错误。
- 参数:
ti (airflow.models.taskinstance.TaskInstance) – 要移除的任务实例