airflow.sdk API 参考

本页面记录了通过 Task SDK Python 模块在 Airflow 3.0+ 中公开的完整公共 API。

如果某项内容未出现在此页面上,最好认为它不属于公共 API,且使用它的风险完全由您自负 —— 我们不会刻意破坏其用法,但也未做任何承诺。

定义 Dag

class airflow.sdk.DAG

DAG 是具有方向性依赖的任务集合。

DAG 还有调度计划、开始日期和(可选的)结束日期。对于每个调度(例如每天或每小时),当任务的依赖关系满足时,DAG 需要运行每个单独的任务。某些任务具有依赖于自身过去执行情况的属性,这意味着在之前的调度周期(及上游任务)完成之前,它们无法运行。

DAG 本质上充当任务的命名空间。一个 task_id 在一个 DAG 中只能添加一次。

请注意,如果您计划使用时区,则提供的所有日期都应为 pendulum 日期。请参阅 时区感知型 DAG

版本 2.4 新增: schedule 参数用于指定基于时间的调度逻辑(timetable)或基于数据集的触发器。

版本 3.0 修改: schedule 的默认值已更改为 None(无调度)。之前的默认值为 timedelta(days=1)

参数:
  • dag_id – DAG 的标识符;必须仅由字母数字字符、破折号、点和下划线组成(全部为 ASCII 字符)

  • description – DAG 的描述,例如会在 Web 服务器界面上显示

  • schedule – 如果提供,此参数定义 DAG 运行的调度规则。可能的值包括 cron 表达式字符串、timedelta 对象、Timetable(调度表)或 Asset(资产)对象列表。另请参阅 使用 Timetable 自定义 DAG 调度

  • start_date – 调度程序将尝试从此时间戳开始回填(backfill)。如果未提供,则必须通过明确的时间范围手动执行回填。

  • end_date – DAG 不会运行超过此日期,留空为 None 表示开放式调度。

  • template_searchpath – 此文件夹列表(非相对路径)定义了 Jinja 查找模板的位置。顺序很重要。请注意,Jinja/Airflow 默认包含您 DAG 文件的路径

  • template_undefined – 模板未定义类型。

  • user_defined_macros – 将在您的 Jinja 模板中公开的宏字典。例如,将 dict(foo='bar') 传递给此参数,允许您在所有与此 DAG 相关的 Jinja 模板中使用 {{ foo }}。注意,您可以在此处传递任何类型的对象。

  • user_defined_filters – 将在您的 Jinja 模板中公开的过滤器字典。例如,将 dict(hello=lambda name: 'Hello %s' % name) 传递给此参数,允许您在所有与此 DAG 相关的 Jinja 模板中使用 {{ 'world' | hello }}

  • default_args – 一个默认参数字典,在初始化算子时用作构造函数的关键字参数。请注意,算子有相同的钩子,且优先级高于此处定义的参数。这意味着如果您的字典在此处包含 ‘depends_on_past’: True,而在算子的 default_args 调用中包含 ‘depends_on_past’: False,则实际值将为 False

  • params – 一个 DAG 级别的参数字典,可在模板中访问,命名空间为 params。这些参数可以在任务级别被覆盖。

  • max_active_tasks – 每次 DAG 运行允许并发执行的任务实例数。请注意,在 Airflow 2 中,这是对 DAG 的全局限制;从 Airflow 3 开始,这是针对单次运行的限制。

  • max_active_runs – 活动 DAG 运行的最大数量;超过此数量的正在运行的 DAG 运行后,调度程序将不再创建新的活动 DAG 运行

  • max_consecutive_failed_dag_runs – (实验性)连续失败的 DAG 运行的最大数量,超过此数量后调度程序将禁用该 DAG

  • dagrun_timeout – 指定 DagRun 在超时或失败前允许运行的持续时间。当 DagRun 超时时,正在运行的任务实例将被标记为已跳过。

  • sla_miss_callback – 已弃用 - SLA 功能在 Airflow 3.0 中被移除,将在 3.1 中由 DeadlineAlerts 替代

  • deadline – DAG 的可选 DeadlineAlert。

  • catchup – 是否执行调度程序追赶(或仅运行最新一次)?默认为 False

  • on_failure_callback – 当此 DAG 的 DagRun 失败时调用的函数或函数列表。上下文字典作为单个参数传递给此函数。

  • on_success_callback – 类似于 on_failure_callback,不同之处在于它在 DAG 成功时执行。

  • access_control – 指定可选的 DAG 级别操作,例如 “{‘role1’: {‘can_read’}, ‘role2’: {‘can_read’, ‘can_edit’, ‘can_delete’}}”,或者如果存在 DAGs Run 资源,则可以指定资源名称,例如 “{‘role1’: {‘DAG Runs’: {‘can_create’}}, ‘role2’: {‘DAGs’: {‘can_read’, ‘can_edit’, ‘can_delete’}}”

  • is_paused_upon_creation – 指定 DAG 在首次创建时是否处于暂停状态。如果 DAG 已经存在,此标志将被忽略。如果未指定此可选参数,将使用全局配置设置。

  • jinja_environment_kwargs

    传递给 Jinja Environment 以进行模板渲染的其他配置选项

    示例: 防止 Jinja 从模板字符串中移除尾随换行符

    DAG(
        dag_id="my-dag",
        jinja_environment_kwargs={
            "keep_trailing_newline": True,
            # some other jinja2 Environment options here
        },
    )
    

    参见: Jinja Environment 文档

  • render_template_as_native_obj – 如果为 True,使用 Jinja NativeEnvironment 将模板渲染为原生 Python 类型。如果为 False,使用 Jinja Environment 将模板渲染为字符串值。

  • tags – 标签列表,有助于在 UI 中过滤 DAG。

  • owner_links – 所有者及其链接的字典,在 DAG 视图 UI 中可点击。可用作 HTTP 链接(例如指向您的 Slack 频道)或 mailto 链接。例如: {"dag_owner": "https://airflow.org.cn/"}

  • auto_register – 当在 with 块中使用时,自动注册此 DAG

  • fail_fast – 当 DAG 中的任务失败时,使当前正在运行的任务失败。警告:快速失败(fail stop)的 DAG 只能包含具有默认触发规则(“all_success”)的任务。如果快速失败 DAG 中的任何任务具有非默认触发规则,将抛出异常。

  • allowed_run_types – 一个可选的列表或单个 DagRunType,指定此 DAG 允许的运行类型。设置后,调度程序和 API 将仅允许指定类型的运行。

  • dag_display_name – 出现在 UI 上的 DAG 显示名称。

配置

conf 对象作为 Task SDK 的一部分提供。它提供了访问配置的接口,允许您读取和交互 Airflow 配置值。

宏 (Macros)

macros 模块作为 Task SDK 的一部分提供。它为 Jinja 模板和任务代码中的日期操作及其他常见操作提供了内置实用函数。

可用函数包括

  • ds_add(ds, days) - 对日期字符串进行加减天数运算

  • ds_format(ds, input_format, output_format) - 格式化日期时间字符串

  • ds_format_locale(ds, input_format, output_format, locale) - 支持本地化的日期时间字符串格式化

  • datetime_diff_for_humans(dt, since) - 人类可读的日期时间差

该模块还提供了直接访问常用标准库模块的途径:json, time, uuid, dateutilrandom

装饰器

airflow.sdk.dag(dag_id='', *, description=None, schedule=None, start_date=None, end_date=None, template_searchpath=None, template_undefined=jinja2.StrictUndefined, user_defined_macros=None, user_defined_filters=None, default_args=None, max_active_tasks=..., max_active_runs=..., max_consecutive_failed_dag_runs=..., dagrun_timeout=None, catchup=..., on_success_callback=None, on_failure_callback=None, deadline=None, doc_md=None, params=None, access_control=None, is_paused_upon_creation=None, jinja_environment_kwargs=None, render_template_as_native_obj=False, tags=None, owner_links=None, auto_register=True, fail_fast=False, allowed_run_types=None, dag_display_name=None, disable_bundle_versioning=False)
参数:
  • dag_id (str)

  • description (str | None)

  • schedule (ScheduleArg)

  • start_date (datetime.datetime | None)

  • end_date (datetime.datetime | None)

  • template_searchpath (str | collections.abc.Iterable[str] | None)

  • template_undefined (type[jinja2.StrictUndefined])

  • user_defined_macros (dict | None)

  • user_defined_filters (dict | None)

  • default_args (dict[str, Any] | None)

  • max_active_tasks (int)

  • max_active_runs (int)

  • max_consecutive_failed_dag_runs (int)

  • dagrun_timeout (datetime.timedelta | None)

  • catchup (bool)

  • on_success_callback (None | DagStateChangeCallback | list[DagStateChangeCallback])

  • on_failure_callback (None | DagStateChangeCallback | list[DagStateChangeCallback])

  • deadline (list[airflow.sdk.definitions.deadline.DeadlineAlert] | airflow.sdk.definitions.deadline.DeadlineAlert | None)

  • doc_md (str | None)

  • params (airflow.sdk.definitions.param.ParamsDict | dict[str, Any] | None)

  • access_control (dict[str, dict[str, collections.abc.Collection[str]]] | dict[str, collections.abc.Collection[str]] | None)

  • is_paused_upon_creation (bool | None)

  • jinja_environment_kwargs (dict | None)

  • render_template_as_native_obj (bool)

  • tags (collections.abc.Collection[str] | None)

  • owner_links (dict[str, str] | None)

  • auto_register (bool)

  • fail_fast (bool)

  • allowed_run_types (airflow.sdk.api.datamodels._generated.DagRunType | collections.abc.Collection[airflow.sdk.api.datamodels._generated.DagRunType] | None)

  • dag_display_name (str | None)

  • disable_bundle_versioning (bool)

返回类型:

collections.abc.Callable[[collections.abc.Callable], collections.abc.Callable[Ellipsis, DAG]]

任务装饰器 (Task Decorators)

  • @task.run_if(condition, skip_message=None) 仅在满足给定条件时运行任务;否则跳过该任务。条件是一个可调用对象,它接收任务执行上下文并返回布尔值或元组 (bool, message)

  • @task.skip_if(condition, skip_message=None) 如果满足给定条件,则跳过任务,并抛出一个带有可选消息的跳过异常。

  • 提供程序特定的任务装饰器位于 @task.<provider> 下,例如 @task.python, @task.docker 等,从已注册的提供程序中动态加载。

airflow.sdk.task_group(group_id: str | None = None, prefix_group_id: bool = True, parent_group: airflow.sdk.definitions.taskgroup.TaskGroup | None = None, dag: airflow.sdk.definitions.dag.DAG | None = None, default_args: dict[str, Any] | None = None, tooltip: str = '', ui_color: str = 'CornflowerBlue', ui_fgcolor: str = '#000', add_suffix_on_collision: bool = False, group_display_name: str = '') collections.abc.Callable[[collections.abc.Callable[FParams, FReturn]], _TaskGroupFactory[FParams, FReturn]]

Python TaskGroup 装饰器。

将函数包装为 Airflow TaskGroup。当用作 @task_group() 形式时,所有参数都转发给底层的 TaskGroup 类。可用于参数化 TaskGroup。

参数:
  • python_callable – 待装饰的函数。

  • tg_kwargs – TaskGroup 对象的关键字参数。

airflow.sdk.task(*args, **kwargs)

提供 @task 语法的实现。

airflow.sdk.setup(func)

装饰函数以将其标记为 setup 任务。

setup 任务在其 DAG 或 TaskGroup 上下文中所有其他任务之前运行,并且可以执行初始化或资源准备工作。

示例

@setup
def initialize_context(...):
    ...
参数:

func (Callable)

返回类型:

Callable

airflow.sdk.teardown(_func=None, *, on_failure_fail_dagrun=False)

装饰函数以将其标记为 teardown 任务。

teardown 任务在其 DAG 或 TaskGroup 上下文中所有主任务之后运行。如果 on_failure_fail_dagrun=True,则 teardown 中的失败将把 DAG 运行标记为失败。

示例

@teardown(on_failure_fail_dagrun=True)
def cleanup(...):
    ...
参数:

on_failure_fail_dagrun (bool)

返回类型:

Callable

airflow.sdk.asset(*, schedule, is_paused_upon_creation=None, dag_id=None, dag_display_name=None, description=None, params=None, on_success_callback=None, on_failure_callback=None, access_control=None, owner_links=NOTHING, tags=NOTHING, name=None, uri=None, group='asset', extra=NOTHING, watchers=NOTHING)

通过装饰实例化函数来创建资产。

参数:
  • schedule (ScheduleArg)

  • is_paused_upon_creation (bool | None)

  • dag_id (str | None)

  • dag_display_name (str | None)

  • description (str | None)

  • params (ParamsDict | None)

  • on_success_callback (None | DagStateChangeCallback | list[DagStateChangeCallback])

  • on_failure_callback (None | DagStateChangeCallback | list[DagStateChangeCallback])

  • access_control (dict[str, dict[str, Collection[str]]] | None)

  • owner_links (dict[str, str])

  • tags (Collection[str])

  • name (str | None)

  • uri (str | ObjectStoragePath | None)

  • group (str)

  • extra (dict[str, JsonValue])

  • watchers (list[BaseTrigger])

返回类型:

基类

class airflow.sdk.BaseAsyncOperator(*, task_id, owner=DEFAULT_OWNER, email=None, email_on_retry=DEFAULT_EMAIL_ON_RETRY, email_on_failure=DEFAULT_EMAIL_ON_FAILURE, retries=DEFAULT_RETRIES, retry_delay=DEFAULT_RETRY_DELAY, retry_exponential_backoff=0, max_retry_delay=None, start_date=None, end_date=None, depends_on_past=False, ignore_first_depends_on_past=DEFAULT_IGNORE_FIRST_DEPENDS_ON_PAST, wait_for_past_depends_before_skipping=DEFAULT_WAIT_FOR_PAST_DEPENDS_BEFORE_SKIPPING, wait_for_downstream=False, dag=None, params=None, default_args=None, priority_weight=DEFAULT_PRIORITY_WEIGHT, weight_rule=DEFAULT_WEIGHT_RULE, queue=DEFAULT_QUEUE, pool=None, pool_slots=DEFAULT_POOL_SLOTS, sla=None, execution_timeout=DEFAULT_TASK_EXECUTION_TIMEOUT, on_execute_callback=None, on_failure_callback=None, on_success_callback=None, on_retry_callback=None, on_skipped_callback=None, pre_execute=None, post_execute=None, trigger_rule=DEFAULT_TRIGGER_RULE, resources=None, run_as_user=None, map_index_template=None, max_active_tis_per_dag=None, max_active_tis_per_dagrun=None, executor=None, executor_config=None, do_xcom_push=True, multiple_outputs=False, inlets=None, outlets=None, task_group=None, doc=None, doc_md=None, doc_json=None, doc_yaml=None, doc_rst=None, task_display_name=None, logger_name=None, allow_nested_operators=True, render_template_as_native_obj=None, **kwargs)

用于异步算子的基类。

与在触发器 (triggerer) 上执行的延迟算子不同,异步算子是在工作节点 (worker) 上执行的。

参数:
  • task_id (str)

  • owner (str)

  • email (str | collections.abc.Sequence[str] | None)

  • email_on_retry (bool)

  • email_on_failure (bool)

  • retries (int | None)

  • retry_delay (datetime.timedelta | float)

  • retry_exponential_backoff (float)

  • max_retry_delay (datetime.timedelta | float | None)

  • start_date (datetime.datetime | None)

  • end_date (datetime.datetime | None)

  • depends_on_past (bool)

  • ignore_first_depends_on_past (bool)

  • wait_for_past_depends_before_skipping (bool)

  • wait_for_downstream (bool)

  • dag (airflow.sdk.definitions.dag.DAG | None)

  • params (collections.abc.MutableMapping[str, Any] | None)

  • default_args (dict | None)

  • priority_weight (int)

  • weight_rule (airflow.sdk.types.WeightRuleParam)

  • queue (str)

  • pool (str | None)

  • pool_slots (int)

  • sla (datetime.timedelta | None)

  • execution_timeout (datetime.timedelta | None)

  • on_execute_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_failure_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | list[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_success_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_retry_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | list[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_skipped_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • pre_execute (TaskPreExecuteHook | None)

  • post_execute (TaskPostExecuteHook | None)

  • trigger_rule (str)

  • resources (dict[str, Any] | None)

  • run_as_user (str | None)

  • map_index_template (str | None)

  • max_active_tis_per_dag (int | None)

  • max_active_tis_per_dagrun (int | None)

  • executor (str | None)

  • executor_config (dict | None)

  • do_xcom_push (bool)

  • multiple_outputs (bool)

  • inlets (Any | None)

  • outlets (Any | None)

  • task_group (airflow.sdk.definitions.taskgroup.TaskGroup | None)

  • doc (str | None)

  • doc_md (str | None)

  • doc_json (str | None)

  • doc_yaml (str | None)

  • doc_rst (str | None)

  • task_display_name (str | None)

  • logger_name (str | None)

  • allow_nested_operators (bool)

  • render_template_as_native_obj (bool | None)

  • kwargs (Any)

class airflow.sdk.BaseBranchOperator(*, task_id, owner=DEFAULT_OWNER, email=None, email_on_retry=DEFAULT_EMAIL_ON_RETRY, email_on_failure=DEFAULT_EMAIL_ON_FAILURE, retries=DEFAULT_RETRIES, retry_delay=DEFAULT_RETRY_DELAY, retry_exponential_backoff=0, max_retry_delay=None, start_date=None, end_date=None, depends_on_past=False, ignore_first_depends_on_past=DEFAULT_IGNORE_FIRST_DEPENDS_ON_PAST, wait_for_past_depends_before_skipping=DEFAULT_WAIT_FOR_PAST_DEPENDS_BEFORE_SKIPPING, wait_for_downstream=False, dag=None, params=None, default_args=None, priority_weight=DEFAULT_PRIORITY_WEIGHT, weight_rule=DEFAULT_WEIGHT_RULE, queue=DEFAULT_QUEUE, pool=None, pool_slots=DEFAULT_POOL_SLOTS, sla=None, execution_timeout=DEFAULT_TASK_EXECUTION_TIMEOUT, on_execute_callback=None, on_failure_callback=None, on_success_callback=None, on_retry_callback=None, on_skipped_callback=None, pre_execute=None, post_execute=None, trigger_rule=DEFAULT_TRIGGER_RULE, resources=None, run_as_user=None, map_index_template=None, max_active_tis_per_dag=None, max_active_tis_per_dagrun=None, executor=None, executor_config=None, do_xcom_push=True, multiple_outputs=False, inlets=None, outlets=None, task_group=None, doc=None, doc_md=None, doc_json=None, doc_yaml=None, doc_rst=None, task_display_name=None, logger_name=None, allow_nested_operators=True, render_template_as_native_obj=None, **kwargs)

用于创建具有分支功能的算子的基类,类似于 BranchPythonOperator。

用户应从此算子创建子类并实现函数 choose_branch(self, context)。此函数应运行确定分支所需的任何业务逻辑,并返回以下内容之一:- 单个 task_id(作为字符串)- 单个 task_group_id(作为字符串)- 包含 task_ids 和 task_group_ids 组合的列表

该算子将继续执行返回的 task_id(s) 和/或 task_group_id(s),而所有其他直接位于该算子下游的任务将被跳过。

参数:
  • task_id (str)

  • owner (str)

  • email (str | collections.abc.Sequence[str] | None)

  • email_on_retry (bool)

  • email_on_failure (bool)

  • retries (int | None)

  • retry_delay (datetime.timedelta | float)

  • retry_exponential_backoff (float)

  • max_retry_delay (datetime.timedelta | float | None)

  • start_date (datetime.datetime | None)

  • end_date (datetime.datetime | None)

  • depends_on_past (bool)

  • ignore_first_depends_on_past (bool)

  • wait_for_past_depends_before_skipping (bool)

  • wait_for_downstream (bool)

  • dag (airflow.sdk.definitions.dag.DAG | None)

  • params (collections.abc.MutableMapping[str, Any] | None)

  • default_args (dict | None)

  • priority_weight (int)

  • weight_rule (airflow.sdk.types.WeightRuleParam)

  • queue (str)

  • pool (str | None)

  • pool_slots (int)

  • sla (datetime.timedelta | None)

  • execution_timeout (datetime.timedelta | None)

  • on_execute_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_failure_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | list[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_success_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_retry_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | list[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_skipped_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • pre_execute (TaskPreExecuteHook | None)

  • post_execute (TaskPostExecuteHook | None)

  • trigger_rule (str)

  • resources (dict[str, Any] | None)

  • run_as_user (str | None)

  • map_index_template (str | None)

  • max_active_tis_per_dag (int | None)

  • max_active_tis_per_dagrun (int | None)

  • executor (str | None)

  • executor_config (dict | None)

  • do_xcom_push (bool)

  • multiple_outputs (bool)

  • inlets (Any | None)

  • outlets (Any | None)

  • task_group (airflow.sdk.definitions.taskgroup.TaskGroup | None)

  • doc (str | None)

  • doc_md (str | None)

  • doc_json (str | None)

  • doc_yaml (str | None)

  • doc_rst (str | None)

  • task_display_name (str | None)

  • logger_name (str | None)

  • allow_nested_operators (bool)

  • render_template_as_native_obj (bool | None)

  • kwargs (Any)

class airflow.sdk.BaseOperator(*, task_id, owner=DEFAULT_OWNER, email=None, email_on_retry=DEFAULT_EMAIL_ON_RETRY, email_on_failure=DEFAULT_EMAIL_ON_FAILURE, retries=DEFAULT_RETRIES, retry_delay=DEFAULT_RETRY_DELAY, retry_exponential_backoff=0, max_retry_delay=None, start_date=None, end_date=None, depends_on_past=False, ignore_first_depends_on_past=DEFAULT_IGNORE_FIRST_DEPENDS_ON_PAST, wait_for_past_depends_before_skipping=DEFAULT_WAIT_FOR_PAST_DEPENDS_BEFORE_SKIPPING, wait_for_downstream=False, dag=None, params=None, default_args=None, priority_weight=DEFAULT_PRIORITY_WEIGHT, weight_rule=DEFAULT_WEIGHT_RULE, queue=DEFAULT_QUEUE, pool=None, pool_slots=DEFAULT_POOL_SLOTS, sla=None, execution_timeout=DEFAULT_TASK_EXECUTION_TIMEOUT, on_execute_callback=None, on_failure_callback=None, on_success_callback=None, on_retry_callback=None, on_skipped_callback=None, pre_execute=None, post_execute=None, trigger_rule=DEFAULT_TRIGGER_RULE, resources=None, run_as_user=None, map_index_template=None, max_active_tis_per_dag=None, max_active_tis_per_dagrun=None, executor=None, executor_config=None, do_xcom_push=True, multiple_outputs=False, inlets=None, outlets=None, task_group=None, doc=None, doc_md=None, doc_json=None, doc_yaml=None, doc_rst=None, task_display_name=None, logger_name=None, allow_nested_operators=True, render_template_as_native_obj=None, **kwargs)

所有算子的抽象基类。

由于算子创建的对象会成为 DAG 中的节点,因此 BaseOperator 包含了许多用于 DAG 遍历行为的递归方法。要从此类派生,您需要覆盖构造函数和 ‘execute’ 方法。

从此类派生的算子应同步执行或触发某些任务(等待完成)。算子的示例可以是运行 Pig 作业的算子(PigOperator),等待分区在 Hive 中落地的传感器算子(HiveSensorOperator),或者将数据从 Hive 移动到 MySQL 的算子(Hive2MySqlOperator)。这些算子(任务)的实例针对特定操作,运行特定的脚本、函数或数据传输。

此类是抽象的,不应实例化。实例化从该类派生的类会导致创建一个任务对象,该对象最终成为 DAG 对象中的节点。应使用 set_upstream 和/或 set_downstream 方法设置任务依赖关系。

参数:
  • task_id (str) – 任务的唯一且有意义的 ID

  • owner (str) – 任务的所有者。建议使用有意义的描述(例如用户/个人/团队/角色名称)以明确所有权。

  • email (str | collections.abc.Sequence[str] | None) – 电子邮件警报中使用的“收件人”电子邮件地址。这可以是一个或多个电子邮件。多个地址可以指定为以逗号或分号分隔的字符串,或通过传递字符串列表来指定。(已弃用)

  • email_on_retry (bool) – 指示任务重试时是否应发送电子邮件警报(已弃用)

  • email_on_failure (bool) – 指示任务失败时是否应发送电子邮件警报(已弃用)

  • retries (int | None) – 任务失败前应执行的重试次数

  • retry_delay (datetime.timedelta | float) – 重试之间的延迟,可以设置为 timedeltafloat 秒,后者将被转换为 timedelta,默认值为 timedelta(seconds=300)

  • retry_exponential_backoff (float) – 重试之间的指数退避乘数。设置为 0 以禁用(恒定延迟)。设置为 2.0 以进行标准指数退避(延迟随每次重试翻倍)。例如,如果 retry_delay=4min 且 retry_exponential_backoff=5,重试将在 4 分钟、20 分钟、100 分钟后发生,依此类推。

  • max_retry_delay (datetime.timedelta | float | None) – 重试之间的最大延迟间隔,可以设置为 timedeltafloat 秒,后者将被转换为 timedelta

  • start_date (datetime.datetime | None) – 任务的 start_date,确定第一个任务实例的 logical_date。最佳做法是将 start_date 四舍五入到 DAG 的 schedule_interval。每日任务的 start_date 应设为某天的 00:00:00,小时任务的 start_date 应设为特定小时的 00:00。注意,Airflow 只是查找最新的 logical_date 并加上 schedule_interval 来确定下一个 logical_date。同样非常重要的是,不同任务的依赖关系需要在时间上对齐。如果任务 A 依赖于任务 B,且它们的 start_date 偏移量导致它们的 logical_date 无法对齐,则 A 的依赖关系将永远无法满足。如果您想延迟任务,例如在凌晨 2 点运行每日任务,请查看 TimeSensorTimeDeltaSensor。我们建议不要使用动态 start_date,而建议使用固定的日期。请阅读 FAQ 中关于 start_date 的条目以获取更多信息。

  • end_date (datetime.datetime | None) – 如果指定,调度程序将不会超过此日期

  • depends_on_past (bool) – 当设置为 true 时,任务实例将按顺序运行,并且仅在上一个实例已成功或已被跳过的情况下运行。start_date 的任务实例被允许运行。

  • wait_for_past_depends_before_skipping (bool) – 当设置为 true 时,如果任务实例应被标记为已跳过,并且 depends_on_past 为 true,则 TI 将保持在 None 状态,等待前一次运行的任务

  • wait_for_downstream (bool) – 当设置为 true 时,任务 X 的一个实例将等待任务 X 前一个实例的紧接下游任务成功完成或被跳过,然后再运行。如果任务 X 的不同实例修改了同一个资产,并且此资产被任务 X 下游的任务使用,这非常有用。注意,凡是使用 wait_for_downstream 的地方,depends_on_past 都被强制为 True。还要注意,只等待上一个任务实例紧接下游的任务;忽略任何更下游任务的状态。

  • dag (airflow.sdk.definitions.dag.DAG | None) – 对任务所属 DAG 的引用(如果有)

  • priority_weight (int) – 该任务相对于其他任务的优先级权重。这允许执行程序在积压时优先触发更高优先级的任务。将 priority_weight 设置为更大的数字表示更重要的任务。由于并非所有数据库引擎都支持 64 位整数,因此值被限制为 32 位。有效范围是从 -2,147,483,648 到 2,147,483,647。

  • weight_rule (airflow.sdk.types.WeightRuleParam) – 用于任务有效总优先级权重的加权方法。选项有:{ downstream | upstream | absolute },默认值为 downstream。当设置为 downstream 时,任务的有效权重是所有下游后代的聚合总和。因此,上游任务将具有更高的权重,在使用正权重值时会更激进地被调度。当您有多个 DAG 运行实例并希望在每个 DAG 继续处理下游任务之前完成所有运行的所有上游任务时,这非常有用。当设置为 upstream 时,有效权重是所有上游祖先的聚合总和。这与上面相反,下游任务权重更高,在使用正权重值时会更激进地被调度。当您有多个 DAG 运行实例并希望在启动其他 DAG 的上游任务之前完成每个 DAG 时,这很有用。当设置为 absolute 时,有效权重就是指定的 priority_weight,没有额外的加权。当您确切知道每个任务应具有的优先级权重时,可能需要这样做。此外,设置为 absolute 时,对于非常大的 DAG,还有一个显著加快任务创建过程的额外好处。选项可以设置为字符串或使用静态类 airflow.utils.WeightRule 中定义的常量。无论权重规则如何,生成的优先级值都上限为 32 位。这是一个 实验性功能。自 2.9.0 版本起,Airflow 允许通过创建 airflow.task.priority_strategy.PriorityWeightStrategy 的子类并将其注册到插件中,然后通过 weight_rule 参数提供类路径或类实例,来定义自定义优先级权重策略。自定义优先级权重策略将用于计算任务实例的有效总优先级权重。

  • queue (str) – 运行此作业时要瞄准的队列。并非所有执行程序都实现队列管理,CeleryExecutor 确实支持瞄准特定队列。

  • pool (str | None) – 此任务应在其中运行的槽位池(slot pool),槽位池是限制某些任务并发性的一种方式

  • pool_slots (int) – 此任务应使用的池槽位数量(>= 1)。不允许小于 1 的值。

  • sla (datetime.timedelta | None) – 已弃用 - SLA 功能在 Airflow 3.0 中被移除,将在 Airflow >=3.1 中由新的实现替代。

  • execution_timeout (datetime.timedelta | None) – 此任务实例执行允许的最大时间,如果超过该时间,它将引发错误并失败。

  • on_failure_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | list[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback]) – 当此任务的任务实例失败时调用的函数或函数列表。上下文字典作为单个参数传递给此函数。上下文包含对与任务实例相关的对象的引用,并记录在 API 的宏部分下。

  • on_execute_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback]) – 类似于 on_failure_callback,不同之处在于它在任务执行前立即执行。

  • on_retry_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | list[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback]) – 类似于 on_failure_callback,不同之处在于它在发生重试时执行。

  • on_success_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback]) – 类似于 on_failure_callback,不同之处在于它在任务成功时执行。

  • on_skipped_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback]) – 类似于 on_failure_callback,不同之处在于它在发生跳过时执行;此回调仅在引发 AirflowSkipException 时调用。明确地说,如果由于 DAG 中的前一个分支决策或导致执行跳过的触发规则而未启动任务执行,从而导致任务调度从未安排,则此回调不会被调用。

  • pre_execute (TaskPreExecuteHook | None) – 在任务执行前立即调用的函数,接收一个上下文字典;抛出异常将阻止任务执行。

  • post_execute (TaskPostExecuteHook | None) – 在任务执行后立即调用的函数,接收一个上下文字典和任务结果;抛出异常将阻止任务成功。

  • trigger_rule (str) – 定义应用任务触发依赖关系的规则。选项为:{ all_success | all_failed | all_done | all_skipped | one_success | one_done | one_failed | none_failed | none_failed_min_one_success | none_skipped | always},默认值为 all_success。选项可以设置为字符串或使用静态类 airflow.utils.TriggerRule 中定义的常量

  • resources (dict[str, Any] | None) – 资源参数名称(Resources 构造函数的参数名称)到其值的映射。

  • run_as_user (str | None) – 运行任务时要模拟的 unix 用户名

  • max_active_tis_per_dag (int | None) – 设置后,任务将能够限制跨逻辑日期 (logical_dates) 的并发运行。

  • max_active_tis_per_dagrun (int | None) – 设置后,任务将能够限制每次 DAG 运行的并发任务实例数。

  • executor (str | None) – 运行此任务时要指定的目标执行器。尚不支持

  • executor_config (dict | None) –

    由特定执行器解释的其他任务级配置参数。参数由执行器名称进行命名空间划分。

    示例:通过 KubernetesExecutor 在特定的 docker 容器中运行此任务

    MyOperator(..., executor_config={"KubernetesExecutor": {"image": "myCustomDockerImage"}})
    

  • do_xcom_push (bool) – 如果为 True,则推送包含 Operator 结果的 XCom

  • multiple_outputs (bool) – 如果为 True 且 do_xcom_push 为 True,则推送多个 XCom,返回字典结果中的每个键对应一个。如果为 False 且 do_xcom_push 为 True,则推送单个 XCom。

  • task_group (airflow.sdk.definitions.taskgroup.TaskGroup | None) – 任务所属的 TaskGroup。通常在不使用 TaskGroup 作为上下文管理器时提供。

  • doc (str | None) – 将文档或注释添加到您的任务对象中,这些内容可在 Web 服务器的任务实例详细信息视图中查看

  • doc_md (str | None) – 将文档(Markdown 格式)或注释添加到您的任务对象中,这些内容可在 Web 服务器的任务实例详细信息视图中查看

  • doc_rst (str | None) – 将文档(RST 格式)或注释添加到您的任务对象中,这些内容可在 Web 服务器的任务实例详细信息视图中查看

  • doc_json (str | None) – 将文档(JSON 格式)或注释添加到您的任务对象中,这些内容可在 Web 服务器的任务实例详细信息视图中查看

  • doc_yaml (str | None) – 将文档(YAML 格式)或注释添加到您的任务对象中,这些内容可在 Web 服务器的任务实例详细信息视图中查看

  • task_display_name (str | None) – 在 UI 上显示的任务名称。

  • logger_name (str | None) – Operator 用于输出日志的记录器名称。如果设置为 None(默认值),记录器名称将回退到 airflow.task.operators.{class.__module__}.{class.__name__}(例如,HttpOperator 的记录器将是 airflow.task.operators.airflow.providers.http.operators.http.HttpOperator)。

  • allow_nested_operators (bool) –

    如果为 True,当一个算子在另一个算子中执行时,将记录一条警告消息。如果为 False,则在算子使用不当(例如嵌套在另一个算子中)时将引发异常。在未来的 Airflow 版本中,此参数将被删除,并且当算子相互嵌套时,将始终抛出异常(默认值为 True)。

    示例:错误使用算子 mixin 的示例

    @task(provide_context=True)
    def say_hello_world(**context):
        hello_world_task = BashOperator(
            task_id="hello_world_task",
            bash_command="python -c \"print('Hello, world!')\"",
            dag=dag,
        )
        hello_world_task.execute(context)
    

  • render_template_as_native_obj (bool | None) – 如果为 True,则使用 Jinja NativeEnvironment 将模板呈现为原生 Python 类型。如果为 False,则使用 Jinja Environment 将模板呈现为字符串值。如果为 None(默认值),则从 DAG 设置继承。

  • ignore_first_depends_on_past (bool)

  • params (collections.abc.MutableMapping[str, Any] | None)

  • default_args (dict | None)

  • map_index_template (str | None)

  • inlets (Any | None)

  • outlets (Any | None)

  • kwargs (Any)

