airflow.providers.snowflake.triggers.snowpark_containers

Classes

SnowparkContainerJobTrigger

Poll a Snowpark Container Services job until it reaches a terminal status.

Module Contents

class airflow.providers.snowflake.triggers.snowpark_containers.SnowparkContainerJobTrigger(job_name, snowflake_conn_id, poll_interval, end_time, execution_deadline=None, database=None, schema=None, role=None, warehouse=None)[source]

Bases: airflow.triggers.base.BaseTrigger

Poll a Snowpark Container Services job until it reaches a terminal status.

Parameters:
  • job_name (str) – name of the submitted job service to poll.

  • snowflake_conn_id (str) – reference to the Snowflake connection id.

  • poll_interval (float) – seconds to sleep between DESCRIBE SERVICE polls.

  • end_time (float) – epoch deadline (time.time() seconds) after which a timeout event is emitted.

  • execution_deadline (float | None) – (Optional) absolute timestamp (in seconds since the epoch) after which the task is considered timed out. (default: None)

  • database (str | None) – (Optional) name of database. (default: None)

  • schema (str | None) – (Optional) name of schema. (default: None)

  • role (str | None) – (Optional) name of role. (default: None)

  • warehouse (str | None) – (Optional) name of warehouse. (default: None)

job_name[source]
snowflake_conn_id[source]
poll_interval[source]
end_time[source]
execution_deadline = None[source]
database = None[source]
schema = None[source]
role = None[source]
warehouse = None[source]
serialize()[source]

Serialize SnowparkContainerJobTrigger arguments and class path.

async run()[source]

Poll the job status and yield exactly one terminal event.

async on_kill()[source]

Drop the job service when a deferred task is killed.

Was this entry helpful?