airflow.providers.edge3.executors.edge_executor

Attributes

CommandType

Classes

EdgeExecutor

Implementation of the EdgeExecutor to distribute work to Edge Workers via HTTP.

Module Contents

airflow.providers.edge3.executors.edge_executor.CommandType[source]
class airflow.providers.edge3.executors.edge_executor.EdgeExecutor(*args, **kwargs)[source]

Bases: airflow.executors.base_executor.BaseExecutor

Implementation of the EdgeExecutor to distribute work to Edge Workers via HTTP.

supports_multi_team: bool = True[source]
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.

sync(*, session=NEW_SESSION)[source]

Sync will get called periodically by the heartbeat method.

end()[source]

End the executor.

terminate()[source]

Terminate the executor is not doing anything.

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 running set 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.

Was this entry helpful?