airflow.providers.edge3.executors.edge_executor¶
Attributes¶
Classes¶
Implementation of the EdgeExecutor to distribute work to Edge Workers via HTTP. |
Module Contents¶
- class airflow.providers.edge3.executors.edge_executor.EdgeExecutor(*args, **kwargs)[source]¶
Bases:
airflow.executors.base_executor.BaseExecutorImplementation of the EdgeExecutor to distribute work to Edge Workers via HTTP.
- last_reported_state: dict[airflow.models.taskinstancekey.TaskInstanceKey | airflow.models.callback.CallbackKey, airflow.utils.state.TaskInstanceState | str][source]¶
- start(*, session=NEW_SESSION)[source]¶
If EdgeExecutor provider is loaded first time, ensure table exists.
- queue_workload(workload, session)[source]¶
Put new workload to queue. Airflow 3 entry point to execute a task.
- revoke_task(*, ti, session=NEW_SESSION)[source]¶
Revoke a task instance from the executor.
This method removes the task from the executor’s internal state and deletes the corresponding EdgeJobModel record to prevent edge workers from picking it up.
- Parameters:
ti (airflow.models.taskinstance.TaskInstance) – Task instance to revoke
session (sqlalchemy.orm.Session) – Database session
- try_adopt_task_instances(tis, *, session=NEW_SESSION)[source]¶
Adopt the task instances whose job is still in flight in the edge_job table.
The
runningset is empty after a scheduler restart, so the adopted keys go back into it to keep slot accounting accurate. Task instances whose job is finished or missing are returned so the scheduler clears and re-schedules them.- Returns:
any TaskInstances that were unable to be adopted
- Return type:
collections.abc.Sequence[airflow.models.taskinstance.TaskInstance]
- static get_cli_commands()[source]¶
Vends CLI commands to be included in Airflow CLI.
Override this method to expose commands via Airflow CLI to manage this executor. This can be commands to setup/teardown the executor, inspect state, etc. Make sure to choose unique names for those commands, to avoid collisions.