与 OpenLineage 集成

本文档介绍了如何将 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 SparkManaged Airflow)包含 内置的 OpenLineage 生产者,这些生产者可以将事件发送到此端点, 从而自动捕获来自这些服务的沿袭。

  • 不支持的功能: Data Lineage API 不支持以下情况:

    • 任何后续 OpenLineage 版本(消息格式发生更改)
    • DatasetEvent
    • JobEvent
  • 消息大小: 单条消息的大小上限为 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 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证

using Google.Cloud.DataCatalog.Lineage.V1;
using Google.Protobuf.WellKnownTypes;

public sealed partial class GeneratedLineageClientSnippets
{
    /// <summary>Snippet for ProcessOpenLineageRunEvent</summary>
    /// <remarks>
    /// This snippet has been automatically generated and should be regarded as a code template only.
    /// It will require modifications to work:
    /// - It may require correct/in-range values for request initialization.
    /// - It may require specifying regional endpoints when creating the service client as shown in
    ///   https://cloud.google.com/dotnet/docs/reference/help/client-configuration#endpoint.
    /// </remarks>
    public void ProcessOpenLineageRunEventRequestObject()
    {
        // Create client
        LineageClient lineageClient = LineageClient.Create();
        // Initialize request argument(s)
        ProcessOpenLineageRunEventRequest request = new ProcessOpenLineageRunEventRequest
        {
            Parent = "",
            OpenLineage = new Struct(),
        };
        // Make the request
        ProcessOpenLineageRunEventResponse response = lineageClient.ProcessOpenLineageRunEvent(request);
    }
}

Go

Go

在试用此示例之前,请按照《Data Lineage 快速入门:使用客户端库》中的 Go 设置说明进行操作。如需了解详情,请参阅 Data Lineage Go API 参考文档

如需向 Data Lineage 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证


//go:build examples

package main

import (
	"context"

	lineage "cloud.google.com/go/datacatalog/lineage/apiv1"
	lineagepb "cloud.google.com/go/datacatalog/lineage/apiv1/lineagepb"
)

func main() {
	ctx := context.Background()
	// This snippet has been automatically generated and should be regarded as a code template only.
	// It will require modifications to work:
	// - It may require correct/in-range values for request initialization.
	// - It may require specifying regional endpoints when creating the service client as shown in:
	//   https://pkg.go.dev/cloud.google.com/go#hdr-Client_Options
	c, err := lineage.NewClient(ctx)
	if err != nil {
		// TODO: Handle error.
	}
	defer c.Close()

	req := &lineagepb.ProcessOpenLineageRunEventRequest{
		// TODO: Fill request struct fields.
		// See https://pkg.go.dev/cloud.google.com/go/datacatalog/lineage/apiv1/lineagepb#ProcessOpenLineageRunEventRequest.
	}
	resp, err := c.ProcessOpenLineageRunEvent(ctx, req)
	if err != nil {
		// TODO: Handle error.
	}
	// TODO: Use resp.
	_ = resp
}

Java

Java

在试用此示例之前,请按照《Data Lineage 快速入门:使用客户端库》中的 Java 设置说明进行操作。如需了解详情,请参阅 Data Lineage Java API 参考文档

如需向 Data Lineage 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证

import com.google.cloud.datacatalog.lineage.v1.LineageClient;
import com.google.cloud.datacatalog.lineage.v1.ProcessOpenLineageRunEventRequest;
import com.google.cloud.datacatalog.lineage.v1.ProcessOpenLineageRunEventResponse;
import com.google.protobuf.Struct;

public class SyncProcessOpenLineageRunEvent {

  public static void main(String[] args) throws Exception {
    syncProcessOpenLineageRunEvent();
  }

  public static void syncProcessOpenLineageRunEvent() throws Exception {
    // This snippet has been automatically generated and should be regarded as a code template only.
    // It will require modifications to work:
    // - It may require correct/in-range values for request initialization.
    // - It may require specifying regional endpoints when creating the service client as shown in
    // https://cloud.google.com/java/docs/setup#configure_endpoints_for_the_client_library
    try (LineageClient lineageClient = LineageClient.create()) {
      ProcessOpenLineageRunEventRequest request =
          ProcessOpenLineageRunEventRequest.newBuilder()
              .setParent("parent-995424086")
              .setOpenLineage(Struct.newBuilder().build())
              .setRequestId("requestId693933066")
              .build();
      ProcessOpenLineageRunEventResponse response =
          lineageClient.processOpenLineageRunEvent(request);
    }
  }
}

Python

Python

在试用此示例之前,请按照《Data Lineage 快速入门:使用客户端库》中的 Python 设置说明进行操作。如需了解详情,请参阅 Data Lineage Python API 参考文档

如需向 Data Lineage 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证

# This snippet has been automatically generated and should be regarded as a
# code template only.
# It will require modifications to work:
# - It may require correct/in-range values for request initialization.
# - It may require specifying regional endpoints when creating the service
#   client as shown in:
#   https://googleapis.dev/python/google-api-core/latest/client_options.html
from google.cloud import datacatalog_lineage_v1


def sample_process_open_lineage_run_event():
    # Create a client
    client = datacatalog_lineage_v1.LineageClient()

    # Initialize request argument(s)
    request = datacatalog_lineage_v1.ProcessOpenLineageRunEventRequest(
        parent="parent_value",
    )

    # Make the request
    response = client.process_open_lineage_run_event(request=request)

    # Handle the response
    print(response)

Ruby

Ruby

在试用此示例之前,请按照Ruby设置说明进行操作,具体请参阅《Data Lineage 快速入门:使用客户端库》。如需了解详情,请参阅 Data Lineage Ruby API 参考文档

如需向 Data Lineage 进行身份验证,请设置应用默认凭据。 如需了解详情,请参阅 为本地开发环境设置身份验证

require "google/cloud/data_catalog/lineage/v1"

##
# Snippet for the process_open_lineage_run_event call in the Lineage service
#
# This snippet has been automatically generated and should be regarded as a code
# template only. It will require modifications to work:
# - It may require correct/in-range values for request initialization.
# - It may require specifying regional endpoints when creating the service
# client as shown in https://cloud.google.com/ruby/docs/reference.
#
# This is an auto-generated example demonstrating basic usage of
# Google::Cloud::DataCatalog::Lineage::V1::Lineage::Client#process_open_lineage_run_event.
#
def process_open_lineage_run_event
  # Create a client object. The client can be reused for multiple calls.
  client = Google::Cloud::DataCatalog::Lineage::V1::Lineage::Client.new

  # Create a request. To set request fields, pass in keyword arguments.
  request = Google::Cloud::DataCatalog::Lineage::V1::ProcessOpenLineageRunEventRequest.new

  # Call the process_open_lineage_run_event method.
  result = client.process_open_lineage_run_event request

  # The returned object is of type Google::Cloud::DataCatalog::Lineage::V1::ProcessOpenLineageRunEventResponse.
  p result
end

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 库现已开源。 该库可在 GitHubMaven上找到。
  • 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_version
    • openlineage-spark-version
    • spark-version
  • 所有 spark.logicalPlan.*
  • environment-properties(自定义 Google Cloud 沿袭分面)
    • origin.sourcetype”和“origin.name
    • spark.app.id
    • spark.app.name
    • spark.batch.id
    • spark.batch.uuid
    • spark.cluster.name
    • spark.cluster.region
    • spark.job.id
    • spark.job.uuid
    • spark.project.id
    • spark.query.node.name
    • spark.session.id
    • spark.session.uuid

Data Lineage API 会存储以下信息:

  • eventTime
  • run.runId
  • job.namespace
  • job.name

后续步骤