Airflow 为 Google Cloud Platform 提供广泛支持。但请注意,大多数 Hooks 和 Operators 都位于 contrib 部分,这意味着它们处于 beta 状态,这意味着它们可以在次要版本之间进行重大更改。
客户方面的数量:
可以有尽可能多的 DAG,每个 DAG 都可以提及多个任务。建议在一个 DAG 文件中保留一个逻辑工作流,并尽量保持精简(例如配置文件)。它允许 Airflow 调度程序在每次心跳时花费更少的时间和资源来处理它们。
可以根据任意数量的配置参数动态创建 DAG(具有相同的基本代码),这在拥有大量客户端时非常有用且节省时间。
要创建新的 DAG,请在 create_dag 函数中创建一个 DAG 模板。代码可以包装在允许传入自定义参数的方法中。此外,输入参数不必存在于 dag 文件本身中。另一种生成 DAG 的常见形式是在 Variable 对象中设置值。请参考here 了解更多信息。
特定客户端配置:
您可以使用宏在运行时将动态信息传递给任务实例。可以在 here 找到所有模板中可访问的默认变量列表。
Airflow 对Jinja templating 的内置支持使用户能够传递可在模板化字段中使用的参数。
用户界面概览
如果您的 dag 需要很长时间才能加载,您可以将 airflow.cfg 中 default_dag_run_display_number 配置的值减小到更小的值。这个可配置控制在 UI 中显示的 dag 运行的数量,默认值为 25。
模块化
如果将 default_args 的字典传递给 DAG,它会将它们应用于其任何运算符。这样可以轻松地将一个通用参数应用于许多运算符,而无需多次键入。
看看例子:
from datetime import datetime, timedelta
default_args = {
'owner': 'Airflow',
'depends_on_past': False,
'start_date': datetime(2015, 6, 1),
'email': ['airflow@example.com'],
'email_on_failure': False,
'email_on_retry': False,
'retries': 1,
'retry_delay': timedelta(minutes=5),
# 'queue': 'bash_queue',
# 'pool': 'backfill',
# 'priority_weight': 10,
# 'end_date': datetime(2016, 1, 1),
}
dag = DAG('my_dag', default_args=default_args)
op = DummyOperator(task_id='dummy', dag=dag)
print(op.owner) # Airflow
有关 BaseOperator 的参数及其作用的更多信息,请参阅airflow.models.BaseOperator 文档。
性能
可以使用可以控制的变量来提高气流 DAG 性能(可以在 airflow.cfg 中设置。):
-
parallelism:控制在整个 Airflow 集群中同时运行的任务实例的数量。
-
concurrency:Airflow 调度程序将在任何给定时间为您的 DAG 运行不超过并发任务实例。并发在您的 Airflow DAG 中定义。如果您未在 DAG 上设置并发,调度程序将使用您的 airflow.cfg 中 dag_concurrency 条目中的默认值。
-
task_concurrency:此变量控制每个任务跨 dag_runs 并发运行的任务实例数。
-
max_active_runs:Airflow 调度程序在给定时间运行的 DAG 不会超过 max_active_runs DagRuns。
-
pool:此变量控制分配给池的并发运行任务实例的数量。
您可以在作曲家实例存储桶gs://composer_instance_bucket/airflow.cfg 中看到气流配置。您可以根据需要调整此配置,但请记住 Cloud Composer 已阻止某些配置。
缩放
请记住,建议节点数必须大于 3,将此数字保持在 3 以下可能会导致一些问题,如果您想升级节点数,可以使用gcloud command来指定这个值。另请注意,有一些与自动缩放相关的气流配置被阻止并且不能是overridden。
一些 Airflow 配置是为 Cloud Composer 预先配置的,您无法更改它们。
容错
请参考以下documentation。
重新执行
就像对象是类的实例一样,Airflow 任务是Operator (BaseOperator) 的实例。因此,只需传递不同的参数,编写一个“可重用”运算符并在您的管道中使用它数百次。
延迟
可以通过以下方式减少生产中的气流 DAG 调度延迟:
-
max_threads:调度程序将并行生成多个线程来调度 dag。这由 max_threads 控制,默认值为 2。
-
scheduler_heartbeat_sec:用户应考虑将 scheduler_heartbeat_sec 配置增加到更高的值(例如 60 秒),以控制气流调度程序获取心跳的频率并更新数据库中的作业条目。
请参阅以下有关最佳实践的文章:
希望对你有所帮助。