class airflow.sdk.BaseSensorOperator(*, poke_interval=60, timeout=conf.getfloat('sensors', 'default_timeout'), soft_fail=False, mode='poke', exponential_backoff=False, max_wait=None, silent_fail=False, never_fail=False, **kwargs)

传感器算子派生自此类并继承这些属性。

传感器算子持续以特定的时间间隔执行,当满足标准时成功,并在超时时失败。

参数:
  • soft_fail (bool) – 设置为 true 可在失败时将任务标记为 SKIPPED。与 never_fail 互斥。

  • poke_interval (datetime.timedelta | float) – 作业在每次尝试之间应等待的时间。可以是 timedeltafloat 秒数。

  • timeout (datetime.timedelta | float) – 任务超时并失败前经过的时间。可以是 timedeltafloat 秒数。这不应与 BaseOperator 类的 execution_timeout 混淆。timeout 测量的是第一次 poke 到当前时间之间经过的时间(考虑到每次 poke 之间的任何重调度延迟),而 execution_timeout 检查任务的 运行 时间(不包括任何重调度延迟)。如果 modepoke(见下文),两者是等效的(因为传感器从未被重调度),而在 reschedule 模式下则不是。

  • mode (str) – 传感器的操作方式。选项为:{ poke | reschedule },默认值为 poke。当设置为 poke 时,传感器在整个执行时间内占用工作节点插槽,并在 poke 之间休眠。如果传感器的预期运行时间很短,或者需要较短的 poke 间隔,请使用此模式。请注意,在此模式下,传感器将占用一个工作节点插槽和一个池插槽,持续时间为其运行时间。当设置为 reschedule 时,如果尚未满足标准,传感器任务会在 poke 之间释放工作节点插槽,并在稍后重调度。如果预期满足标准前的时间较长,请使用此模式。poke 间隔应超过一分钟,以防止调度程序负载过重。

  • exponential_backoff (bool) – 通过使用指数退避算法,允许在 poke 之间逐步延长等待时间

  • max_wait (datetime.timedelta | float | None) – poke 之间的最大等待间隔,可以是 timedeltafloat 秒数

  • silent_fail (bool) – 如果为 true,且 poke 方法引发了不同于 AirflowSensorTimeout、AirflowTaskTimeout、AirflowSkipException 和 AirflowFailException 的异常,则传感器将记录错误并继续执行。否则,传感器任务失败,并可根据提供的 retries 参数进行重试。

  • never_fail (bool) – 如果为 true,且 poke 方法引发异常,则跳过传感器。与 soft_fail 互斥。

