本文档介绍了如何将 OpenLineage 与 Knowledge Catalog(以前称为 Dataplex Universal Catalog)集成,以导入和直观呈现外部系统的数据沿袭。通过使用 ProcessOpenLineageRunEvent REST API 充当 OpenLineage 消费者
,Knowledge Catalog 可让您统一自定义流水线
沿袭以及来自 Google Cloud 服务的内置沿袭。
概览
OpenLineage 是一个用于收集和 分析数据沿袭信息的开放平台。OpenLineage 使用沿袭数据的开放标准,从使用 OpenLineage API 报告运行、作业和数据集的数据流水线组件捕获沿袭事件。
通过 Data Lineage API,您可以导入 OpenLineage 事件,以便在 Knowledge Catalog 网页界面中与来自 Google Cloud 服务(例如 BigQuery、Managed Service for Apache Airflow、 Cloud Data Fusion 和 Managed Service for Apache Spark)的沿袭信息一起显示 。
如需导入使用
OpenLineage 规范的 OpenLineage 事件,
请使用
ProcessOpenLineageRunEvent
REST API 方法,并将 OpenLineage 分面映射到 Data Lineage API 属性。
OpenLineage 集成限制
支持的版本: Data Lineage API 支持 OpenLineage 主要版本 1。
API 操作: Data Lineage API 端点
ProcessOpenLineageRunEvent仅充当 OpenLineage 消息的消费者 ,而不是生产者 。借助该 API,您可以将任何符合 OpenLineage 规范的工具或系统生成的沿袭信息发送到 Knowledge Catalog。某些 Google Cloud 服务(例如 Managed Service for Apache Spark 和 Managed Airflow)包含 内置的 OpenLineage 生产者,这些生产者可以将事件发送到此端点, 从而自动捕获来自这些服务的沿袭。不支持的功能: Data Lineage API 不支持以下情况:
- 任何后续 OpenLineage 版本(消息格式发生更改)
DatasetEventJobEvent
消息大小: 单条消息的大小上限为 5 MB。
名称长度: 输入和输出中每个完全限定名称 的长度上限为 4,000 个字符。
链接限制: 链接 按事件分组,每个事件最多包含 100 个链接。表级链接的总数上限为 1,000。如果消息包含的列级链接超过 1, 500 个,则系统会跳过列级信息。
图表范围: Knowledge Catalog 会为每个作业运行显示沿袭图,其中显示了沿袭事件的输入和输出。它不支持较低级层的进程,例如 Spark 阶段。
OpenLineage 分面属性映射
如需了解 OpenLineage 映射,请参阅 OpenLineage 映射。
导入 OpenLineage 事件
如果您尚未设置 OpenLineage,请参阅 使用入门。
如需将 OpenLineage 事件导入 Knowledge Catalog,请调用 API 方法
ProcessOpenLineageRunEvent。
C#
C#
在试用此示例之前,请按照C#设置说明进行操作,具体请参阅《Data Lineage 快速入门:使用客户端库》。如需了解详情,请参阅 Data Lineage C# API 参考文档。
如需向 Data Lineage 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证。
Go
Go
在试用此示例之前,请按照《Data Lineage 快速入门:使用客户端库》中的 Go 设置说明进行操作。如需了解详情,请参阅 Data Lineage Go API 参考文档。
如需向 Data Lineage 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证。
Java
Java
在试用此示例之前,请按照《Data Lineage 快速入门:使用客户端库》中的 Java 设置说明进行操作。如需了解详情,请参阅 Data Lineage Java API 参考文档。
如需向 Data Lineage 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证。
Python
Python
在试用此示例之前,请按照《Data Lineage 快速入门:使用客户端库》中的 Python 设置说明进行操作。如需了解详情,请参阅 Data Lineage Python API 参考文档。
如需向 Data Lineage 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证。
Ruby
Ruby
在试用此示例之前,请按照Ruby设置说明进行操作,具体请参阅《Data Lineage 快速入门:使用客户端库》。如需了解详情,请参阅 Data Lineage Ruby API 参考文档。
如需向 Data Lineage 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证。
REST
如需导入 OpenLineage 事件,请使用
processOpenLineageRunEvent 方法。
在使用任何请求数据之前, 请先进行以下替换:
PROJECT_ID:您的 Google Cloud 项目 ID。LOCATION_ID: Google Cloud 位置,例如us-central1。
HTTP 方法和网址:
POST https://datalineage.googleapis.com/v1/projects/PROJECT_ID/locations/LOCATION_ID:processOpenLineageRunEvent
请求 JSON 正文:
{
"eventTime": "2023-04-04T13:21:16.098Z",
"eventType": "COMPLETE",
"inputs": [
{
"name": "somename",
"namespace": "customnamespace"
}
],
"job": {
"name": "somename",
"namespace": "customnamespace"
},
"outputs": [
{
"name": "somename",
"namespace": "customnamespace"
}
],
"producer": "someproducer",
"run": {
"runId": "somerunid"
},
"schemaURL": "https://openlineage.io/spec/1-0-5/OpenLineage.json#/$defs/RunEvent"
}
如需发送您的请求,请展开以下选项之一:
您应该收到类似以下内容的 JSON 响应:
{
"process": "projects/my-project/locations/us-central1/processes/my-process",
"run": "projects/my-project/locations/us-central1/processes/my-process/runs/my-run",
"lineageEvents": [
"projects/my-project/locations/us-central1/processes/my-process/runs/my-run/lineageEvents/my-lineage-event"
]
}
用于发送 OpenLineage 消息的工具
为了简化向 Data Lineage API 发送事件的过程,您可以使用各种工具和库:
- Data Lineage 的 Google Cloud 客户端库: Google 提供了客户端库,以便以编程方式与 Data Lineage API 进行交互。如需了解 安装说明,请参阅 客户端库。
- Google Cloud Java 生产者库: Google 提供了一个开源 Java 库,可帮助构建 OpenLineage 事件并将其发送到 Data Lineage API。如需了解详情,请参阅博文 Data Lineage 的生产者 Java 库现已开源。 该库可在 GitHub 和 Maven上找到。
- OpenLineage GCP 传输: 对于基于 Java 的 OpenLineage 生产者,可以使用专用的
GcpLineage 传输
。它通过最大限度地减少向 Data Lineage API 发送事件所需的代码,简化了与 Data Lineage API 的集成。
GcpLineageTransport可以配置为任何现有 OpenLineage 生产者(例如 Airflow、Spark 和 Flink)的事件接收器。如需了解详情 和示例,请参阅 GcpLineage。
分析来自 OpenLineage 的信息
如需分析导入的 OpenLineage 事件,请参阅 在 Knowledge Catalog 界面中查看沿袭图。
存储的 OpenLineage 分面数据
Data Lineage API 不会存储 OpenLineage 消息中的所有分面数据。Data Lineage API 会存储以下分面字段:
spark_versionopenlineage-spark-versionspark-version
- 所有
spark.logicalPlan.* environment-properties(自定义 Google Cloud 沿袭分面)- “
origin.sourcetype”和“origin.name” spark.app.idspark.app.namespark.batch.idspark.batch.uuidspark.cluster.namespark.cluster.regionspark.job.idspark.job.uuidspark.project.idspark.query.node.namespark.session.idspark.session.uuid
- “
Data Lineage API 会存储以下信息:
eventTimerun.runIdjob.namespacejob.name
后续步骤
- 详细了解 Managed Service for Apache Spark 和 Hive 数据沿袭 集成的数据沿袭。
- 在互动式实验中试用: 使用 Data Lineage 和 OpenLineage 捕获和探索数据更新