Integrar com o OpenLineage

Este documento explica como integrar o OpenLineage ao Knowledge Catalog (antigo Dataplex Universal Catalog) para importar e visualizar a linhagem de dados de sistemas externos. Ao atuar como um consumidor do OpenLineage usando a API REST ProcessOpenLineageRunEvent, o Knowledge Catalog permite unificar a linhagem de pipeline personalizada com a linhagem integrada dos Google Cloud serviços.

Visão geral

OpenLineage é uma plataforma aberta para coletar e analisar informações de linhagem de dados. Usando um padrão aberto para dados de linhagem, o OpenLineage captura eventos de linhagem de componentes de pipeline de dados que usam uma API OpenLineage para gerar relatórios sobre execuções, jobs e conjuntos de dados.

Com a API Data Lineage, é possível importar eventos do OpenLineage para mostrar na interface da Web do Knowledge Catalog junto com informações de linhagem de Google Cloud serviços, como BigQuery, Serviço Gerenciado para Apache Airflow, Cloud Data Fusion e Serviço Gerenciado para Apache Spark.

Para importar eventos do OpenLineage que usam a especificação do OpenLineage, use o método da API REST ProcessOpenLineageRunEvent e mapeie as facetas do OpenLineage para os atributos da API Data Lineage.

Limitações da integração do OpenLineage

  • Versões compatíveis:a API Data Lineage é compatível com a versão principal 1 do OpenLineage.

  • Ações da API: o endpoint de API Data Lineage ProcessOpenLineageRunEvent atua apenas como um consumidor de mensagens do OpenLineage, não como um produtor. A API permite enviar informações de linhagem geradas por qualquer ferramenta ou sistema compatível com o OpenLineage para o Knowledge Catalog. Alguns Google Cloud serviços, como o Serviço Gerenciado para Apache Spark e o Airflow Gerenciado, incluem produtores do OpenLineage integrados que podem enviar eventos para esse endpoint, automatizando a captura de linhagem desses serviços.

  • Recursos não compatíveis:a API Data Lineage não oferece suporte a:

    • Qualquer versão subsequente do OpenLineage com mudanças no formato da mensagem
    • DatasetEvent
    • JobEvent
  • Tamanho da mensagem:o tamanho máximo de uma única mensagem é de 5 MB.

  • Comprimento do nome: O comprimento de cada Nome Totalmente Qualificado nas entradas e saídas é limitado a 4.000 caracteres.

  • Limites de links: Os links são agrupados por eventos, com um máximo de 100 links por evento. O número máximo agregado de links no nível da tabela é de 1.000. Se uma mensagem contiver mais de 1.500 links no nível da coluna, as informações no nível da coluna serão ignoradas.

  • Escopo do gráfico:o Knowledge Catalog mostra um gráfico de linhagem para cada execução de job, mostrando as entradas e saídas de eventos de linhagem. Ele não oferece suporte a processos de nível inferior, como estágios do Spark.

Mapeamento de atributos de faceta do OpenLineage

Para informações sobre o mapeamento do OpenLineage, consulte Mapeamento do OpenLineage.

Importar um evento do OpenLineage

Se você ainda não configurou o OpenLineage, consulte Introdução.

Para importar um evento do OpenLineage para o Knowledge Catalog, chame o método da API ProcessOpenLineageRunEvent.

C#

C#

Antes de testar essa amostra, siga as instruções de configuração para C# no Guia de início rápido do Data Lineage: como usar as bibliotecas de cliente. Para mais informações, consulte a documentação de referência da API Data Lineage C#.

Para autenticar no Data Lineage, configure o Application Default Credentials. Se quiser mais informações, consulte Configurar a autenticação para um ambiente de desenvolvimento local.

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

Antes de testar essa amostra, siga as instruções de configuração Go no Guia de início rápido do Data Lineage: como usar as bibliotecas de cliente. Para mais informações, consulte a documentação de referência da API Data Lineage Go API reference documentation.