class airflow.sdk.BaseNotifier(context=None)

用于发送通知的 BaseNotifier 类。

如果实现了 async_notify,可以异步使用(推荐);如果实现了 notify,可以同步使用。

目前,DAG/任务状态更改回调在 DAG 处理器上运行,仅支持同步用法。

用法:

# 异步用法 await Notifier(context)

# 同步用法 notifier = Notifier() notifier(context)

参数:

context (airflow.sdk.definitions.context.Context | None)

定义如何获取算子链接的抽象基类。

class airflow.sdk.BaseXCom

BaseXcom 现在是一个与 XCom 后端交互的接口。

class airflow.sdk.BranchMixIn(context=None)

将分支处理为单行代码的实用辅助类。

class airflow.sdk.PokeReturnValue(is_done, xcom_value=None)

poke 方法的可选返回值。

传感器可以在 poke 方法中选择返回 PokeReturnValue 类的实例。如果传感器完成时提供了 XCom 值,则 XCom 值将通过算子返回值进行推送。:param is_done: 设置为 true 以指示传感器可以停止 poking。:param xcom_value: 一个可选的 XCOM 值,将由算子返回。

参数:
  • is_done (bool)

  • xcom_value (Any | None)

class airflow.sdk.SkipMixin(context=None)

用于跳过任务实例的 Mixin。

class airflow.sdk.BaseHook(logger_name=None)

Hook 的抽象基类。

Hook 被设计为与外部系统交互的接口。MySqlHook、HiveHook、PigHook 返回可处理连接并与这些系统的特定实例进行交互的对象,并公开一致的方法与它们交互。

参数:

logger_name (str | None) – Hook 用于输出日志的记录器名称。如果设置为 None(默认值),记录器名称将回退到 airflow.task.hooks.{class.__module__}.{class.__name__}(例如,DbApiHook 的记录器将是 airflow.task.hooks.airflow.providers.common.sql.hooks.sql.DbApiHook)。

回调 (Callbacks)

class airflow.sdk.AsyncCallback(callback_callable, kwargs=None)

在触发器中运行的异步回调。

callback_callable 可以是 Python 可调用类型,也可以是包含可调用路径的字符串,该字符串可用于导入可调用对象。它必须是存在于触发器模块中的顶级可等待可调用对象。

当错过截止日期时,它将使用 Airflow 上下文和指定的 kwargs 被调用。

参数:
  • callback_callable (Callable | str)

  • kwargs (dict)

class airflow.sdk.SyncCallback(callback_callable, kwargs=None, executor=None)

在指定的或默认执行器中运行的同步回调。

callback_callable 可以是 Python 可调用类型,也可以是包含可调用路径的字符串,该字符串可用于导入可调用对象。它必须是存在于执行器模块中的顶级可调用对象。

当错过截止日期时,它将使用 Airflow 上下文和指定的 kwargs 被调用。

参数:
  • callback_callable (Callable | str)

  • kwargs (dict)

  • executor (str | None)

截止日期提醒 (Deadline Alerts)

class airflow.sdk.DeadlineAlert(reference, interval, callback)

存储计算需求时间戳和回调信息所需的 Deadline 值。

参数:
  • reference (DeadlineReferenceType)

  • interval (timedelta)

  • callback (Callback)

class airflow.sdk.DeadlineReference

所有 DeadlineReference 选项的公共接口类。

此类为处理截止日期提供了统一接口,支持计算截止日期(从数据库获取值)和固定截止日期(返回预定义的日期时间)。

用法:

  1. 截止日期参考示例

fixed = DeadlineReference.FIXED_DATETIME(datetime(2025, 5, 4))
logical = DeadlineReference.DAGRUN_LOGICAL_DATE
queued = DeadlineReference.DAGRUN_QUEUED_AT
  1. 在 DAG 中使用

DAG(
    dag_id="dag_with_deadline",
    deadline=DeadlineAlert(
        reference=DeadlineReference.DAGRUN_LOGICAL_DATE,
        interval=timedelta(hours=1),
        callback=hello_callback,
    ),
)
  1. 评估截止日期时将忽略意外参数

# For deadlines requiring parameters:
deadline = DeadlineReference.DAGRUN_LOGICAL_DATE
deadline.evaluate_with(dag_id=dag.dag_id)

# For deadlines with no required parameters:
deadline = DeadlineReference.FIXED_DATETIME(datetime(2025, 5, 4))
deadline.evaluate_with()

连接和变量 (Connections & Variables)

class airflow.sdk.Connection(*, conn_id: str, uri: str)

到外部数据源的连接。

参数:
  • conn_id – 连接 ID。

  • conn_type – 连接类型。

  • description – 连接描述。

  • host – 主机。

  • login – 登录名。

  • password – 密码。

  • schema – 模式。

  • port – 端口号。

  • extra – 额外元数据。私钥/SSH 密钥等非标准数据可保存于此。JSON 编码对象。

  • uri – 描述连接参数的 URI 地址。

class airflow.sdk.Variable

