airflow.providers.celery.executors.celery_kubernetes_executor

属性

CommandType

CeleryKubernetesExecutor

CeleryKubernetesExecutor 由 CeleryExecutor 和 KubernetesExecutor 组成。

模块内容

airflow.providers.celery.executors.celery_kubernetes_executor.CommandType[source]
class airflow.providers.celery.executors.celery_kubernetes_executor.CeleryKubernetesExecutor(celery_executor=None, kubernetes_executor=None)[source]

Bases: airflow.executors.base_executor.BaseExecutor

CeleryKubernetesExecutor 由 CeleryExecutor 和 KubernetesExecutor 组成。

它根据任务上定义的队列来选择要使用的 executor。当队列的值为配置文件

[celery_kubernetes_executor] 中的 kubernetes_queue(默认值:kubernetes)时,选择 KubernetesExecutor 来运行任务;否则使用 CeleryExecutor。

supports_ad_hoc_ti_run: bool = True[source]
supports_pickling: bool = True[source]
supports_sentry: bool = False[source]
is_local: bool = False[source]
is_single_threaded: bool = False[source]
is_production: bool = True[source]
serve_logs: bool = False[source]
callback_sink: airflow.callbacks.base_callback_sink.BaseCallbackSink | None = None[source]
property kubernetes_queue: str[source]
celery_executor = None[source]
kubernetes_executor = None[source]
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。

start()[source]

启动 Celery 和 Kubernetes executor。

property slots_available: int[source]

此 executor 实例可接受的新任务数量。

property slots_occupied[source]

此 executor 实例当前管理的任务数量。

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 将任务实例排队。

get_task_log(ti, try_number)[source]

从 Kubernetes executor 获取任务日志。

has_task(task_instance)[source]

检查任务是否在 Celery 或 Kubernetes executor 中排队或运行。

参数:

task_instance (airflow.models.taskinstance.TaskInstance) – 任务实例

返回:

如果此 executor 已知该任务则返回 True

返回类型:

bool

heartbeat()[source]

发送心跳以触发 Celery 和 Kubernetes executor 中的新作业。

get_event_buffer(dag_ids=None)[source]

返回并清空来自 Celery 和 Kubernetes executor 的事件缓冲区。

参数:

dag_ids (list[str] | None) – 要返回事件的 dag_id 列表,如果为 None 则返回全部

返回:

事件字典

返回类型:

dict[airflow.models.taskinstancekey.TaskInstanceKey, airflow.executors.base_executor.EventBufferValueType]

try_adopt_task_instances(tis)[source]

尝试收养因 SchedulerJob 失败而被遗弃的运行中任务实例。

未被收养的任务将被调度器清除(随后可重新调度)。

返回:

未能被收养的 TaskInstance 列表

返回类型:

collections.abc.Sequence[airflow.models.taskinstance.TaskInstance]

cleanup_stuck_queued_tasks(tis)[source]
revoke_task(*, ti)[source]

尝试从执行器中移除任务。

它应尝试确保该任务不再在工作节点上运行,并确保其已从内部数据结构中清除。

不应 更改 Airflow 中任务的状态,也不应向事件缓冲区添加任何事件。

它不应抛出任何错误。

参数:

ti (airflow.models.taskinstance.TaskInstance) – 要移除的任务实例

end()[source]

结束 Celery 和 Kubernetes executor。

terminate()[source]

终止 Celery 和 Kubernetes executor。

debug_dump()[source]

调试转储;由调度器在收到 SIGUSR2 信号时调用。

send_callback(request)[source]

发送执行回调。

参数:

request (airflow.callbacks.callback_requests.CallbackRequest) – 要执行的回调请求。

static get_cli_commands()[source]

提供要包含在 Airflow CLI 中的命令。

重写此方法,以通过 Airflow CLI 暴露管理此执行器的命令。这些命令可以用于设置/拆除执行器、检查状态等。请确保为这些命令选择唯一名称,以避免冲突。

此条目是否有帮助?