Managed Airflow(第 3 代) | Managed Airflow(第 2 代) | Managed Airflow(旧版第 1 代)
本页介绍了如何在 Managed Airflow 中启用数据沿袭集成。
数据沿袭集成简介
数据沿袭是 Knowledge Catalog的一项功能,可跟踪数据在系统中的移动方式:数据来自何处、传递到何处以及对其应用了哪些转换。
Managed Airflow 使用 apache-airflow-providers-openlineage 软件包生成沿袭事件,这些事件会发送到 Data Lineage API。
此软件包已安装在 Managed Airflow 环境中。如果您 安装此软件包的其他版本,受支持的运算符列表 可能会发生变化。我们建议您仅在必要时这样做,否则请保留预安装的软件包版本。
数据沿袭适用于与 Knowledge Catalog 区域 支持数据沿袭的区域位于同一区域的环境。
如果在 Managed Service for Apache Airflow 环境中启用了数据沿袭, Managed Service for Apache Airflow 会向 Data Lineage API 报告使用任何 受支持运算符的 DAG 的沿袭信息。如果您想报告不受支持的运算符的沿袭,还可以 发送自定义沿袭事件。
您可以使用以下方式访问沿袭信息:
- Data Lineage API
- Knowledge Catalog 中受支持条目的沿袭图。 如需了解详情,请参阅 Knowledge Catalog 文档中的 沿袭图 。
创建环境时,如果满足以下条件,系统会自动启用 数据沿袭集成:
您的项目中启用了 Data Lineage API。如需了解详情,请参阅 Knowledge Catalog 文档中的 启用 Data Lineage API。
Airflow 中未配置自定义 沿袭后端。
您可以在创建环境时停用数据沿袭集成。
Managed Service for Apache Airflow 中的功能注意事项
在以下情况下,Managed Airflow 会进行 RPC 调用以创建沿袭事件:
- Airflow 任务开始或结束时
- DAG 运行开始或结束时
如需详细了解这些实体,请参阅 Knowledge Catalog 文档中的 沿袭信息模型和 沿袭 API 参考文档。
发出的沿袭流量受 Data Lineage API 中的配额限制。 Managed Airflow 会消耗写入配额。
处理沿袭数据相关的价格受沿袭价格的影响。 请参阅数据沿袭注意事项。
Managed Service for Apache Airflow 中的性能注意事项
数据沿袭会在 Airflow 任务执行结束时报告。 平均而言,数据沿袭报告大约需要 1-2 秒。
这不会影响任务本身的性能:如果沿袭未成功报告给沿袭 API,Airflow 任务不会失败。 这不会影响主要运算符逻辑,但整个任务实例的执行时间会稍长一些,以考虑报告沿袭数据。
报告数据沿袭的环境的相关费用会略有增加,因为报告数据沿袭需要额外的时间。
合规性
数据沿袭为 VPC Service Controls 等功能提供不同的支持级别。查看 数据沿袭注意事项 ,确保支持级别符合您的环境要求。
准备工作
此功能提供不同的合规性支持 。请务必先查看特定于 Managed Service for Apache Airflow 的功能注意事项和数据沿袭功能注意事项。
Managed Airflow 2.1.2 版及更高版本支持数据沿袭集成,Airflow 版本为 2.2.5 及更高版本。
数据沿袭所需的所有 IAM 权限都已包含在 Composer Worker (
roles/composer.worker) 角色中。此角色 是环境的服务账号所需的 角色。如需详细了解数据沿袭权限,请参阅 Knowledge Catalog 文档中的 沿袭角色和权限。
检查运算符是否受支持
数据沿袭支持由运算符所在的提供程序软件包提供:
检查运算符所在的提供程序软件包的更改日志,了解添加 OpenLineage 支持的条目。
例如,BigQueryToBigQueryOperator 从
apache-airflow-providers-google11.0.0 版开始支持 OpenLineage。检查您的环境使用的提供程序软件包的版本。为此,请参阅您的环境中使用的 Managed Airflow 版本的预安装软件包列表。您还可以在 环境中安装软件包的其他版本。
此外,受支持的类
页面在 apache-airflow-providers-openlineage 文档中列出了最新
受支持的运算符。
配置数据沿袭集成
Managed Service for Apache Airflow 的数据沿袭集成按环境进行管理。这意味着启用该功能需要两个步骤:
- 在项目中启用 Data Lineage API。
- 在特定的 Managed Service for Apache Airflow 环境中启用数据沿袭集成。
在 Managed Service for Apache Airflow 中启用数据沿袭
控制台
在 Google Cloud 控制台中,前往环境 页面。
在环境列表中,点击您的环境名称。环境详情 页面会打开。
选择环境配置 标签页。
在 Knowledge Catalog 沿袭集成 部分中,点击修改 。
在 Knowledge Catalog 沿袭集成 面板中,选择启用与 Knowledge Catalog 沿袭的集成 。
点击保存 。
gcloud
使用 --enable-cloud-data-lineage-integration 实参。
gcloud composer environments update ENVIRONMENT_NAME \
--location LOCATION \
--enable-cloud-data-lineage-integration
替换以下内容:
ENVIRONMENT_NAME:您的环境的名称。LOCATION:环境所在的区域。
示例:
gcloud composer environments update example-environment \
--location us-central1 \
--enable-cloud-data-lineage-integration
在 Managed Service for Apache Airflow 中停用数据沿袭
在 Managed Service for Apache Airflow 环境中停用沿袭集成不会停用 Data Lineage API。如果您想完全停用项目的沿袭报告,请同时停用 Data Lineage API。请参阅 停用服务。
控制台
在 Google Cloud 控制台中,前往环境 页面。
在环境列表中,点击您的环境名称。环境详情 页面会打开。
选择环境配置 标签页。
在 Knowledge Catalog 沿袭集成 部分中,点击修改 。
在 Knowledge Catalog 沿袭集成 面板中,选择停用与 Knowledge Catalog 沿袭的集成 。
点击保存 。
gcloud
使用 --disable-cloud-data-lineage-integration 实参。
gcloud composer environments update ENVIRONMENT_NAME \
--location LOCATION \
--disable-cloud-data-lineage-integration
替换以下内容:
ENVIRONMENT_NAME:您的环境的名称。LOCATION:环境所在的区域。
示例:
gcloud composer environments update example-environment \
--location us-central1 \
--disable-cloud-data-lineage-integration
在受支持的运算符中发送沿袭事件
如果启用了数据沿袭,受支持的运算符会自动发送沿袭事件。您无需更改 DAG 代码。
例如,运行以下任务:
task = BigQueryInsertJobOperator(
task_id='snapshot_task',
dag=dag,
location='<dataset-location>',
configuration={
'query': {
'query': 'SELECT * FROM dataset.tableA',
'useLegacySql': False,
'destinationTable': {
'project_id': 'example-project',
'dataset_id': 'dataset',
'table_id': 'tableB',
},
}
},
)
会在 Knowledge Catalog 界面中创建以下沿袭图:
发送自定义沿袭事件
如果您想报告不受自动沿袭报告支持的运算符的沿袭,可以发送自定义沿袭事件。
例如,如需发送具有以下内容的自定义事件:
- BashOperator:修改任务定义中的
inlets或outlets形参。 - PythonOperator:修改任务定义中的
task.inlets或task.outlets形参。 - 您可以为
inlets形参使用AUTO。这会将其值设置为等于其上游任务的outlets。
以下示例演示了 inlet 和 outlet 的用法:
from airflow.composer.data_lineage.entities import BigQueryTable
from airflow.lineage import AUTO
...
bash_task = BashOperator(
task_id="bash_task",
dag=dag,
bash_command="sleep 0",
inlets=[
BigQueryTable(
project_id="example-project",
dataset_id="dataset",
table_id="table1",
)
],
outlets=[
BigQueryTable(
project_id="example-project",
dataset_id="dataset",
table_id="table2",
)
],
)
def _python_task(task):
print("Python task")
python_task = PythonOperator(
task_id="python_task",
dag=dag,
python_callable=_python_task,
inlets=[
AUTO,
BigQueryTable(
project_id="example-project",
dataset_id="dataset",
table_id="table3",
),
],
outlets=[
BigQueryTable(
project_id="example-project",
dataset_id="dataset",
table_id="table4",
)
],
)
bash_task >> python_task
因此,系统会在 Knowledge Catalog 界面中创建以下沿袭图:
在 Managed Service for Apache Airflow 中查看沿袭日志
您可以使用 Knowledge Catalog 沿袭集成 部分中环境配置 页面上的链接检查与数据沿袭相关的日志。
问题排查
如果沿袭数据未报告给沿袭 API,或者您在 Knowledge Catalog 中看不到沿袭数据,请尝试执行以下问题排查步骤:
- 确保在 Managed Service for Apache Airflow 环境的项目中启用了 Data Lineage API。
- 检查 Managed Service for Apache Airflow 环境中是否启用了数据沿袭集成。
- 检查您使用的运算符是否包含在自动沿袭报告支持中。如需了解详情,请参阅受支持的 Airflow 运算符。
- 检查 Managed Service for Apache Airflow 中的沿袭日志,了解可能存在的问题。