一种用于以简单的键/值存储方式存储和检索任意内容或设置的通用方法。

参数:
  • key – 变量键。

  • value – 变量值。

  • description – 变量描述。

任务与运算符

class airflow.sdk.TaskGroup

任务的集合。

当在 TaskGroup 上调用 set_downstream() 或 set_upstream() 时,如有必要,该操作将应用于组内的所有任务。

参数:
  • group_id – TaskGroup 的唯一且有意义的 ID。group_id 不得与 Dag 中的其他 TaskGroup 的 group_id 或任务的 task_id 冲突。根 TaskGroup 的 group_id 设置为 None。

  • prefix_group_id – 如果设置为 True,子任务的 task_id 和 group_id 将以该 TaskGroup 的 group_id 为前缀。如果设置为 False,则不添加前缀。默认为 True。

  • parent_group – 此 TaskGroup 的父 TaskGroup。根 TaskGroup 的 parent_group 设置为 None。

  • dag – 此 TaskGroup 所属的 Dag。

  • default_args – 初始化运算符时用作构造函数关键字参数的默认参数字典,将覆盖 Dag 级别定义的 default_args。请注意,运算符具有相同的钩子(hook),并且优先级高于此处定义的参数,这意味着如果您的字典包含 ‘depends_on_past’: True,而运算符调用中的 default_args 包含 ‘depends_on_past’: False,则实际值将为 False

  • tooltip – 在 UI 中显示时 TaskGroup 节点的工具提示

  • ui_color – 在 UI 中显示时 TaskGroup 节点的填充颜色

  • ui_fgcolor – 在 UI 中显示时 TaskGroup 节点的标签颜色

  • add_suffix_on_collision – 如果此任务组名称已存在,则自动添加 __1 等后缀

  • group_display_name – 如果设置,这将是 UI 中 TaskGroup 节点的显示名称。

class airflow.sdk.TaskInstance(*args, **kwargs)

运行时可用的 TaskInstance 协议。

此类提供了在 Task SDK 中与 TaskInstance 属性和方法(如 xcom_pull/push)进行交互的接口。

class airflow.sdk.XComArg

对从另一个运算符推送的 XCom 值的引用。

该实现支持

xcomarg >> op
xcomarg << op
op >> xcomarg  # By BaseOperator code
op << xcomarg  # By BaseOperator code

示例:当您从任何运算符(装饰过或常规的)获得结果时,您可以

any_op = AnyOperator()
xcomarg = XComArg(any_op)
# or equivalently
xcomarg = any_op.output
my_op = MyOperator()
my_op >> xcomarg

此对象可以通过 Jinja 在遗留运算符中使用。

示例:您可以使此结果成为任何生成字符串的一部分

any_op = AnyOperator()
xcomarg = any_op.output
op1 = MyOperator(my_text_message=f"the value is {xcomarg}")
op2 = MyOperator(my_text_message=f"the value is {xcomarg['topic']}")
参数:
  • operator – XComArg 引用的运算符实例。

  • key – 用于拉取 XCom 值的键。默认为 XCOM_RETURN_KEY,即引用运算符的返回值。

airflow.sdk.literal(value)

包装一个值以确保其按原样呈现,而不对其内容应用 Jinja 模板。

专为在运算符的模板字段中使用而设计。

参数:

value (Any) – 要呈现且不应用模板的值

返回类型:

airflow.sdk.definitions._internal.templater.LiteralValue

class airflow.sdk.Param(default=NOTSET, description=None, source=None, **kwargs)

用于保存 Param 默认值和执行验证的规则集的类。

如果没有规则集,它始终验证并返回默认值。

参数:
  • default (Any) – 此 Param 对象持有的值

  • description (str | None) – Param 的可选帮助文本

  • schema – Param 的验证模式,如果未提供,则除 default 和 description 之外的所有 kwargs 将构成模式

  • source (Literal['dag', 'task'] | None)

class airflow.sdk.ParamsDict(dict_obj=None, suppress_exception=False)

用于保存 dag 或任务的所有参数的类。

所有键均为严格字符串,如果值尚不是 Param 对象,则会转换为 Param 对象。此类旨在隐式替换参数字典,理想情况下无需直接使用。

参数:
  • dict_obj (Mapping[str, Any] | None) – 用于初始化 ParamsDict 的字典或类字典对象

  • suppress_exception (bool) – 在初始化 ParamsDict 时抑制值异常的标志

class airflow.sdk.TriggerRule(value)

一种枚举。

airflow.sdk.get_current_context()

检索执行上下文字典,而不更改用户方法的签名。

这是检索执行上下文字典的最简单方法。

旧式

def my_task(**context):
    ti = context["ti"]

新式

from airflow.sdk import get_current_context


def my_task():
    context = get_current_context()
    ti = context["ti"]

仅当在运算符开始执行后调用此方法时,当前上下文才会有值。

返回类型:

Context

airflow.sdk.get_parsing_context()

返回当前的 (Dag) 解析上下文信息。

返回类型:

AirflowParsingContext

状态枚举

class airflow.sdk.TaskInstanceState(value)

任务实例可能处于的所有状态。

请注意,也允许为 None,因此在类型提示中请务必将其与 Optional 一起使用。

class airflow.sdk.DagRunState(value)

DagRun 可能处于的所有状态。

在代码的某些部分中,这些状态与 TaskInstanceState 是“共享”的,因此请确保它们的值始终与 TaskInstanceState 中同名的值匹配。

class airflow.sdk.WeightRule(value)

一种枚举。

设置依赖项

airflow.sdk.chain(*tasks)

给定多个任务,构建依赖链。

此函数接受 BaseOperator(即任务)、EdgeModifiers(即标签)、XComArg、TaskGroups 或包含这些类型中任何混合的列表(或同一列表中的混合)。如果要链接两个列表,则必须确保它们具有相同的长度。

使用传统运算符/传感器

chain(t1, [t2, t3], [t4, t5], t6)

等同于

  / -> t2 -> t4 \
t1               -> t6
  \ -> t3 -> t5 /
t1.set_downstream(t2)
t1.set_downstream(t3)
t2.set_downstream(t4)
t3.set_downstream(t5)
t4.set_downstream(t6)
t5.set_downstream(t6)

使用任务装饰函数即 XComArgs

chain(x1(), [x2(), x3()], [x4(), x5()], x6())

等同于

  / -> x2 -> x4 \
x1               -> x6
  \ -> x3 -> x5 /
x1 = x1()
x2 = x2()
x3 = x3()
x4 = x4()
x5 = x5()
x6 = x6()
x1.set_downstream(x2)
x1.set_downstream(x3)
x2.set_downstream(x4)
x3.set_downstream(x5)
x4.set_downstream(x6)
x5.set_downstream(x6)

使用 TaskGroups

chain(t1, task_group1, task_group2, t2)

t1.set_downstream(task_group1)
task_group1.set_downstream(task_group2)
task_group2.set_downstream(t2)

也可以在传统运算符/传感器、EdgeModifiers、XComArg 和 TaskGroups 之间进行混合

chain(t1, [Label("branch one"), Label("branch two")], [x1(), x2()], task_group1, x3())

等同于

  / "branch one" -> x1 \
t1                      -> task_group1 -> x3
  \ "branch two" -> x2 /
x1 = x1()
x2 = x2()
x3 = x3()
label1 = Label("branch one")
label2 = Label("branch two")
t1.set_downstream(label1)
label1.set_downstream(x1)
t2.set_downstream(label2)
label2.set_downstream(x2)
x1.set_downstream(task_group1)
x2.set_downstream(task_group1)
task_group1.set_downstream(x3)

# or

x1 = x1()
x2 = x2()
x3 = x3()
t1.set_downstream(x1, edge_modifier=Label("branch one"))
t1.set_downstream(x2, edge_modifier=Label("branch two"))
x1.set_downstream(task_group1)
x2.set_downstream(task_group1)
task_group1.set_downstream(x3)
参数:

tasks (airflow.sdk.definitions._internal.abstractoperator.DependencyMixin | collections.abc.Sequence[airflow.sdk.definitions._internal.abstractoperator.DependencyMixin]) – 设置依赖项的单个任务和/或任务列表、EdgeModifiers、XComArgs 或 TaskGroups

返回类型:

airflow.sdk.chain_linear(*elements)

简化任务依赖定义。

例如:假设您想要这样的优先级

    ╭─op2─╮ ╭─op4─╮
op1─┤     ├─├─op5─┤─op7
    ╰-op3─╯ ╰-op6─╯

那么您可以这样完成

chain_linear(op1, [op2, op3], [op4, op5, op6], op7)
参数:

elements (airflow.sdk.definitions._internal.abstractoperator.DependencyMixin | collections.abc.Sequence[airflow.sdk.definitions._internal.abstractoperator.DependencyMixin]) – 运算符列表/运算符列表

airflow.sdk.cross_downstream(from_tasks, to_tasks)

将 from_tasks 中所有任务的下游依赖项设置为 to_tasks 中的所有任务。

使用传统运算符/传感器

cross_downstream(from_tasks=[t1, t2, t3], to_tasks=[t4, t5, t6])

等同于

t1 ---> t4
   \ /
t2 -X -> t5
   / \
t3 ---> t6
t1.set_downstream(t4)
t1.set_downstream(t5)
t1.set_downstream(t6)
t2.set_downstream(t4)
t2.set_downstream(t5)
t2.set_downstream(t6)
t3.set_downstream(t4)
t3.set_downstream(t5)
t3.set_downstream(t6)

使用任务装饰函数即 XComArgs

cross_downstream(from_tasks=[x1(), x2(), x3()], to_tasks=[x4(), x5(), x6()])

等同于

x1 ---> x4
   \ /
x2 -X -> x5
   / \
x3 ---> x6
x1 = x1()
x2 = x2()
x3 = x3()
x4 = x4()
x5 = x5()
x6 = x6()
x1.set_downstream(x4)
x1.set_downstream(x5)
x1.set_downstream(x6)
x2.set_downstream(x4)
x2.set_downstream(x5)
x2.set_downstream(x6)
x3.set_downstream(x4)
x3.set_downstream(x5)
x3.set_downstream(x6)

也可以在传统运算符/传感器和 XComArg 任务之间进行混合

cross_downstream(from_tasks=[t1, x2(), t3], to_tasks=[x1(), t2, x3()])

等同于

t1 ---> x1
   \ /
x2 -X -> t2
   / \
t3 ---> x3
x1 = x1()
x2 = x2()
x3 = x3()
t1.set_downstream(x1)
t1.set_downstream(t2)
t1.set_downstream(x3)
x2.set_downstream(x1)
x2.set_downstream(t2)
x2.set_downstream(x3)
t3.set_downstream(x1)
t3.set_downstream(t2)
t3.set_downstream(x3)
参数:
  • from_tasks (collections.abc.Sequence[airflow.sdk.definitions._internal.abstractoperator.DependencyMixin]) – 要开始的任务或 XComArgs 列表。

  • to_tasks (airflow.sdk.definitions._internal.abstractoperator.DependencyMixin | collections.abc.Sequence[airflow.sdk.definitions._internal.abstractoperator.DependencyMixin]) – 要设置为下游依赖项的任务或 XComArgs 列表。

边缘与标签

class airflow.sdk.EdgeModifier(label=None)

表示要在两个任务/运算符之间添加的边缘信息的类。

具有简写工厂函数,如 Label(“hooray”)。

当前实现支持

t1 >> Label(“Success route”) >> t2 t2 << Label(“Success route”) << t2

请注意,由于可能在任一方向上使用,因此它会等待双方都声明后再在双方之间建立实际连接,如果添加了多个向上/向下操作,则会逐步进行。

这个类与 EdgeInfo 相关——EdgeModifier 是您用来向(可能多个)边缘添加信息的 Python 对象,而 EdgeInfo 是一个特定边缘的信息表示。

参数:

label (str | None)

class airflow.sdk.Label(label)

创建一个在边缘上设置人类可读标签的 EdgeModifier。

参数:

label (str)

资产

class airflow.sdk.Asset(name: str, uri: str | airflow.sdk.io.path.ObjectStoragePath, *, group: str = ..., extra: dict[str, pydantic.types.JsonValue] | None = None, watchers: list[AssetWatcher] = ...)

工作流之间数据资产依赖关系的表示。

class airflow.sdk.AssetAlias

资产别名的表示。

资产别名可用于在任务执行时创建资产。

class airflow.sdk.AssetAll(*objects)

用于在“与”关系中组合资产计划引用。

参数:

objects (BaseAsset)

class airflow.sdk.AssetAny(*objects)

用于在“或”关系中组合资产计划引用。

参数:

objects (BaseAsset)

class airflow.sdk.AssetWatcher

资产观察者的表示。名称唯一标识该监视。

class airflow.sdk.Metadata

要附加到 AssetEvent 的元数据。

时间表(Timetables)

class airflow.sdk.AssetOrTimeSchedule

将基于时间的调度与基于事件的调度相结合。

参数:
  • assets – 资产或资产列表,使用基于事件的调度时格式与 DAG(schedule=...) 中的格式相同。这用于评估基于事件的调度。

  • timetable – 用于评估基于时间的调度的时间表实例。

class airflow.sdk.CronDataIntervalTimetable

使用 cron 表达式调度数据区间的时间表。

这对应于 schedule=<cron>,其中 <cron> 是五/六段表示法或 cron_presets 之一。

该实现扩展了 croniter 以增加时区感知能力。这是因为 croniter 仅适用于原始时间戳,无法在确定下一个/上一个时间时考虑夏令时。

使用此类等同于直接提供 cron 表达式

不要在此传入 @once;请使用 OnceTimetable 替代。

class airflow.sdk.CronTriggerTimetable

根据 cron 表达式触发 Dag 运行的时间表。

这与 CronDataIntervalTimetable 不同,后者中 cron 表达式指定 DAG 运行的数据区间。使用此时间表,数据区间与 cron 表达式独立指定。同样出于同样的原因,此时间表会在周期开始时立即启动 DAG 运行(类似于 POSIX cron),而无需等待一个数据区间过去。

不要在此传入 @once;请使用 OnceTimetable 替代。

参数:
  • cron – 定义何时运行的 cron 字符串

  • timezone – 用于解释 cron 字符串的时区

  • interval – 定义数据区间开始的 timedelta。默认为 0。

