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 serviços do Google Cloud .

Visão geral

O 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 serviços doGoogle Cloud , 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. Com a API, é possível enviar informações de linhagem geradas por qualquer ferramenta ou sistema compatível com OpenLineage para o Knowledge Catalog. Alguns Google Cloud serviços do Google Cloud, como o Serviço Gerenciado para Apache Spark e o Airflow Gerenciado, incluem produtores OpenLineage integrados que podem enviar eventos para esse endpoint, automatizando a captura de linhagem desses serviços.

  • Recursos não disponíveis:a API Data Lineage não é compatível com o seguinte:

    • 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 em entradas e saídas é limitado a 4.000 caracteres.

  • Linhagem no nível da coluna:o endpoint ProcessOpenLineageRunEvent materializa arestas de linhagem no nível da tabela apenas para eventos OpenLineage personalizados. Os gráficos de linhagem no nível da coluna são gerados apenas para transformações de SQL do BigQuery e jobs do Serviço Gerenciado para Apache Spark. As facetas personalizadas no nível da coluna não são renderizadas no console.

  • 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 é 1.000. Para mecanismos automatizados compatíveis que coletam linhagem no nível da coluna, se uma mensagem tiver 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 etapas 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 Como começar.

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

C#

C#

Antes de testar esta amostra, siga as instruções de configuração do C# no Guia de início rápido do Data Lineage: como usar 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. Para 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 esta amostra, siga as instruções de configuração do Go no Guia de início rápido do Data Lineage: como usar bibliotecas de cliente. Para mais informações, consulte a documentação de referência da API Data Lineage Go.

Para autenticar no Data Lineage, configure o Application Default Credentials. Para 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 esta amostra, siga as instruções de configuração do Java no Guia de início rápido do Data Lineage: como usar bibliotecas de cliente. Para mais informações, consulte a documentação de referência da API Data Lineage Java.

Para autenticar no Data Lineage, configure o Application Default Credentials. Para 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 esta amostra, siga as instruções de configuração do Python no Guia de início rápido do Data Lineage: como usar bibliotecas de cliente. Para mais informações, consulte a documentação de referência da API Data Lineage Python.

Para autenticar no Data Lineage, configure o Application Default Credentials. Para 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 esta amostra, siga as instruções de configuração do Ruby no Guia de início rápido do Data Lineage: como usar bibliotecas de cliente. Para mais informações, consulte a documentação de referência da API Data Lineage Ruby.

Para autenticar no Data Lineage, configure o Application Default Credentials. Para 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 método processOpenLineageRunEvent.

Antes de usar os dados da solicitação abaixo, faça as substituições a seguir:

  • PROJECT_ID: o ID do projeto Google Cloud .
  • LOCATION_ID: o local Google Cloud , 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, use várias ferramentas e bibliotecas:

  • Bibliotecas de cliente do Google Cloud para linhagem de dados: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 à API Data Lineage. Para mais informações, consulte a postagem do blog A biblioteca Java do produtor para linhagem de dados 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. Ela simplifica a integração com a API Data Lineage, minimizando o código necessário para enviar eventos a essa API. O GcpLineageTransport pode ser configurado como o coletor de eventos para qualquer produtor OpenLineage atual, 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 Ver 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 os spark.logicalPlan.*
  • environment-properties (aspecto de linhagem Google Cloud personalizado)
    • 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