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 }, )
render_template_as_native_obj – 如果为 True,使用 Jinja
NativeEnvironment将模板渲染为原生 Python 类型。如果为 False,使用 JinjaEnvironment将模板渲染为字符串值。tags – 标签列表,有助于在 UI 中过滤 DAG。
owner_links – 所有者及其链接的字典,在 DAG 视图 UI 中可点击。可用作 HTTP 链接(例如指向您的 Slack 频道)或 mailto 链接。例如:
{"dag_owner": "https://airflow.org.cn/"}auto_register – 当在
with块中使用时,自动注册此 DAGfail_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, dateutil 和 random。
装饰器¶
- 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) – 重试之间的延迟,可以设置为
timedelta或float秒,后者将被转换为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) – 重试之间的最大延迟间隔,可以设置为
timedelta或float秒,后者将被转换为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 点运行每日任务,请查看TimeSensor和TimeDeltaSensor。我们建议不要使用动态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,则使用 JinjaEnvironment将模板呈现为字符串值。如果为 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) – 作业在每次尝试之间应等待的时间。可以是
timedelta或float秒数。timeout (datetime.timedelta | float) – 任务超时并失败前经过的时间。可以是
timedelta或float秒数。这不应与BaseOperator类的execution_timeout混淆。timeout测量的是第一次 poke 到当前时间之间经过的时间(考虑到每次 poke 之间的任何重调度延迟),而execution_timeout检查任务的 运行 时间(不包括任何重调度延迟)。如果mode是poke(见下文),两者是等效的(因为传感器从未被重调度),而在reschedule模式下则不是。mode (str) – 传感器的操作方式。选项为:
{ poke | reschedule },默认值为poke。当设置为poke时,传感器在整个执行时间内占用工作节点插槽,并在 poke 之间休眠。如果传感器的预期运行时间很短,或者需要较短的 poke 间隔,请使用此模式。请注意,在此模式下,传感器将占用一个工作节点插槽和一个池插槽,持续时间为其运行时间。当设置为reschedule时,如果尚未满足标准,传感器任务会在 poke 之间释放工作节点插槽,并在稍后重调度。如果预期满足标准前的时间较长,请使用此模式。poke 间隔应超过一分钟,以防止调度程序负载过重。exponential_backoff (bool) – 通过使用指数退避算法,允许在 poke 之间逐步延长等待时间
max_wait (datetime.timedelta | float | None) – poke 之间的最大等待间隔,可以是
timedelta或float秒数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.BaseOperatorLink¶
定义如何获取算子链接的抽象基类。
- 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 选项的公共接口类。
此类为处理截止日期提供了统一接口,支持计算截止日期(从数据库获取值)和固定截止日期(返回预定义的日期时间)。
用法:¶
截止日期参考示例
fixed = DeadlineReference.FIXED_DATETIME(datetime(2025, 5, 4)) logical = DeadlineReference.DAGRUN_LOGICAL_DATE queued = DeadlineReference.DAGRUN_QUEUED_AT
在 DAG 中使用
DAG( dag_id="dag_with_deadline", deadline=DeadlineAlert( reference=DeadlineReference.DAGRUN_LOGICAL_DATE, interval=timedelta(hours=1), callback=hello_callback, ), )
评估截止日期时将忽略意外参数
# 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"]
仅当在运算符开始执行后调用此方法时,当前上下文才会有值。
- 返回类型:
- 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.timedelta或dateutil.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_run、task_instance、logical_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_callback 和 on_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) – 作业在每次尝试之间应等待的时间。可以是
timedelta或float秒数。timeout (datetime.timedelta | float) – 任务超时并失败前经过的时间。可以是
timedelta或float秒数。这不应与BaseOperator类的execution_timeout混淆。timeout测量的是第一次 poke 到当前时间之间经过的时间(考虑到每次 poke 之间的任何重调度延迟),而execution_timeout检查任务的 运行 时间(不包括任何重调度延迟)。如果mode是poke(见下文),两者是等效的(因为传感器从未被重调度),而在reschedule模式下则不是。mode (str) – 传感器的操作方式。选项为:
{ poke | reschedule },默认值为poke。当设置为poke时,传感器在整个执行时间内占用工作节点插槽,并在 poke 之间休眠。如果传感器的预期运行时间很短,或者需要较短的 poke 间隔,请使用此模式。请注意,在此模式下,传感器将占用一个工作节点插槽和一个池插槽,持续时间为其运行时间。当设置为reschedule时,如果尚未满足标准,传感器任务会在 poke 之间释放工作节点插槽,并在稍后重调度。如果预期满足标准前的时间较长,请使用此模式。poke 间隔应超过一分钟,以防止调度程序负载过重。exponential_backoff (bool) – 通过使用指数退避算法,允许在 poke 之间逐步延长等待时间
max_wait (datetime.timedelta | float | None) – poke 之间的最大等待间隔,可以是
timedelta或float秒数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)
- 返回类型:
- classmethod from_uri(uri, conn_id)
根据 URI 字符串创建 Connection。
- 参数:
uri (str) – 要解析的 URI 字符串
conn_id (str) – 分配给连接的 Connection ID
- 返回:
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.timedelta或dateutil.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)
- 返回类型:
- 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,则我们确保读取的起点和终点位于offset和offset + 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)
- 返回类型:
- 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