run_immediately 控制(如果未向 Dag 提供 start_time)何时应调度 Dag 的第一次运行。如果此 Dag 已存在运行,则无效。

  • 如果为 True,则始终立即运行最近可能发生的 Dag 运行。

  • 如果为 False,则等待直到未来的下一个预定时间运行。

  • 如果传递了 timedelta,则如果该运行的 data_interval_end 在现在后的 timedelta 时间内,将运行最近可能发生的 Dag 运行。

  • 如果为 None,则 timedelta 计算为最近一次过去的预定时间与下一次预定时间之间时间的 10%。例如,如果每小时运行一次,如果自上次运行时间以来过去的时间少于 6 分钟,则这将运行上一次,否则它将等待直到下一个小时。

class airflow.sdk.CronPartitionTimetable(cron, *, timezone, run_offset=None, run_immediately=False, key_format='%Y-%m-%dT%H:%M:%S')

根据 cron 表达式触发 Dag 运行的时间表。

为分区键创建运行。

cron 表达式确定运行日期的序列。分区日期根据 run_offset 从这些日期派生而来。然后使用分区日期格式化分区键。

run_offset 为 1 意味着 partition_date 将在运行日期之后的一个 cron 区间;负数意味着 partition_date 将在运行日期之前的一个 cron 区间。

参数:
  • cron (str) – 定义何时运行的 cron 字符串

  • timezone (str | pendulum.tz.timezone.Timezone | pendulum.tz.timezone.FixedTimezone) – 用于解释 cron 字符串的时区

  • run_offset (int | datetime.timedelta | dateutil.relativedelta.relativedelta | None) – 确定要为哪个分区日期运行的整数偏移量。分区键将从分区日期派生。

  • key_format (str) – 如何将分区日期转换为字符串分区键。

  • run_immediately (bool | datetime.timedelta)

run_immediately 控制(如果未向 Dag 提供 start_time)何时应调度 Dag 的第一次运行。如果此 Dag 已存在运行,则无效。

  • 如果为 True,则始终立即运行最近可能发生的 Dag 运行。

  • 如果为 False,则等待直到未来的下一个预定时间运行。

  • 如果传递了 timedelta,则如果该运行的 data_interval_end 在现在后的 timedelta 时间内,将运行最近可能发生的 Dag 运行。

  • 如果为 None,则 timedelta 计算为最近一次过去的预定时间与下一次预定时间之间时间的 10%。例如,如果每小时运行一次,如果自上次运行时间以来过去的时间少于 6 分钟,则这将运行上一次,否则它将等待直到下一个小时。

# todo: AIP-76 讨论我们如何实现分区的自动再处理 # todo: AIP-76 我们可以允许整数 + 基于时间的元组

class airflow.sdk.DeltaDataIntervalTimetable

使用时间增量调度数据区间的时间表。

这对应于 schedule=<delta>,其中 <delta>datetime.timedeltadateutil.relativedelta.relativedelta 实例。

class airflow.sdk.DeltaTriggerTimetable

根据 cron 表达式触发 DAG 运行的时间表。

这与 DeltaDataIntervalTimetable 不同,后者中增量值指定 DAG 运行的数据区间。使用此时间表,数据区间是独立指定的。同样出于同样的原因,此时间表会在周期开始时立即启动 DAG 运行,而无需等待一个数据区间过去。

参数:
  • delta – 每次运行之间等待的时间。

  • interval – 每次运行的数据区间。默认为 0。

class airflow.sdk.EventsTimetable(event_dates, *, restrict_to_events=False, presorted=False, description=None)

在特定列出的日期时间调度 DAG 运行的时间表。

适用于可预测但真正不规则的调度,例如体育赛事,或针对国定假日进行调度。

参数:
  • event_dates (collections.abc.Iterable[pendulum.DateTime]) – DAG 运行的日期时间列表。重复项将被忽略。这必须是有限且合理的大小,因为它将被全部加载。

  • restrict_to_events (bool) – 手动运行应该使用最近的事件还是当前时间

  • presorted (bool) – 如果为 True,则假定 event_dates 按升序排列。为较大的 event_dates 列表提供适度的性能改进。

  • description (str | None) – 在 UI 中显示的时间表名称。如果未明确提供(或为 None),UI 将显示“X Events”,其中 X 是 event_dates 的长度。

class airflow.sdk.MultipleCronTriggerTimetable(*crons, timezone, interval=datetime.timedelta(), run_immediately=False)

根据多个 cron 表达式触发 Dag 运行的时间表。

这在底层结合了多个 CronTriggerTimetable 实例,并且只要其中一个时间表想要触发运行,就会触发 Dag 运行。

对于任何给定的时间,最多只触发一次运行,即使有多个时间表同时触发也是如此。

参数:
  • crons (str)

  • timezone (str | pendulum.tz.timezone.Timezone | pendulum.tz.timezone.FixedTimezone)

  • interval (datetime.timedelta | dateutil.relativedelta.relativedelta)

  • run_immediately (bool | datetime.timedelta)

class airflow.sdk.PartitionedAssetTimetable

监听分区资产的资产驱动时间表。

分区映射器(Partition Mapper)

class airflow.sdk.PartitionMapper

基础分区映射器类。

将来自资产事件的键映射到目标 dag 运行分区。

class airflow.sdk.ChainMapper(mapper0, mapper1, /, *mappers)

顺序应用多个映射器的分区映射器。