Para autenticar no Data Lineage, configure o Application Default Credentials. Se quiser mais informações, consulte Configurar a autenticação para um ambiente de desenvolvimento local.


//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

Antes de testar essa amostra, siga as instruções de configuração Java no Guia de início rápido do Data Lineage: como usar as bibliotecas de cliente. Para mais informações, consulte a Data Lineage Java API documentação de referência.

Para autenticar no Data Lineage, configure o Application Default Credentials. Se quiser mais informações, consulte Configurar a autenticação para um ambiente de desenvolvimento local.

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

Antes de testar essa amostra, siga as instruções de configuração Python no Guia de início rápido do Data Lineage: como usar as bibliotecas de cliente. Para mais informações, consulte a Data Lineage Python API documentação de referência.

Para autenticar no Data Lineage, configure o Application Default Credentials. Se quiser mais informações, consulte Configurar a autenticação para um ambiente de desenvolvimento local.

# 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

Antes de testar essa amostra, siga as instruções de configuração para Ruby no Guia de início rápido do Data Lineage: como usar as bibliotecas de cliente. Para mais informações, consulte a documentação de referência da APIRuby Data Lineage.

Para autenticar no Data Lineage, configure o Application Default Credentials. Se quiser mais informações, consulte Configurar a autenticação para um ambiente de desenvolvimento local.

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

Para importar um evento do OpenLineage, use o processOpenLineageRunEvent método.

Antes de usar qualquer um dos dados da solicitação, faça as seguintes substituições:

  • PROJECT_ID: o ID do Google Cloud projeto.
  • LOCATION_ID: a Google Cloud localização, como us-central1.

Método HTTP e URL:

POST https://datalineage.googleapis.com/v1/projects/PROJECT_ID/locations/LOCATION_ID:processOpenLineageRunEvent

Corpo JSON da solicitação:

{
  "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"
}

Para enviar a solicitação, expanda uma destas opções:

Você receberá uma resposta JSON semelhante a esta:

{
  "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"
  ]
}

Ferramentas para enviar mensagens do OpenLineage

Para simplificar o envio de eventos à API Data Lineage, é possível usar várias ferramentas e bibliotecas:

  • Bibliotecas de cliente do Google Cloud para Data Lineage:o Google oferece bibliotecas de cliente para interagir com a API Data Lineage de maneira programática. Para instruções de instalação, consulte Bibliotecas de cliente.
  • Biblioteca de produtor Java do Google Cloud:o Google oferece uma biblioteca Java de código aberto para ajudar a criar e enviar eventos do OpenLineage para a API Data Lineage. Para mais informações, consulte a postagem do blog A biblioteca Java do produtor para Data Lineage agora é de código aberto. A biblioteca está disponível no GitHub e no Maven.
  • Transporte do OpenLineage GCP:para produtores do OpenLineage baseados em Java, um transporte GcpLineage dedicado está disponível. Ele simplifica a integração com a API Data Lineage, minimizando o código necessário para enviar eventos à API Data Lineage. O GcpLineageTransport pode ser configurado como o coletor de eventos para qualquer produtor do OpenLineage, como Airflow, Spark e Flink. Para mais informações e exemplos, consulte GcpLineage.

Analisar informações do OpenLineage

Para analisar os eventos importados do OpenLineage, consulte Visualizar gráficos de linhagem na interface do Knowledge Catalog.

Dados de faceta do OpenLineage armazenados

A API Data Lineage não armazena todos os dados de facetas das mensagens do OpenLineage. A API Data Lineage armazena os seguintes campos de faceta:

  • spark_version
    • openlineage-spark-version
    • spark-version
  • todos spark.logicalPlan.*
  • environment-properties (custom Google Cloud lineage faceta)
    • origin.sourcetype e 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

A API Data Lineage armazena as seguintes informações:

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

A seguir