airflow.providers.dbt.cloud.operators.dbt

DbtCloudRunJobOperatorLink

允许用户直接在 dbt Cloud 中监控触发的任务运行。

DbtCloudRunJobOperator

执行一个 dbt Cloud 任务。

DbtCloudGetJobRunArtifactOperator

从 dbt Cloud 任务运行中下载产物(Artifacts)。

DbtCloudListJobsOperator

列出 dbt Cloud 项目中的任务。

模块内容

Bases: airflow.providers.common.compat.sdk.BaseOperatorLink

允许用户直接在 dbt Cloud 中监控触发的任务运行。

name = 'Monitor Job Run'[源代码]

链接的名称。这将是任务 UI 上按钮的名称。

指向外部系统的链接。

参数:
  • operator (airflow.providers.common.compat.sdk.BaseOperator) – 与此链接关联的 Airflow 操作符对象。

  • ti_key – 要返回链接的任务实例(TaskInstance)ID。

返回:

指向外部系统的链接

class airflow.providers.dbt.cloud.operators.dbt.DbtCloudRunJobOperator(*, dbt_cloud_conn_id=DbtCloudHook.default_conn_name, job_id=None, project_name=None, environment_name=None, job_name=None, account_id=None, trigger_reason=None, steps_override=None, schema_override=None, wait_for_termination=True, timeout=60 * 60 * 24 * 7, check_interval=60, additional_run_config=None, reuse_existing_run=False, retry_from_failure=False, deferrable=conf.getboolean('operators', 'default_deferrable', fallback=False), hook_params=None, **kwargs)[源代码]

Bases: airflow.providers.common.compat.sdk.BaseOperator

执行一个 dbt Cloud 任务。

另请参阅

有关如何使用此操作符的更多信息,请参阅指南:触发 dbt Cloud 任务

参数:
  • dbt_cloud_conn_id (str) – 连接到 dbt Cloud 的连接 ID。

  • job_id (int | None) – dbt Cloud 任务的 ID。如果未提供 project_name、environment_name 和 job_name,则必须提供此项。

  • project_name (str | None) – 可选。dbt Cloud 项目名称。仅在 job_id 为 None 时使用。

  • environment_name (str | None) – 可选。dbt Cloud 环境名称。仅在 job_id 为 None 时使用。

  • job_name (str | None) – 可选。dbt Cloud 任务名称。仅在 job_id 为 None 时使用。

  • account_id (int | None) – 可选。dbt Cloud 账户 ID。

  • trigger_reason (str | None) – 可选。触发任务原因的描述。默认为“Triggered via Apache Airflow by task <task_id> in the <dag_id> DAG.”

  • steps_override (list[str] | None) – 可选。触发任务时执行的 dbt 命令列表,将覆盖 dbt Cloud 中配置的命令。

  • schema_override (str | None) – 可选。覆盖此任务在已配置目标中的目标 schema。

  • wait_for_termination (bool) – 等待任务运行终止的标志。默认情况下启用此功能,但可以禁用以使用 DbtCloudJobRunSensor 对长时间运行的任务执行异步等待。

  • timeout (int) – 在非异步等待模式下,等待任务运行达到终止状态的超时时间(秒)。仅在 wait_for_termination 为 True 时使用。这限制了操作符等待任务完成的时间,并不意味着会取消任务。任务级别的超时应通过 execution_timeout 强制执行。默认为 7 天。

  • check_interval (int) – 在非异步等待模式下,检查任务运行状态的时间间隔(秒)。仅在 wait_for_termination 为 True 时使用。默认为 60 秒。

  • additional_run_config (dict[str, Any] | None) – 可选。触发任务时 API 请求中应包含的任何附加参数。

  • reuse_existing_run (bool) – 决定是否重用现有非终止状态任务运行的标志。如果设置为 true 且找到了非终止的任务运行,它将使用最新的运行,而不会触发新的任务运行。

  • retry_from_failure (bool) – 决定是否从失败状态重试任务运行的标志。如果设置为 true 且上一次任务运行失败,它将以与失败运行相同的配置触发新的任务运行。有关重试逻辑的更多信息,请参阅:https://docs.getdbt.com/dbt-cloud/api-v2#/operations/Retry%20Failed%20Job

  • deferrable (bool) – 在可延迟模式下运行算子

  • hook_params (dict[str, Any] | None) – 传递给 DbtCloudHook 构造函数的额外参数。

  • execution_timeout – 任务允许运行的最长时间。如果超过此时间,dbt Cloud 任务将被取消,任务将失败。当同时设置了 execution_timeouttimeout 时,较早的截止时间优先。

返回:

触发的 dbt Cloud 任务运行 ID。

