OpenLineage との統合

このドキュメントでは、OpenLineage を Knowledge Catalog(以前の Dataplex Universal Catalog)と統合して、外部システムからデータリネージをインポートして可視化する方法について説明します。Knowledge Catalog は、ProcessOpenLineageRunEvent REST API を使用して OpenLineage コンシューマーとして機能することで、 Google Cloud サービスからの組み込みリネージとともにカスタム パイプライン リネージを統合できます。

概要

OpenLineage は、データリネージ情報を収集して分析するためのオープン プラットフォームです。OpenLineage は、リネージデータに対してオープン標準を使用し、OpenLineage API を使用して実行、ジョブ、データセットを報告するデータ パイプライン コンポーネントからリネージ イベントを取得します。

Data Lineage API を使用して OpenLineage イベントをインポートし、BigQuery、Managed Service for Apache Airflow、Cloud Data Fusion、Managed Service for Apache Spark などのGoogle Cloud サービスからのリネージ情報とともに Knowledge Catalog ウェブ インターフェースに表示できます。

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 に送信できます。Managed Service for Apache SparkManaged Airflow などの一部の Google Cloud サービスには、このエンドポイントにイベントを送信できる OpenLineage プロデューサーが組み込まれており、これらのサービスからのリネージ キャプチャを自動化できます。

  • サポートされていない機能: Data Lineage API は、次の機能をサポートしていません。

    • メッセージ形式が変更された、以降の OpenLineage リリース
    • DatasetEvent
    • JobEvent
  • メッセージ サイズ: 1 つのメッセージの最大サイズは 5 MB です。

  • 名前の長さ: 入力と出力の完全修飾名の長さは 4,000 文字に制限されています。

  • リンクの上限: リンクはイベント別にグループ化され、イベントあたりのリンク数は最大 100 個です。テーブルレベルのリンクの最大数は 1,000 個です。メッセージに 1, 500 個を超える列レベルのリンクが含まれている場合、列レベルの情報はスキップされます。

  • グラフのスコープ: Knowledge Catalog には、ジョブの実行ごとにリネージグラフが表示され、リネージ イベントの入力と出力が示されます。Spark ステージなどの下位レベルのプロセスはサポートされていません。

OpenLineage ファセット属性のマッピング

OpenLineage マッピングについては、OpenLineage マッピングをご覧ください。

OpenLineage イベントをインポートする

OpenLineage をまだ設定していない場合は、スタートガイドをご覧ください。

OpenLineage イベントを Knowledge Catalog にインポートするには、API メソッド ProcessOpenLineageRunEvent を呼び出します。

C#

C#

このサンプルを試す前に、クライアント ライブラリを使用した Data Lineage のクイックスタートにある C# の設定手順を行ってください。詳細については、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

このサンプルを試す前に、クライアント ライブラリを使用した Data Lineage のクイックスタートにある Ruby の設定手順を行ってください。詳細については、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 メソッドと URL:

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 Producer Library: Google は、OpenLineage イベントを Data Lineage API に構築して送信するのに役立つオープンソースの Java ライブラリを提供しています。詳細については、ブログ投稿の Data Lineage の Producer Java ライブラリがオープンソースになりましたをご覧ください。このライブラリは、GitHubMaven で入手できます。
  • OpenLineage GCP Transport: Java ベースの OpenLineage プロデューサーには、専用の GcpLineage Transport が用意されています。Data Lineage API にイベントを送信するために必要なコードを最小限に抑えることで、Data Lineage API との統合を簡素化します。GcpLineageTransport は、Airflow、Spark、Flink などの既存の OpenLineage プロデューサーのイベントシンクとして構成できます。詳細と例については、GcpLineage をご覧ください。

OpenLineage の情報を分析する

インポートされた OpenLineage イベントを分析するには、Knowledge Catalog UI でリネージグラフを表示するをご覧ください。

保存された OpenLineage ファセット データ

Data Lineage API では、OpenLineage メッセージのすべてのファセット データが格納されるわけではありません。Data Lineage API には、次のファセット フィールドが格納されます。

  • spark_version
    • openlineage-spark-version
    • spark-version
  • すべての spark.logicalPlan.*
  • environment-properties(カスタム Google Cloud リネージ ファセット)
    • origin.sourcetypeorigin.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

次のステップ