参数:
  • mapper0 (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • mapper1 (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • mappers (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

class airflow.sdk.IdentityMapper

不更改键的分区映射器。

class airflow.sdk.StartOfHourMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到小时。

参数:
  • input_format (str)

  • output_format (str | None)

class airflow.sdk.StartOfDayMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到天。

参数:
  • input_format (str)

  • output_format (str | None)

class airflow.sdk.StartOfWeekMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到周。

参数:
  • input_format (str)

  • output_format (str | None)

class airflow.sdk.StartOfMonthMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到月。

参数:
  • input_format (str)

  • output_format (str | None)

class airflow.sdk.StartOfQuarterMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到季度。

参数:
  • input_format (str)

  • output_format (str | None)

class airflow.sdk.StartOfYearMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到年。

参数:
  • input_format (str)

  • output_format (str | None)

class airflow.sdk.ProductMapper(mapper0, mapper1, /, *mappers, delimiter='|')

将多个映射器组合成多维键的分区映射器。

参数:
  • mapper0 (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • mapper1 (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • mappers (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • delimiter (str)

class airflow.sdk.AllowedKeyMapper(allowed_keys)

根据一组允许的键验证键的分区映射器。

参数:

allowed_keys (list[str])

I/O 助手

class airflow.sdk.ObjectStoragePath(*args, protocol=None, conn_id=None, **storage_options)

用于对象存储的类路径对象。

参数:
  • args (upath.types.JoinablePathLike)

  • protocol (str | None)

  • conn_id (str | None)

  • storage_options (Any)

执行时间组件

Context

class airflow.sdk.Context

用于任务呈现的 Jinja2 模板上下文。

Context 对象表示任务可用的执行时间上下文。它对应于任务执行期间暴露给 Jinja 模板的相同上下文。

有关可用上下文变量(如 dag_runtask_instancelogical_date 等)的完整列表,请参阅 模板参考

日志记录

airflow.sdk.log.mask_secret(secret, name=None)

在任务进程和管理进程中屏蔽机密信息。

对于从后端(Vault、环境变量等)加载的机密信息,这可确保它们在任务子进程和管理器的日志输出中都被屏蔽。在同步和异步上下文中均可安全工作。

参数:
  • secret (JsonValue)

  • name (str | None)

返回类型:

可观测性

其他所有内容

class airflow.sdk.AirflowSDKConfigParser(default_config=None, *args, **kwargs)

扩展共享解析器的 SDK 配置解析器。

在第一阶段,这会读取 Core 的 config.yml 并可选择读取 airflow.cfg。最终,SDK 将拥有自己的仅包含与创作相关配置的 config.yml。

参数:

default_config (str | None)

expand_all_configuration_values()

使用 SDK 特定的扩展变量扩展所有配置值。

load_test_config()

使用测试配置而不是 Airflow 默认值。

单元测试从 unit_tests.cfg 加载值以确保一致的行为。实际上我们不应该需要这样做,但这只是暂时的,旨在帮助修复使用 dag_maker 并依赖少量配置的测试。

SDK 不会扩展模板变量(FERNET_KEY、JWT_SECRET_KEY 等),因为它不使用需要扩展的配置字段。

remove_all_read_configurations()

删除所有已读取的配置,仅在配置中保留默认值。

class airflow.sdk.AllowedKeyMapper(allowed_keys)

根据一组允许的键验证键的分区映射器。

参数:

allowed_keys (list[str])

allowed_keys
class airflow.sdk.AssetOrTimeSchedule

将基于时间的调度与基于事件的调度相结合。

参数:
  • assets – 资产或资产列表,使用基于事件的调度时格式与 DAG(schedule=...) 中的格式相同。这用于评估基于事件的调度。

  • timetable – 用于评估基于时间的调度的时间表实例。

asset_condition: airflow.sdk.definitions.asset.BaseAsset
timetable: airflow.sdk.bases.timetable.BaseTimetable
class airflow.sdk.BaseBranchOperator(*, task_id, owner=DEFAULT_OWNER, email=None, email_on_retry=DEFAULT_EMAIL_ON_RETRY, email_on_failure=DEFAULT_EMAIL_ON_FAILURE, retries=DEFAULT_RETRIES, retry_delay=DEFAULT_RETRY_DELAY, retry_exponential_backoff=0, max_retry_delay=None, start_date=None, end_date=None, depends_on_past=False, ignore_first_depends_on_past=DEFAULT_IGNORE_FIRST_DEPENDS_ON_PAST, wait_for_past_depends_before_skipping=DEFAULT_WAIT_FOR_PAST_DEPENDS_BEFORE_SKIPPING, wait_for_downstream=False, dag=None, params=None, default_args=None, priority_weight=DEFAULT_PRIORITY_WEIGHT, weight_rule=DEFAULT_WEIGHT_RULE, queue=DEFAULT_QUEUE, pool=None, pool_slots=DEFAULT_POOL_SLOTS, sla=None, execution_timeout=DEFAULT_TASK_EXECUTION_TIMEOUT, on_execute_callback=None, on_failure_callback=None, on_success_callback=None, on_retry_callback=None, on_skipped_callback=None, pre_execute=None, post_execute=None, trigger_rule=DEFAULT_TRIGGER_RULE, resources=None, run_as_user=None, map_index_template=None, max_active_tis_per_dag=None, max_active_tis_per_dagrun=None, executor=None, executor_config=None, do_xcom_push=True, multiple_outputs=False, inlets=None, outlets=None, task_group=None, doc=None, doc_md=None, doc_json=None, doc_yaml=None, doc_rst=None, task_display_name=None, logger_name=None, allow_nested_operators=True, render_template_as_native_obj=None, **kwargs)

用于创建具有分支功能的算子的基类,类似于 BranchPythonOperator。

用户应从此算子创建子类并实现函数 choose_branch(self, context)。此函数应运行确定分支所需的任何业务逻辑,并返回以下内容之一:- 单个 task_id(作为字符串)- 单个 task_group_id(作为字符串)- 包含 task_ids 和 task_group_ids 组合的列表

该算子将继续执行返回的 task_id(s) 和/或 task_group_id(s),而所有其他直接位于该算子下游的任务将被跳过。

参数:
  • task_id (str)

  • owner (str)

  • email (str | collections.abc.Sequence[str] | None)

  • email_on_retry (bool)

  • email_on_failure (bool)

  • retries (int | None)

  • retry_delay (datetime.timedelta | float)

  • retry_exponential_backoff (float)

  • max_retry_delay (datetime.timedelta | float | None)

  • start_date (datetime.datetime | None)

  • end_date (datetime.datetime | None)

  • depends_on_past (bool)

  • ignore_first_depends_on_past (bool)

  • wait_for_past_depends_before_skipping (bool)

  • wait_for_downstream (bool)

  • dag (airflow.sdk.definitions.dag.DAG | None)

  • params (collections.abc.MutableMapping[str, Any] | None)

  • default_args (dict | None)

  • priority_weight (int)

  • weight_rule (airflow.sdk.types.WeightRuleParam)

  • queue (str)

  • pool (str | None)

  • pool_slots (int)

  • sla (datetime.timedelta | None)

  • execution_timeout (datetime.timedelta | None)

  • on_execute_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_failure_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | list[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_success_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_retry_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | list[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • on_skipped_callback (None | airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback | collections.abc.Collection[airflow.sdk.definitions._internal.abstractoperator.TaskStateChangeCallback])

  • pre_execute (TaskPreExecuteHook | None)

  • post_execute (TaskPostExecuteHook | None)

  • trigger_rule (str)

  • resources (dict[str, Any] | None)

  • run_as_user (str | None)

  • map_index_template (str | None)

  • max_active_tis_per_dag (int | None)

  • max_active_tis_per_dagrun (int | None)

  • executor (str | None)

  • executor_config (dict | None)

  • do_xcom_push (bool)

  • multiple_outputs (bool)

  • inlets (Any | None)

  • outlets (Any | None)

  • task_group (airflow.sdk.definitions.taskgroup.TaskGroup | None)

  • doc (str | None)

  • doc_md (str | None)

  • doc_json (str | None)

  • doc_yaml (str | None)

  • doc_rst (str | None)

  • task_display_name (str | None)

  • logger_name (str | None)

  • allow_nested_operators (bool)

  • render_template_as_native_obj (bool | None)

  • kwargs (Any)

abstract choose_branch(context)

抽象方法,用于选择要运行的分支。

子类应实现此方法,执行选择分支所需的逻辑,并返回 task_id 或 task_id 列表。如果返回 None,则所有下游任务将被跳过。

参数:

context (airflow.sdk.definitions.context.Context) – 传递给 execute() 的上下文字典

返回类型:

str | collections.abc.Iterable[str] | None

execute(context)

在创建算子时派生。

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

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

参数:

context (airflow.sdk.definitions.context.Context)

inherits_from_skipmixin = True

用于判断运算符是否继承自 SkipMixin 或其子类(例如 BranchMixin)。

class airflow.sdk.BaseHook(logger_name=None)

Hook 的抽象基类。

Hook 被设计为与外部系统交互的接口。MySqlHook、HiveHook、PigHook 返回可处理连接并与这些系统的特定实例进行交互的对象,并公开一致的方法与它们交互。

参数:

logger_name (str | None) – Hook 用于输出日志的记录器名称。如果设置为 None(默认值),记录器名称将回退到 airflow.task.hooks.{class.__module__}.{class.__name__}(例如,DbApiHook 的记录器将是 airflow.task.hooks.airflow.providers.common.sql.hooks.sql.DbApiHook)。

async classmethod aget_connection(conn_id)

获取连接(异步),给定连接 ID。

参数:

conn_id (str) – 连接 ID

返回:

connection

返回类型:

airflow.sdk.definitions.connection.Connection

abstract get_conn()

返回此 Hook 的连接。

返回类型:

Any

classmethod get_connection(conn_id)

获取连接,给定连接 ID。

参数:

conn_id (str) – 连接 ID

返回:

connection

返回类型:

airflow.sdk.definitions.connection.Connection

classmethod get_connection_form_widgets()
返回类型:

dict[str, Any]

classmethod get_hook(conn_id, hook_params=None)

返回此连接 ID 的默认钩子。

参数:
  • conn_id (str) – 连接 ID

  • hook_params (dict | None) – 钩子参数

返回:

此连接的默认钩子

classmethod get_ui_field_behaviour()
返回类型:

dict[str, Any]

class airflow.sdk.BaseNotifier(context=None)

用于发送通知的 BaseNotifier 类。

如果实现了 async_notify,可以异步使用(推荐);如果实现了 notify,可以同步使用。

目前,DAG/任务状态更改回调在 DAG 处理器上运行,仅支持同步用法。

用法:

# 异步用法 await Notifier(context)

# 同步用法 notifier = Notifier() notifier(context)

参数:

context (airflow.sdk.definitions.context.Context | None)

abstract async async_notify(context)

发送通知(异步)。

实现此功能是在触发器(triggerer)中运行此通知程序的要求,这是使用截止日期警报(Deadline Alerts)的推荐方法。

参数:

context (airflow.sdk.definitions.context.Context) – Airflow 上下文

返回类型:

注意:当前版本中上下文不可用。

context: airflow.sdk.definitions.context.Context
abstract notify(context)

发送通知(同步)。

实现此方法是该通知器在 Dag 处理器中运行的必要条件,Dag 处理器负责执行 on_success_callbackon_failure_callback

参数:

context (airflow.sdk.definitions.context.Context) – Airflow 上下文

返回类型:

render_template_fields(context, jinja_env=None)

self.template_fields 中列出的所有属性进行模板渲染。

此操作会直接修改属性,且不可逆。

参数:
  • context (airflow.sdk.definitions.context.Context) – 包含要应用于内容的上下文值的字典。

  • jinja_env (jinja2.Environment | None) – 用于渲染的 Jinja 环境。

返回类型:

template_ext: collections.abc.Sequence[str] = ()
template_fields: collections.abc.Sequence[str] = ()
class airflow.sdk.BaseOperatorLink

定义如何获取算子链接的抽象基类。

abstract get_link(operator, *, ti_key)

指向外部系统的链接。

参数:
  • operator (airflow.sdk.BaseOperator) – 此链接关联的 Airflow 运算符对象。

  • ti_key (airflow.sdk.types.TaskInstanceKey) – 要返回链接的任务实例 ID。

返回:

指向外部系统的链接

返回类型:

str

abstract property name: str

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

返回类型:

str

operators: ClassVar[list[type[airflow.sdk.BaseOperator]]] = []

此属性将被 Airflow 插件用于查找您要为其分配此操作员链接的运算符

返回:

任务所使用的操作员类列表,您想为其创建额外链接

property xcom_key: str

用于存储此算子链接完整 “link” 的 XCom 键。

使用此键检索时,将返回完整的链接。

如果未提供,默认值为 _link_<类名>

返回类型:

str

class airflow.sdk.BaseSensorOperator(*, poke_interval=60, timeout=conf.getfloat('sensors', 'default_timeout'), soft_fail=False, mode='poke', exponential_backoff=False, max_wait=None, silent_fail=False, never_fail=False, **kwargs)

传感器算子派生自此类并继承这些属性。

传感器算子持续以特定的时间间隔执行,当满足标准时成功,并在超时时失败。

参数:
  • soft_fail (bool) – 设置为 true 可在失败时将任务标记为 SKIPPED。与 never_fail 互斥。

  • poke_interval (datetime.timedelta | float) – 作业在每次尝试之间应等待的时间。可以是 timedeltafloat 秒数。

  • timeout (datetime.timedelta | float) – 任务超时并失败前经过的时间。可以是 timedeltafloat 秒数。这不应与 BaseOperator 类的 execution_timeout 混淆。timeout 测量的是第一次 poke 到当前时间之间经过的时间(考虑到每次 poke 之间的任何重调度延迟),而 execution_timeout 检查任务的 运行 时间(不包括任何重调度延迟)。如果 modepoke(见下文),两者是等效的(因为传感器从未被重调度),而在 reschedule 模式下则不是。

  • mode (str) – 传感器的操作方式。选项为:{ poke | reschedule },默认值为 poke。当设置为 poke 时,传感器在整个执行时间内占用工作节点插槽,并在 poke 之间休眠。如果传感器的预期运行时间很短,或者需要较短的 poke 间隔,请使用此模式。请注意,在此模式下,传感器将占用一个工作节点插槽和一个池插槽,持续时间为其运行时间。当设置为 reschedule 时,如果尚未满足标准,传感器任务会在 poke 之间释放工作节点插槽,并在稍后重调度。如果预期满足标准前的时间较长,请使用此模式。poke 间隔应超过一分钟,以防止调度程序负载过重。

  • exponential_backoff (bool) – 通过使用指数退避算法,允许在 poke 之间逐步延长等待时间

  • max_wait (datetime.timedelta | float | None) – poke 之间的最大等待间隔,可以是 timedeltafloat 秒数

  • silent_fail (bool) – 如果为 true,且 poke 方法引发了不同于 AirflowSensorTimeout、AirflowTaskTimeout、AirflowSkipException 和 AirflowFailException 的异常,则传感器将记录错误并继续执行。否则,传感器任务失败,并可根据提供的 retries 参数进行重试。

  • never_fail (bool) – 如果为 true,且 poke 方法引发异常,则跳过传感器。与 soft_fail 互斥。

execute(context)

在创建算子时派生。

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

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

参数:

context (airflow.sdk.definitions.context.Context)

返回类型:

Any

exponential_backoff = False
classmethod get_serialized_fields()

字符串化的 DAG 和运算符仅包含这些字段。

max_wait = None
mode = 'poke'
never_fail = False
poke(context)

在子类中覆盖此方法。

参数:

context (airflow.sdk.definitions.context.Context)

返回类型:

bool | PokeReturnValue

poke_interval
property reschedule

定义模式为“重新调度”的传感器。

resume_execution(next_method, next_kwargs, context)

任务恢复时由任务运行器(Task Runner)调用的入口点方法(替代 execute)。

参数:
  • next_method (str)

  • next_kwargs (dict[str, Any] | None)

  • context (airflow.sdk.definitions.context.Context)

silent_fail = False
soft_fail = False
timeout: int | float
ui_color: str = '#e6f1f2'
valid_modes: collections.abc.Iterable[str] = ['poke', 'reschedule']
class airflow.sdk.BaseXCom

BaseXcom 现在是一个与 XCom 后端交互的接口。

XCOM_RETURN_KEY = 'return_value'
classmethod delete(key, task_id, dag_id, run_id, map_index=None)

删除 XCom 条目;对于自定义 XCom 后端,它会获取后端数据关联的路径并将其清除。

参数:
  • key (str)

  • task_id (str)

  • dag_id (str)

  • run_id (str)

  • map_index (int | None)

返回类型:

static deserialize_value(result)

从 str 对象反序列化 XCom 值。

返回类型:

Any

classmethod get_all(*, key, dag_id, task_id, run_id, include_prior_dates=False)

检索任务的所有 XCom 值,通常用于获取所有 map index 的值。

XComSequenceSliceResult 绝不会包含 None,如果未找到值,它将返回一个空列表。

这对于一次性获取映射任务所有 map index 的 XCom 值特别有用。

参数:
  • key (str) – XCom 的键。仅返回具有此键的 XCom。

  • run_id (str) – 该任务的 DAG 运行 ID。

  • dag_id (str) – 要获取 XCom 的 DAG ID。

  • task_id (str) – 要获取 XCom 的任务 ID。

  • include_prior_dates (bool) – 如果为 False(默认),仅返回指定 DAG 运行的 XCom。如果为 True,则返回最近匹配的 XCom,无论其属于哪个运行。

返回:

如果找到,返回所有 XCom 值的列表。

返回类型:

Any

classmethod get_one(*, key, dag_id, task_id, run_id, map_index=None, include_prior_dates=False)

检索一个 XCom 值,可选地满足特定条件。

此方法返回“完整”的 XCom 值(即使用 XCom 后端的 deserialize_value)。

如果没有结果,则返回 None。如果多个 XCom 条目符合条件,则返回其中任意一个。

另请参阅

如果您已有结构化的 TaskInstance 或 TaskInstanceKey 对象,get_value() 是一个方便的函数。

参数:
  • run_id (str) – 该任务的 DAG 运行 ID。

  • dag_id (str) – 仅从此 DAG 获取 XCom。传递 None(默认)以移除过滤器。

  • task_id (str) – 仅从具有匹配 ID 的任务获取 XCom。传递 None(默认)以移除过滤器。

  • map_index (int | None) – 仅从具有匹配 ID 的任务获取 XCom。传递 None(默认)以移除过滤器。

  • key (str) – XCom 的键。如果提供,仅返回匹配键的 XCom。传递 None(默认)以移除过滤器。

  • include_prior_dates (bool) – 如果为 False(默认),仅返回指定 DAG 运行的 XCom。如果为 True,则返回最近匹配的 XCom,无论其属于哪个运行。

返回类型:

Any | None

classmethod get_value(*, ti_key, key)

检索任务实例的 XCom 值。

此方法返回“完整”的 XCom 值(即使用 XCom 后端的 deserialize_value)。

如果没有结果,则返回 None。如果多个 XCom 条目符合条件,则返回其中任意一个。

参数:
  • ti_key (TIKeyProtocol) – 用于查找 XCom 的 TaskInstanceKey。

  • key (str) – XCom 的键。如果提供,仅返回匹配键的 XCom。传递 None(默认)以移除过滤器。

返回类型:

Any

classmethod purge(xcom, *args)

从底层存储实现中清除 XCom 条目。

参数:

xcom (airflow.sdk.execution_time.comms.XComResult)

返回类型:

static serialize_value(value, *, key=None, task_id=None, dag_id=None, run_id=None, map_index=None)

将 XCom 值序列化为 JSON 字符串。

参数:
  • value (Any)

  • key (str | None)

  • task_id (str | None)

  • dag_id (str | None)

  • run_id (str | None)

  • map_index (int | None)

返回类型:

str

classmethod set(key, value, *, dag_id, task_id, run_id, map_index=-1, _mapped_length=None)

存储 XCom 值。

参数:
  • key (str) – 存储 XCom 的键。

  • value (Any) – 要存储的 XCom 值。

  • dag_id (str) – DAG ID。

  • task_id (str) – 任务 ID。

  • run_id (str) – 该任务的 DAG 运行 ID。

  • map_index (int) – 可选的 map index,用于为映射任务指定 XCom。默认值为 -1(非映射任务设置)。

  • _mapped_length (int | None)

返回类型:

class airflow.sdk.BranchMixIn(context=None)

将分支处理为单行代码的实用辅助类。

do_branch(context, branches_to_execute)

实现分支处理(包括日志记录)。

参数:
  • context (airflow.sdk.definitions.context.Context)

  • branches_to_execute (str | collections.abc.Iterable[str] | None)

返回类型:

str | collections.abc.Iterable[str] | None

class airflow.sdk.ChainMapper(mapper0, mapper1, /, *mappers)

顺序应用多个映射器的分区映射器。

参数:
  • mapper0 (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • mapper1 (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • mappers (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

mappers
class airflow.sdk.Connection(*, conn_id: str, uri: str)

到外部数据源的连接。

参数:
  • conn_id – 连接 ID。

  • conn_type – 连接类型。

  • description – 连接描述。

  • host – 主机。

  • login – 登录名。

  • password – 密码。

  • schema – 模式。

  • port – 端口号。

  • extra – 额外元数据。私钥/SSH 密钥等非标准数据可保存于此。JSON 编码对象。

  • uri – 描述连接参数的 URI 地址。

EXTRA_KEY = '__extra__'
as_json()

将 Connection 转换为 JSON 字符串对象。

返回类型:

str

async classmethod async_get(conn_id)
参数:

conn_id (str)

返回类型:

Any

conn_id: str
conn_type: str | None = None
description: str | None = None
extra: str | None = None
property extra_dejson: dict

通过反序列化 JSON 返回 extra 属性。

返回类型:

dict

classmethod from_json(value, conn_id=None)
返回类型:

Connection

classmethod from_uri(uri, conn_id)

根据 URI 字符串创建 Connection。

参数:
  • uri (str) – 要解析的 URI 字符串

  • conn_id (str) – 分配给连接的 Connection ID

返回:

Connection 对象

返回类型:

Connection

classmethod get(conn_id)
参数:

conn_id (str)

返回类型:

Any

get_extra_dejson()

将 extra 属性反序列化为 JSON。

返回类型:

dict

get_hook(*, hook_params=None)

根据 conn_type 返回 Hook。

get_uri()

生成并返回 URI 格式的连接。

返回类型:

str

host: str | None = None
login: str | None = None
password: str | None = None
port: int | None = None
schema: str | None = None
to_dict(*, prune_empty=False, validate=True)

将 Connection 转换为可 JSON 序列化的字典。

参数:
  • prune_empty (bool) – 是否移除空值。

  • validate (bool) – 验证字典是否可 JSON 序列化

返回类型:

dict[str, Any]

class airflow.sdk.Context

用于任务呈现的 Jinja2 模板上下文。

conn: Any
dag_run: airflow.sdk.types.DagRunProtocol
data_interval_end: NotRequired[pendulum.DateTime | None]
data_interval_start: NotRequired[pendulum.DateTime | None]
ds: str
ds_nodash: str
exception: NotRequired[None | str | BaseException]
expanded_ti_count: NotRequired[int | None]
inlet_events: airflow.sdk.execution_time.context.InletEventsAccessors
inlets: list
logical_date: pendulum.DateTime
macros: Any
map_index_template: NotRequired[str | None]
outlet_events: airflow.sdk.types.OutletEventAccessorsProtocol
outlets: list
params: dict[str, Any]
prev_data_interval_end_success: NotRequired[pendulum.DateTime | None]
prev_data_interval_start_success: NotRequired[pendulum.DateTime | None]
prev_end_date_success: NotRequired[pendulum.DateTime | None]
prev_start_date_success: NotRequired[pendulum.DateTime | None]
reason: NotRequired[str | None]
run_id: str
task: airflow.sdk.bases.operator.BaseOperator | airflow.sdk.types.Operator
task_instance: airflow.sdk.types.RuntimeTaskInstanceProtocol
task_instance_key_str: str
task_reschedule_count: int
templates_dict: NotRequired[dict[str, Any] | None]
test_mode: bool
ti: airflow.sdk.types.RuntimeTaskInstanceProtocol
triggering_asset_events: Any
try_number: NotRequired[int | None]
ts: str
ts_nodash: str
ts_nodash_with_tz: str
var: Any
class airflow.sdk.CronDataIntervalTimetable

使用 cron 表达式调度数据区间的时间表。

这对应于 schedule=<cron>,其中 <cron> 是五/六段表示法或 cron_presets 之一。

该实现扩展了 croniter 以增加时区感知能力。这是因为 croniter 仅适用于原始时间戳,无法在确定下一个/上一个时间时考虑夏令时。

使用此类等同于直接提供 cron 表达式

不要在此传入 @once;请使用 OnceTimetable 替代。

class airflow.sdk.CronPartitionTimetable(cron, *, timezone, run_offset=None, run_immediately=False, key_format='%Y-%m-%dT%H:%M:%S')

根据 cron 表达式触发 Dag 运行的时间表。

为分区键创建运行。

cron 表达式确定运行日期的序列。分区日期根据 run_offset 从这些日期派生而来。然后使用分区日期格式化分区键。

run_offset 为 1 意味着 partition_date 将在运行日期之后的一个 cron 区间;负数意味着 partition_date 将在运行日期之前的一个 cron 区间。

参数:
  • cron (str) – 定义何时运行的 cron 字符串

  • timezone (str | pendulum.tz.timezone.Timezone | pendulum.tz.timezone.FixedTimezone) – 用于解释 cron 字符串的时区

  • run_offset (int | datetime.timedelta | dateutil.relativedelta.relativedelta | None) – 确定要为哪个分区日期运行的整数偏移量。分区键将从分区日期派生。

  • key_format (str) – 如何将分区日期转换为字符串分区键。

  • run_immediately (bool | datetime.timedelta)

run_immediately 控制(如果未向 Dag 提供 start_time)何时应调度 Dag 的第一次运行。如果此 Dag 已存在运行,则无效。

  • 如果为 True,则始终立即运行最近可能发生的 Dag 运行。

  • 如果为 False,则等待直到未来的下一个预定时间运行。

  • 如果传递了 timedelta,则如果该运行的 data_interval_end 在现在后的 timedelta 时间内,将运行最近可能发生的 Dag 运行。

  • 如果为 None,则 timedelta 计算为最近一次过去的预定时间与下一次预定时间之间时间的 10%。例如,如果每小时运行一次,如果自上次运行时间以来过去的时间少于 6 分钟,则这将运行上一次,否则它将等待直到下一个小时。

# todo: AIP-76 讨论我们如何实现分区的自动再处理 # todo: AIP-76 我们可以允许整数 + 基于时间的元组

key_format: str = '%Y-%m-%dT%H:%M:%S'
run_offset: int | datetime.timedelta | dateutil.relativedelta.relativedelta | None = None
class airflow.sdk.CronTriggerTimetable

根据 cron 表达式触发 Dag 运行的时间表。

这与 CronDataIntervalTimetable 不同,后者中 cron 表达式指定 DAG 运行的数据区间。使用此时间表,数据区间与 cron 表达式独立指定。同样出于同样的原因,此时间表会在周期开始时立即启动 DAG 运行(类似于 POSIX cron),而无需等待一个数据区间过去。

不要在此传入 @once;请使用 OnceTimetable 替代。

参数:
  • cron – 定义何时运行的 cron 字符串

  • timezone – 用于解释 cron 字符串的时区

  • interval – 定义数据区间开始的 timedelta。默认为 0。

run_immediately 控制(如果未向 Dag 提供 start_time)何时应调度 Dag 的第一次运行。如果此 Dag 已存在运行,则无效。

  • 如果为 True,则始终立即运行最近可能发生的 Dag 运行。

  • 如果为 False,则等待直到未来的下一个预定时间运行。

  • 如果传递了 timedelta,则如果该运行的 data_interval_end 在现在后的 timedelta 时间内,将运行最近可能发生的 Dag 运行。

  • 如果为 None,则 timedelta 计算为最近一次过去的预定时间与下一次预定时间之间时间的 10%。例如,如果每小时运行一次,如果自上次运行时间以来过去的时间少于 6 分钟,则这将运行上一次,否则它将等待直到下一个小时。

interval: datetime.timedelta | dateutil.relativedelta.relativedelta
run_immediately: bool | datetime.timedelta
class airflow.sdk.DeltaDataIntervalTimetable

使用时间增量调度数据区间的时间表。

这对应于 schedule=<delta>,其中 <delta>datetime.timedeltadateutil.relativedelta.relativedelta 实例。

class airflow.sdk.DeltaTriggerTimetable

根据 cron 表达式触发 DAG 运行的时间表。

这与 DeltaDataIntervalTimetable 不同,后者中增量值指定 DAG 运行的数据区间。使用此时间表,数据区间是独立指定的。同样出于同样的原因,此时间表会在周期开始时立即启动 DAG 运行,而无需等待一个数据区间过去。

参数:
  • delta – 每次运行之间等待的时间。

  • interval – 每次运行的数据区间。默认为 0。

interval: datetime.timedelta | dateutil.relativedelta.relativedelta
class airflow.sdk.EdgeModifier(label=None)

表示要在两个任务/运算符之间添加的边缘信息的类。

具有简写工厂函数,如 Label(“hooray”)。

当前实现支持

t1 >> Label(“Success route”) >> t2 t2 << Label(“Success route”) << t2

请注意,由于可能在任一方向上使用,因此它会等待双方都声明后再在双方之间建立实际连接,如果添加了多个向上/向下操作,则会逐步进行。

这个类与 EdgeInfo 相关——EdgeModifier 是您用来向(可能多个)边缘添加信息的 Python 对象,而 EdgeInfo 是一个特定边缘的信息表示。

参数:

label (str | None)

add_edge_info(dag, upstream_id, downstream_id)

为此特定任务对在 DAG 上添加或更新任务信息。

既可通过上述关系触发器方法调用,也可直接由运算符中的 set_upstream/set_downstream 调用。

参数:
  • dag (airflow.sdk.definitions.dag.DAG)

  • upstream_id (str)

  • downstream_id (str)

label = None
property leaves

叶子节点列表——仅有上游依赖的节点。

即该子图的“终点”

property roots

根节点列表——没有上游依赖的节点。

即该子图的“起点”

set_downstream(other, edge_modifier=None)

将给定的任务/列表设置到下游属性上,然后尝试解析关系。

提供此项还可以通过 DependencyMixin 提供 >> 操作符。

参数:
  • other (airflow.sdk.definitions._internal.mixins.DependencyMixin | collections.abc.Sequence[airflow.sdk.definitions._internal.mixins.DependencyMixin])

  • edge_modifier (EdgeModifier | None)

set_upstream(other, edge_modifier=None)

将给定的任务/列表设置到上游属性上,然后尝试解析关系。

提供此项还可以通过 DependencyMixin 提供 << 操作符。

参数:
  • other (airflow.sdk.definitions._internal.mixins.DependencyMixin | collections.abc.Sequence[airflow.sdk.definitions._internal.mixins.DependencyMixin])

  • edge_modifier (EdgeModifier | None)

update_relative(other, upstream=True, edge_modifier=None)

如果我们不是关系的“主要”一方,则更新相对关系;逻辑保持不变。

参数:
  • other (airflow.sdk.definitions._internal.mixins.DependencyMixin)

  • upstream (bool)

  • edge_modifier (EdgeModifier | None)

返回类型:

class airflow.sdk.EventsTimetable(event_dates, *, restrict_to_events=False, presorted=False, description=None)

在特定列出的日期时间调度 DAG 运行的时间表。

适用于可预测但真正不规则的调度,例如体育赛事,或针对国定假日进行调度。

参数:
  • event_dates (collections.abc.Iterable[pendulum.DateTime]) – DAG 运行的日期时间列表。重复项将被忽略。这必须是有限且合理的大小,因为它将被全部加载。

  • restrict_to_events (bool) – 手动运行应该使用最近的事件还是当前时间

  • presorted (bool) – 如果为 True,则假定 event_dates 按升序排列。为较大的 event_dates 列表提供适度的性能改进。

  • description (str | None) – 在 UI 中显示的时间表名称。如果未明确提供(或为 None),UI 将显示“X Events”,其中 X 是 event_dates 的长度。

description: str | None
event_dates: list[pendulum.DateTime]
restrict_to_events: bool
class airflow.sdk.IdentityMapper

不更改键的分区映射器。

to_downstream(key)
参数:

key (str)

返回类型:

str

airflow.sdk.Label(label)

创建一个在边缘上设置人类可读标签的 EdgeModifier。

参数:

label (str)

class airflow.sdk.Metadata

要附加到 AssetEvent 的元数据。

alias: airflow.sdk.definitions.asset.AssetAlias | None = None
extra: dict[str, pydantic.types.JsonValue]
class airflow.sdk.MultipleCronTriggerTimetable(*crons, timezone, interval=datetime.timedelta(), run_immediately=False)

根据多个 cron 表达式触发 Dag 运行的时间表。

这在底层结合了多个 CronTriggerTimetable 实例,并且只要其中一个时间表想要触发运行,就会触发 Dag 运行。

对于任何给定的时间,最多只触发一次运行,即使有多个时间表同时触发也是如此。

参数:
  • crons (str)

  • timezone (str | pendulum.tz.timezone.Timezone | pendulum.tz.timezone.FixedTimezone)

  • interval (datetime.timedelta | dateutil.relativedelta.relativedelta)

  • run_immediately (bool | datetime.timedelta)

timetables: list[CronTriggerTimetable]
class airflow.sdk.ObjectStoragePath(*args, protocol=None, conn_id=None, **storage_options)

用于对象存储的类路径对象。

参数:
  • args (upath.types.JoinablePathLike)

  • protocol (str | None)

  • conn_id (str | None)

  • storage_options (Any)

__version__: ClassVar[int] = 1
property bucket: str
返回类型:

str

checksum()

返回此路径下文件的校验和。

返回类型:

int

property conn_id: str | None

返回此路径的连接 ID。

返回类型:

str | None

property container: str
返回类型:

str

copy(dst, recursive=False, **kwargs)

将文件从当前路径复制到另一个位置。

对于远程到远程的复制,目标使用的键将与源相同。因此,s3://src_bucket/foo/bar 将被复制到 gcs://dst_bucket/foo/bar,而不是 gcs://dst_bucket/bar。

参数:
  • dst (str | ObjectStoragePath) – 目标路径

  • recursive (bool) – 如果为 True,则递归复制目录。

返回类型:

kwargs: 传递给底层实现的其他关键字参数。

copy_into(target_dir, recursive=False, **kwargs)

将文件从当前路径复制到另一个目录中。

参数:
  • target_dir (str | ObjectStoragePath) – 目标目录

  • recursive (bool) – 如果为 True,则递归复制目录。

返回类型:

kwargs: 传递给底层实现的其他关键字参数。

classmethod cwd()
返回类型:

Self

classmethod deserialize(data, version)
参数:
  • data (dict)

  • version (int)

返回类型:

ObjectStoragePath

property fs: fsspec.AbstractFileSystem

使用 Airflow 的附加机制返回此路径的文件系统。

返回类型:

fsspec.AbstractFileSystem

classmethod home()
返回类型:

Self

property key: str
返回类型:

str

move(path, recursive=False, **kwargs)

将文件从当前路径移动到另一个位置。

参数:
  • path (str | ObjectStoragePath) – 目标路径

  • recursive (bool) – 如果为 True,则递归移动目录。

返回类型:

kwargs: 传递给底层实现的其他关键字参数。

move_into(target_dir, recursive=False, **kwargs)

将文件从当前路径移动到另一个目录中。

参数:
  • target_dir (str | ObjectStoragePath) – 目标目录

  • recursive (bool) – 如果为 True,则递归移动目录。

返回类型:

kwargs: 传递给底层实现的其他关键字参数。

property namespace: str
返回类型:

str

open(mode='r', **kwargs)

打开当前路径指向的文件。

read_block(offset, length, delimiter=None)

读取一块字节。

从文件的 offset 开始,读取 length 字节。如果设置了 delimiter,则我们确保读取的起点和终点位于 offsetoffset + length 位置之后的定界符边界处。如果 offset 为零,则我们从零开始。返回的字节串将包含结束定界符字符串。

如果 offset+length 超出文件末尾,则读取至 EOF。

参数:
  • offset (int) – 字节偏移量,从此处开始读取

  • length (int) – 读取的字节数。如果为 None,则读取至末尾。

  • delimiter – bytes (可选) 确保读取在定界符字节串处开始和停止

示例

# Read the first 13 bytes (no delimiter)
>>> read_block(0, 13)
b'Alice, 100\nBo'

# Read first 13 bytes, but force newline boundaries
>>> read_block(0, 13, delimiter=b"\n")
b'Alice, 100\nBob, 200\n'

# Read until EOF, but only stop at newline
>>> read_block(0, None, delimiter=b"\n")
b'Alice, 100\nBob, 200\nCharlie, 300'

参见

fsspec.utils.read_block()

replace(target)

将此路径重命名为目标路径,如果目标路径存在则覆盖。

目标路径可以是绝对路径或相对路径。相对路径是相对于当前工作目录解释的,而不是相对于 Path 对象所在的目录。

返回指向目标路径的新 Path 实例。

返回类型:

Self

root_marker: ClassVar[str] = '/'
samefile(other_path)

返回 other_path 是否与当前文件相同。

参数:

other_path (Any)

返回类型:

bool

samestore(other)
参数:

other (Any)

返回类型:

bool

sep: ClassVar[str] = '/'
serialize()
返回类型:

dict[str, Any]

sign(expiration=100, **kwargs)

创建代表给定路径的签名 URL。

某些实现允许生成临时 URL,作为委托凭据的一种方式。

参数:
  • path – 文件系统上的路径

  • expiration (int) – 启用该 URL 的秒数(如果支持)

返回 URL:

str 签名后的 URL

抛出:

NotImplementedError – 如果该方法未针对某存储实现

size()

该路径文件的字节大小。

返回类型:

int

stat()

调用 stat 并返回结果。

返回类型:

airflow.sdk.io.stat.stat_result

ukey()

用于标识文件属性的哈希,以判断文件是否发生变化。

返回类型:

str

class airflow.sdk.Param(default=NOTSET, description=None, source=None, **kwargs)

用于保存 Param 默认值和执行验证的规则集的类。

如果没有规则集,它始终验证并返回默认值。

参数:
  • default (Any) – 此 Param 对象持有的值

  • description (str | None) – Param 的可选帮助文本

  • schema – Param 的验证模式,如果未提供,则除 default 和 description 之外的所有 kwargs 将构成模式

  • source (Literal['dag', 'task'] | None)

CLASS_IDENTIFIER = '__class'
__version__: ClassVar[int] = 1
description = None
static deserialize(data, version)
参数:
  • data (dict[str, Any])

  • version (int)

返回类型:

Param

dump()

将 Param 转储为字典。

返回类型:

dict

property has_value: bool
返回类型:

bool

resolve(value=NOTSET, suppress_exception=False)

运行验证并返回 Param 的最终值。

可能会在验证失败时引发 ValueError,或者在未传递值且不存在现有值时引发 TypeError。我们首先检查值是否可 JSON 序列化;如果不能,则发出警告。在未来的版本中,我们将要求该值必须是可 JSON 序列化的。

参数:
  • value (Any) – 要更新为 Param 的值

  • suppress_exception (bool) – 在验证失败时是否抛出异常。如果为 true 且验证失败,返回值将为 None。

返回类型:

Any

schema
serialize()
返回类型:

dict

source = None
value
class airflow.sdk.PartitionMapper

基础分区映射器类。

将来自资产事件的键映射到目标 dag 运行分区。

class airflow.sdk.PartitionedAssetTimetable

监听分区资产的资产驱动时间表。

asset_condition: airflow.sdk.definitions.asset.BaseAsset
default_partition_mapper: airflow.sdk.definitions.partition_mappers.base.PartitionMapper
partition_mapper_config: dict[airflow.sdk.definitions.asset.BaseAsset, airflow.sdk.definitions.partition_mappers.base.PartitionMapper]
class airflow.sdk.PokeReturnValue(is_done, xcom_value=None)

poke 方法的可选返回值。

传感器可以在 poke 方法中选择返回 PokeReturnValue 类的实例。如果传感器完成时提供了 XCom 值,则 XCom 值将通过算子返回值进行推送。:param is_done: 设置为 true 以指示传感器可以停止 poking。:param xcom_value: 一个可选的 XCOM 值,将由算子返回。

参数:
  • is_done (bool)

  • xcom_value (Any | None)

is_done
xcom_value = None
class airflow.sdk.ProductMapper(mapper0, mapper1, /, *mappers, delimiter='|')

将多个映射器组合成多维键的分区映射器。

参数:
  • mapper0 (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • mapper1 (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • mappers (airflow.sdk.definitions.partition_mappers.base.PartitionMapper)

  • delimiter (str)

delimiter = '|'
mappers
class airflow.sdk.SkipMixin(context=None)

用于跳过任务实例的 Mixin。

skip(ti, tasks)

将同一 DAG 运行中的任务实例设置为已跳过(skipped)。

如果此实例具有 task_id 属性,它会将已跳过的任务 ID 存储到 XCom 中,以便 NotPreviouslySkippedDep 在清除任务时知道这些任务应被跳过。

参数:
  • ti (airflow.sdk.types.RuntimeTaskInstanceProtocol)

  • tasks (collections.abc.Iterable[airflow.sdk.definitions._internal.node.DAGNode])

skip_all_except(ti, branch_task_ids)

实现分支运算符的逻辑。

给定一个要遵循的任务 ID 或任务 ID 列表,此方法会立即跳过该运算符下游的所有其他任务。

branch_task_ids 存储到 XCom 中,以便 NotPreviouslySkippedDep 在清除任务时知道应该跳过已跳过的任务或新添加的任务。

参数:
  • ti (airflow.sdk.types.RuntimeTaskInstanceProtocol)

  • branch_task_ids (None | str | collections.abc.Iterable[str])

class airflow.sdk.StartOfDayMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到天。

参数:
  • input_format (str)

  • output_format (str | None)

default_output_format = '%Y-%m-%d'
class airflow.sdk.StartOfHourMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到小时。

参数:
  • input_format (str)

  • output_format (str | None)

default_output_format = '%Y-%m-%dT%H'
class airflow.sdk.StartOfMonthMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到月。

参数:
  • input_format (str)

  • output_format (str | None)

default_output_format = '%Y-%m'
class airflow.sdk.StartOfQuarterMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到季度。

参数:
  • input_format (str)

  • output_format (str | None)

default_output_format = '%Y-Q{quarter}'
class airflow.sdk.StartOfWeekMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到周。

参数:
  • input_format (str)

  • output_format (str | None)

default_output_format = '%Y-%m-%d (W%V)'
class airflow.sdk.StartOfYearMapper(input_format='%Y-%m-%dT%H:%M:%S', output_format=None)

将基于时间的分区键映射到年。

参数:
  • input_format (str)

  • output_format (str | None)

default_output_format = '%Y'
class airflow.sdk.Variable

一种用于以简单的键/值存储方式存储和检索任意内容或设置的通用方法。

参数:
  • key – 变量键。

  • value – 变量值。

  • description – 变量描述。

classmethod delete(key)
参数:

key (str)

返回类型:

description: str | None = None
classmethod get(key, default=NOTSET, deserialize_json=False)
参数:
  • key (str)

  • default (Any)

  • deserialize_json (bool)

key: str
classmethod set(key, value, description=None, serialize_json=False)
参数:
  • key (str)

  • value (Any)

  • description (str | None)

  • serialize_json (bool)

返回类型:

value: Any | None = None
airflow.sdk.__version__: str
airflow.sdk.chain(*tasks)

给定多个任务,构建依赖链。

此函数接受 BaseOperator(即任务)、EdgeModifiers(即标签)、XComArg、TaskGroups 或包含这些类型中任何混合的列表(或同一列表中的混合)。如果要链接两个列表,则必须确保它们具有相同的长度。

使用传统运算符/传感器

chain(t1, [t2, t3], [t4, t5], t6)

等同于

  / -> t2 -> t4 \
t1               -> t6
  \ -> t3 -> t5 /
t1.set_downstream(t2)
t1.set_downstream(t3)
t2.set_downstream(t4)
t3.set_downstream(t5)
t4.set_downstream(t6)
t5.set_downstream(t6)

使用任务装饰函数即 XComArgs

chain(x1(), [x2(), x3()], [x4(), x5()], x6())

等同于

  / -> x2 -> x4 \
x1               -> x6
  \ -> x3 -> x5 /
x1 = x1()
x2 = x2()
x3 = x3()
x4 = x4()
x5 = x5()
x6 = x6()
x1.set_downstream(x2)
x1.set_downstream(x3)
x2.set_downstream(x4)
x3.set_downstream(x5)
x4.set_downstream(x6)
x5.set_downstream(x6)

使用 TaskGroups

chain(t1, task_group1, task_group2, t2)

t1.set_downstream(task_group1)
task_group1.set_downstream(task_group2)
task_group2.set_downstream(t2)

也可以在传统运算符/传感器、EdgeModifiers、XComArg 和 TaskGroups 之间进行混合

chain(t1, [Label("branch one"), Label("branch two")], [x1(), x2()], task_group1, x3())

等同于

  / "branch one" -> x1 \
t1                      -> task_group1 -> x3
  \ "branch two" -> x2 /
x1 = x1()
x2 = x2()
x3 = x3()
label1 = Label("branch one")
label2 = Label("branch two")
t1.set_downstream(label1)
label1.set_downstream(x1)
t2.set_downstream(label2)
label2.set_downstream(x2)
x1.set_downstream(task_group1)
x2.set_downstream(task_group1)
task_group1.set_downstream(x3)

# or

x1 = x1()
x2 = x2()
x3 = x3()
t1.set_downstream(x1, edge_modifier=Label("branch one"))
t1.set_downstream(x2, edge_modifier=Label("branch two"))
x1.set_downstream(task_group1)
x2.set_downstream(task_group1)
task_group1.set_downstream(x3)
参数:

tasks (airflow.sdk.definitions._internal.abstractoperator.DependencyMixin | collections.abc.Sequence[airflow.sdk.definitions._internal.abstractoperator.DependencyMixin]) – 设置依赖项的单个任务和/或任务列表、EdgeModifiers、XComArgs 或 TaskGroups

返回类型:

airflow.sdk.chain_linear(*elements)

简化任务依赖定义。

例如:假设您想要这样的优先级

    ╭─op2─╮ ╭─op4─╮
op1─┤     ├─├─op5─┤─op7
    ╰-op3─╯ ╰-op6─╯

那么您可以这样完成

chain_linear(op1, [op2, op3], [op4, op5, op6], op7)
参数:

elements (airflow.sdk.definitions._internal.abstractoperator.DependencyMixin | collections.abc.Sequence[airflow.sdk.definitions._internal.abstractoperator.DependencyMixin]) – 运算符列表/运算符列表

airflow.sdk.conf: configuration.AirflowSDKConfigParser
airflow.sdk.cross_downstream(from_tasks, to_tasks)

将 from_tasks 中所有任务的下游依赖项设置为 to_tasks 中的所有任务。

使用传统运算符/传感器

cross_downstream(from_tasks=[t1, t2, t3], to_tasks=[t4, t5, t6])

等同于

t1 ---> t4
   \ /
t2 -X -> t5
   / \
t3 ---> t6
t1.set_downstream(t4)
t1.set_downstream(t5)
t1.set_downstream(t6)
t2.set_downstream(t4)
t2.set_downstream(t5)
t2.set_downstream(t6)
t3.set_downstream(t4)
t3.set_downstream(t5)
t3.set_downstream(t6)

使用任务装饰函数即 XComArgs

cross_downstream(from_tasks=[x1(), x2(), x3()], to_tasks=[x4(), x5(), x6()])

等同于

x1 ---> x4
   \ /
x2 -X -> x5
   / \
x3 ---> x6
x1 = x1()
x2 = x2()
x3 = x3()
x4 = x4()
x5 = x5()
x6 = x6()
x1.set_downstream(x4)
x1.set_downstream(x5)
x1.set_downstream(x6)
x2.set_downstream(x4)
x2.set_downstream(x5)
x2.set_downstream(x6)
x3.set_downstream(x4)
x3.set_downstream(x5)
x3.set_downstream(x6)

也可以在传统运算符/传感器和 XComArg 任务之间进行混合

cross_downstream(from_tasks=[t1, x2(), t3], to_tasks=[x1(), t2, x3()])

等同于

t1 ---> x1
   \ /
x2 -X -> t2
   / \
t3 ---> x3
x1 = x1()
x2 = x2()
x3 = x3()
t1.set_downstream(x1)
t1.set_downstream(t2)
t1.set_downstream(x3)
x2.set_downstream(x1)
x2.set_downstream(t2)
x2.set_downstream(x3)
t3.set_downstream(x1)
t3.set_downstream(t2)
t3.set_downstream(x3)
参数:
  • from_tasks (collections.abc.Sequence[airflow.sdk.definitions._internal.abstractoperator.DependencyMixin]) – 要开始的任务或 XComArgs 列表。

  • to_tasks (airflow.sdk.definitions._internal.abstractoperator.DependencyMixin | collections.abc.Sequence[airflow.sdk.definitions._internal.abstractoperator.DependencyMixin]) – 要设置为下游依赖项的任务或 XComArgs 列表。

airflow.sdk.literal(value)

包装一个值以确保其按原样呈现,而不对其内容应用 Jinja 模板。

专为在运算符的模板字段中使用而设计。

参数:

value (Any) – 要呈现且不应用模板的值

返回类型:

airflow.sdk.definitions._internal.templater.LiteralValue

airflow.sdk.setup: collections.abc.Callable
参数:

func (Callable)

返回类型:

Callable

airflow.sdk.task: TaskDecoratorCollection
airflow.sdk.task_group(group_id: str | None = None, prefix_group_id: bool = True, parent_group: airflow.sdk.definitions.taskgroup.TaskGroup | None = None, dag: airflow.sdk.definitions.dag.DAG | None = None, default_args: dict[str, Any] | None = None, tooltip: str = '', ui_color: str = 'CornflowerBlue', ui_fgcolor: str = '#000', add_suffix_on_collision: bool = False, group_display_name: str = '') collections.abc.Callable[[collections.abc.Callable[FParams, FReturn]], _TaskGroupFactory[FParams, FReturn]]

Python TaskGroup 装饰器。

将函数包装为 Airflow TaskGroup。当用作 @task_group() 形式时,所有参数都转发给底层的 TaskGroup 类。可用于参数化 TaskGroup。

参数:
  • python_callable – 待装饰的函数。

  • tg_kwargs – TaskGroup 对象的关键字参数。

airflow.sdk.teardown: collections.abc.Callable
参数:

on_failure_fail_dagrun (bool)

返回类型:

Callable

此条目是否有帮助?