airflow.providers.edge3.models.edge_job¶
Classes¶
A job which is queued, waiting or running on a Edge Worker. |
Functions¶
|
Build the key the executor layer uses for a job row. |
Module Contents¶
- airflow.providers.edge3.models.edge_job.build_job_key(dag_id, task_id, run_id, try_number, map_index)[source]¶
Build the key the executor layer uses for a job row.
A row is a callback only if it has the full identity
queue_workload()writes for callbacks, sinceExecuteCallbackis a valid Dag id. A task row maps to theairflow.modelsTaskInstanceKey, not theairflow.sdkone, becauseBaseExecutordispatches on it withisinstance.
- class airflow.providers.edge3.models.edge_job.EdgeJobModel(dag_id, task_id, run_id, map_index, try_number, state, queue, concurrency_slots, command, queued_dttm=None, edge_worker=None, last_update=None, team_name=None)[source]¶
Bases:
airflow.providers.edge3.models.edge_base.Base,airflow.utils.log.logging_mixin.LoggingMixinA job which is queued, waiting or running on a Edge Worker.
Each tuple in the database represents and describes the state of one job.
- dag_id: sqlalchemy.orm.Mapped[str][source]¶
- task_id: sqlalchemy.orm.Mapped[str][source]¶
- run_id: sqlalchemy.orm.Mapped[str][source]¶
- map_index: sqlalchemy.orm.Mapped[int][source]¶
- try_number: sqlalchemy.orm.Mapped[int][source]¶
- state: sqlalchemy.orm.Mapped[str][source]¶
- queue: sqlalchemy.orm.Mapped[str][source]¶
- concurrency_slots: sqlalchemy.orm.Mapped[int][source]¶
- command: sqlalchemy.orm.Mapped[str][source]¶
- queued_dttm: sqlalchemy.orm.Mapped[datetime.datetime | None][source]¶
- edge_worker: sqlalchemy.orm.Mapped[str | None][source]¶
- last_update: sqlalchemy.orm.Mapped[datetime.datetime | None][source]¶
- team_name: sqlalchemy.orm.Mapped[str | None][source]¶