template_fields = ('dbt_cloud_conn_id', 'job_id', 'project_name', 'environment_name', 'job_name', 'account_id',...[源代码]
dbt_cloud_conn_id = 'dbt_cloud_default'[源代码]
account_id = None[源代码]
job_id = None[源代码]
project_name = None[源代码]
environment_name = None[源代码]
job_name = None[源代码]
trigger_reason = None[源代码]
steps_override = None[源代码]
schema_override = None[源代码]
wait_for_termination = True[源代码]
timeout = 604800[源代码]
check_interval = 60[源代码]
additional_run_config[源代码]
run_id: int | None = None[源代码]
reuse_existing_run = False[源代码]
retry_from_failure = False[源代码]
deferrable[源代码]
hook_params[源代码]
execute(context)[源代码]

在创建算子时派生。

执行任务的主要方法。Context 是与渲染 jinja 模板时使用的相同字典。

有关更多上下文,请参考 get_template_context。

execute_complete(context, event)[源代码]

当触发器触发时执行 —— 立即返回。

on_kill()[源代码]

当任务实例被终止时,重写此方法以清理子进程。

在算子内部使用 threading、subprocess 或 multiprocessing 模块时,需要进行清理,否则会留下僵尸进程。

property hook[源代码]

返回 DBT Cloud hook。

get_openlineage_facets_on_complete(task_instance)[源代码]

实现 _on_complete,因为 job_run 需要先在 execute 方法中触发。

只有当操作符 wait_for_termination 设置为 True 时,此方法才会发送额外的事件。

class airflow.providers.dbt.cloud.operators.dbt.DbtCloudGetJobRunArtifactOperator(*, dbt_cloud_conn_id=DbtCloudHook.default_conn_name, run_id, path, account_id=None, step=None, output_file_name=None, hook_params=None, **kwargs)[源代码]

Bases: airflow.providers.common.compat.sdk.BaseOperator

从 dbt Cloud 任务运行中下载产物(Artifacts)。

另请参阅

有关如何使用此操作符的更多信息,请参阅指南:下载运行产物

参数:
  • dbt_cloud_conn_id (str) – 连接到 dbt Cloud 的连接 ID。

  • run_id (int) – dbt Cloud 任务运行的 ID。

  • path (str) – 与产物文件相关的文件路径。路径以 target/ 目录为根目录。使用 “manifest.json”、“catalog.json” 或 “run_results.json” 来下载 dbt 生成的运行产物。

  • account_id (int | None) – 可选。dbt Cloud 账户 ID。

  • step (int | None) – 可选。查询产物的运行步骤索引。运行中的第一个步骤索引为 1。如果省略 step 参数,将返回运行中最后一个步骤的产物。

  • output_file_name (str | None) – 可选。下载产物文件所需的名称。默认为 <run_id>_<path>(例如 “728368_run_results.json”)。

  • hook_params (dict[str, Any] | None) – 传递给 DbtCloudHook 构造函数的额外参数。

template_fields = ('dbt_cloud_conn_id', 'run_id', 'path', 'account_id', 'output_file_name')[源代码]
dbt_cloud_conn_id = 'dbt_cloud_default'[源代码]
run_id[源代码]
path[源代码]
account_id = None[源代码]
step = None[源代码]
output_file_name[源代码]
hook_params[源代码]
execute(context)[源代码]

在创建算子时派生。

执行任务的主要方法。Context 是与渲染 jinja 模板时使用的相同字典。

有关更多上下文,请参考 get_template_context。

class airflow.providers.dbt.cloud.operators.dbt.DbtCloudListJobsOperator(*, dbt_cloud_conn_id=DbtCloudHook.default_conn_name, account_id=None, project_id=None, order_by=None, hook_params=None, **kwargs)[源代码]

Bases: airflow.providers.common.compat.sdk.BaseOperator

列出 dbt Cloud 项目中的任务。

另请参阅

有关如何使用此操作符的更多信息,请参阅指南:列出任务

检索关联到指定 dbt Cloud 账户的所有任务的元数据。如果提供了 project_id,则仅检索与该项目 ID 相关的任务。

参数:
  • account_id (int | None) – 可选。如果没有明确提供账户 ID,将使用 dbt Cloud 连接中的账户 ID。

  • order_by (str | None) – 可选。对结果进行排序的字段。使用 ‘-’ 表示反向排序。例如,要按运行 ID 进行反向排序,请使用 order_by=-id

  • project_id (int | None) – 可选。dbt Cloud 项目的 ID。

  • hook_params (dict[str, Any] | None) – 传递给 DbtCloudHook 构造函数的额外参数。

template_fields = ('account_id', 'project_id')[源代码]
dbt_cloud_conn_id = 'dbt_cloud_default'[源代码]
account_id = None[源代码]
project_id = None[源代码]
order_by = None[源代码]
hook_params[源代码]
execute(context)[源代码]

在创建算子时派生。

执行任务的主要方法。Context 是与渲染 jinja 模板时使用的相同字典。

有关更多上下文,请参考 get_template_context。

此条目是否有帮助?