Desenvolva um conetor personalizado para a importação de metadados

Este documento fornece um modelo de referência para criar um conetor personalizado para extrair metadados de origens de terceiros, como MySQL, SQL Server e Oracle. Pode usar este conetor para importar metadados para o Dataplex Universal Catalog através de um pipeline de conetividade gerido. É incluído um exemplo de conetor Python para o Oracle Database Express Edition (XE) como ponto de partida. Também pode desenvolver conetores com Java, Scala ou R.

Como funcionam os conetores

Um conetor extrai metadados de uma origem de dados de terceiros, transforma os metadados no formato do ImportItem Dataplex Universal Catalog e gera ficheiros de importação de metadados que podem ser importados pelo Dataplex Universal Catalog.

O conector faz parte de um pipeline de conetividade gerido. Um pipeline de conetividade gerido é um fluxo de trabalho organizado que usa para importar metadados do catálogo universal do Dataplex. O pipeline de conetividade gerido executa o conetor e realiza outras tarefas no fluxo de trabalho de importação, como executar uma tarefa de importação de metadados e capturar registos.

O pipeline de conetividade gerido executa o conetor através de uma tarefa em lote do Google Cloud Serverless para Apache Spark. O Serverless para Apache Spark oferece um ambiente de execução do Spark sem servidor. Embora possa criar um conetor que não use o Spark, recomendamos que o use porque pode melhorar o desempenho do conetor.

Requisitos do conetor

O conector tem os seguintes requisitos:

  • O conector tem de ser uma imagem do Artifact Registry que possa ser executada no Serverless para Apache Spark.
  • O conetor tem de gerar ficheiros de metadados num formato que possa ser importado por uma tarefa de importação de metadados do catálogo universal do Dataplex (o método da API metadataJobs.create). Para ver os requisitos detalhados, consulte o artigo Ficheiro de importação de metadados.
  • O conector tem de aceitar os seguintes argumentos da linha de comandos para receber informações do pipeline:

    Argumento da linha de comandos Valor que o pipeline oferece
    target_project_id PROJECT_ID
    target_location_id REGION
    target_entry_group_id ENTRY_GROUP_ID
    output_bucket CLOUD_STORAGE_BUCKET_ID
    output_folder FOLDER_ID

    O conetor usa estes argumentos para gerar metadados num grupo de entradas de destino projects/PROJECT_ID/locations/REGION/entryGroups/ENTRY_GROUP_ID, e para escrever num contentor do Cloud Storage gs://CLOUD_STORAGE_BUCKET_ID/FOLDER_ID. Cada execução do pipeline cria uma nova pasta FOLDER_ID no contentor CLOUD_STORAGE_BUCKET_ID. O conetor deve escrever ficheiros de importação de metadados nesta pasta.

Os modelos de pipelines suportam conetores PySpark. Os modelos partem do princípio de que o controlador (mainPythonFileUri) é um ficheiro local na imagem do conector denominado main.py. Pode modificar os modelos de pipeline para outros cenários, como um conector do Spark, um URI do controlador diferente ou outras opções.

Veja como usar o PySpark para criar um item de importação no ficheiro de importação de metadados.

"""PySpark schemas for the data."""
entry_source_schema = StructType([
      StructField("display_name", StringType()),
      StructField("source", StringType())])

aspect_schema = MapType(StringType(),
                        StructType([
                            StructField("aspect_type", StringType()),
                            StructField("data", StructType([
                            ]))
                          ])
                        )

entry_schema = StructType([
  StructField("name", StringType()),
  StructField("entry_type", StringType()),
  StructField("fully_qualified_name", StringType()),
  StructField("parent_entry", StringType()),
  StructField("entry_source", entry_source_schema),
  StructField("aspects", aspect_schema)
])

import_item_schema = StructType([
  StructField("entry", entry_schema),
  StructField("aspect_keys", ArrayType(StringType())),
  StructField("update_mask", ArrayType(StringType()))
])

Antes de começar

Este guia pressupõe que tem conhecimentos de Python e PySpark.

Reveja as seguintes informações:

Faça o seguinte. Crie todos os recursos na mesma Google Cloud localização.

  1. Create or select a Google Cloud project.

    Roles required to select or create a project

    • Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
    • Create a project: To create a project, you need the Project Creator role (roles/resourcemanager.projectCreator), which contains the resourcemanager.projects.create permission. Learn how to grant roles.
    • Create a Google Cloud project:

      gcloud projects create PROJECT_ID

      Replace PROJECT_ID with a name for the Google Cloud project you are creating.

    • Select the Google Cloud project that you created:

      gcloud config set project PROJECT_ID

      Replace PROJECT_ID with your Google Cloud project name.

  2. Verify that billing is enabled for your Google Cloud project.

  3. Enable the Dataplex, Dataproc, Workflows, and Artifact Registry APIs:

    Roles required to enable APIs

    To enable APIs, you need the Service Usage Admin IAM role (roles/serviceusage.serviceUsageAdmin), which contains the serviceusage.services.enable permission. Learn how to grant roles.

    gcloud services enable dataplex.googleapis.com dataproc.googleapis.com workflows.googleapis.com artifactregistry.googleapis.com
  4. Install the Google Cloud CLI.

  5. Se estiver a usar um fornecedor de identidade (IdP) externo, tem primeiro de iniciar sessão na CLI gcloud com a sua identidade federada.

  6. Para inicializar a CLI gcloud, execute o seguinte comando:

    gcloud init
  7. Grant roles to your user account. Run the following command once for each of the following IAM roles: roles/resourcemanager.projectCreator, roles/billing.projectManager, roles/serviceusage.admin, roles/iam.serviceAccountCreator, roles/iam.securityAdmin, roles/storage.admin, roles/artifactregistry.writer, roles/dataplex.entryGroupOwner, roles/dataplex.entryOwner, roles/dataplex.aspectTypeOwner

    gcloud projects add-iam-policy-binding PROJECT_ID --member="user:USER_IDENTIFIER" --role=ROLE

    Replace the following:

    • PROJECT_ID: Your project ID.
    • USER_IDENTIFIER: The identifier for your user account. For example, myemail@example.com.
    • ROLE: The IAM role that you grant to your user account.
  8. Set up authentication:

    1. Ensure that you have the Create Service Accounts IAM role (roles/iam.serviceAccountCreator) and the Project IAM Admin role (roles/resourcemanager.projectIamAdmin). Learn how to grant roles.
    2. Create the service account:

      gcloud iam service-accounts create SERVICE_ACCOUNT_NAME

      Replace SERVICE_ACCOUNT_NAME with a name for the service account.

    3. Grant the roles/owner IAM role to the service account:

      gcloud projects add-iam-policy-binding PROJECT_ID --member="serviceAccount:SERVICE_ACCOUNT_NAME@PROJECT_ID.iam.gserviceaccount.com" --role=roles/owner

      Replace the following:

      • SERVICE_ACCOUNT_NAME: the name of the service account
      • PROJECT_ID: the project ID where you created the service account
  9. Crie um contentor do Cloud Storage para armazenar os ficheiros de importação de metadados.

  10. Crie os seguintes recursos de metadados no mesmo projeto.

    Para ver valores de exemplo, consulte a secção Exemplos de recursos de metadados para uma origem Oracle deste documento.

    1. Crie um grupo de entradas.
    2. Crie tipos de aspetos personalizados para as entradas que quer importar. Use a convenção de nomenclatura SOURCE-ENTITY_TO_IMPORT.

      Por exemplo, para uma base de dados Oracle, crie um tipo de aspeto denominado oracle-database.

      Opcionalmente, pode criar outros tipos de aspetos para armazenar outras informações.

    3. Crie tipos de entradas personalizados para os recursos que quer importar e atribua-lhes os tipos de aspetos relevantes. Use a convenção de nomenclatura SOURCE-ENTITY_TO_IMPORT.

      Por exemplo, para uma base de dados Oracle, crie um tipo de entrada denominado oracle-database. Associe-o ao tipo de aspeto denominado oracle-database.

  11. Certifique-se de que a sua origem de terceiros é acessível a partir do seu Google Cloud projeto. Para mais informações, consulte o artigo Configuração de rede sem servidor para o Apache Spark.
  12. Crie um conetor Python básico

    O conector Python básico de exemplo cria entradas de nível superior para uma origem de dados Oracle através das classes da biblioteca cliente do catálogo universal do Dataplex. Em seguida, indica os valores para os campos de entrada.

    O conetor cria um ficheiro de importação de metadados com as seguintes entradas:

    • Uma entrada instance, com o tipo de entrada projects/PROJECT_ID/locations/LOCATION/entryTypes/oracle-instance. Esta entrada representa um sistema Oracle Database XE.
    • Uma entrada database, que representa uma base de dados no sistema Oracle Database XE.

    Para criar um conector Python básico, faça o seguinte:

    1. Clone o repositório cloud-dataplex.

    2. Configure um ambiente local. Recomendamos que use um ambiente virtual.

      mkdir venv
      python -m venv venv/
      source venv/bin/activate
      

      Use as versões ativas ou de manutenção do Python. As versões 3.7 e posteriores do Python são suportadas.

    3. Crie um projeto Python.

    4. Requisitos de instalação:

      pip install -r requirements.txt
      

      Os seguintes requisitos estão instalados:

      google-cloud-dataplex==2.2.2
      google-cloud-storage
      google-cloud-secret-manager
      
    5. Adicione um ficheiro de pipeline main.py na raiz do projeto.

      from src import bootstrap
      
      
      if __name__ == '__main__':
          bootstrap.run()
      

      Quando implementa o seu código no Serverless para Apache Spark, o ficheiro main.py serve como ponto de entrada para a execução. Recomendamos que minimize a quantidade de informações armazenadas no ficheiro main.py. Use este ficheiro para chamar funções e classes definidas no conector, como a classe src/bootstap.py.

    6. Crie uma pasta src para armazenar a maioria da lógica do conector.

    7. Atualize o ficheiro src/cmd_reader.py com uma classe Python para aceitar argumentos da linha de comando. Pode usar o módulo argeparse para o fazer.

      """Command line reader."""
      import argparse
      
      
      def read_args():
          """Reads arguments from the command line."""
          parser = argparse.ArgumentParser()
      
          # Dataplex arguments
          parser.add_argument("--target_project_id", type=str, required=True,
              help="The name of the target Google Cloud project to import the metadata into.")
          parser.add_argument("--target_location_id", type=str, required=True,
              help="The target Google Cloud location where the metadata will be imported into.")
          parser.add_argument("--target_entry_group_id", type=str, required=True,
              help="The ID of the entry group to import metadata into. "
                   "The metadata will be imported into entry group with the following"
                   "full resource name: projects/${target_project_id}/"
                   "locations/${target_location_id}/entryGroups/${target_entry_group_id}.")
      
          # Oracle arguments
          parser.add_argument("--host_port", type=str, required=True,
              help="Oracle host and port number separated by the colon (:).")
          parser.add_argument("--user", type=str, required=True, help="Oracle User.")
          parser.add_argument("--password-secret", type=str, required=True,
              help="Secret resource name in the Secret Manager for the Oracle password.")
          parser.add_argument("--database", type=str, required=True,
              help="Source Oracle database.")
      
          # Google Cloud Storage arguments
          # It is assumed that the bucket is in the same region as the entry group
          parser.add_argument("--output_bucket", type=str, required=True,
              help="The Cloud Storage bucket to write the generated metadata import file.")
          parser.add_argument("--output_folder", type=str, required=True,
              help="A folder in the Cloud Storage bucket, to write the generated metadata import files.")
      
          return vars(parser.parse_known_args()[0])
      

      Em ambientes de produção, recomendamos que armazene a palavra-passe no Secret Manager.

    8. Atualize o ficheiro src/constants.py com código para criar constantes.

      """Constants that are used in the different files."""
      import enum
      
      SOURCE_TYPE = "oracle"
      
      # Symbols for replacement
      FORBIDDEN = "#"
      ALLOWED = "!"
      
      
      class EntryType(enum.Enum):
          """Types of Oracle entries."""
          INSTANCE: str = "projects/{project}/locations/{location}/entryTypes/oracle-instance"
          DATABASE: str = "projects/{project}/locations/{location}/entryTypes/oracle-database"
          DB_SCHEMA: str = "projects/{project}/locations/{location}/entryTypes/oracle-schema"
          TABLE: str = "projects/{project}/locations/{location}/entryTypes/oracle-table"
          VIEW: str = "projects/{project}/locations/{location}/entryTypes/oracle-view"
      
    9. Atualize o ficheiro src/name_builder.py com métodos para criar os recursos de metadados que quer que o conector crie para os seus recursos Oracle. Use as convenções descritas na secção Exemplos de recursos de metadados para uma origem Oracle deste documento.

      """Builds Dataplex hierarchy identifiers."""
      from typing import Dict
      from src.constants import EntryType, SOURCE_TYPE
      
      
      # Oracle cluster users start with C## prefix, but Dataplex doesn't accept #.
      # In that case in names it is changed to C!!, and escaped with backticks in FQNs
      FORBIDDEN_SYMBOL = "#"
      ALLOWED_SYMBOL = "!"
      
      
      def create_fqn(config: Dict[str, str], entry_type: EntryType,
                     schema_name: str = "", table_name: str = ""):
          """Creates a fully qualified name or Dataplex v1 hierarchy name."""
          if FORBIDDEN_SYMBOL in schema_name:
              schema_name = f"`{schema_name}`"
      
          if entry_type == EntryType.INSTANCE:
              # Requires backticks to escape column
              return f"{SOURCE_TYPE}:`{config['host_port']}`"
          if entry_type == EntryType.DATABASE:
              instance = create_fqn(config, EntryType.INSTANCE)
              return f"{instance}.{config['database']}"
          if entry_type == EntryType.DB_SCHEMA:
              database = create_fqn(config, EntryType.DATABASE)
              return f"{database}.{schema_name}"
          if entry_type in [EntryType.TABLE, EntryType.VIEW]:
              database = create_fqn(config, EntryType.DATABASE)
              return f"{database}.{schema_name}.{table_name}"
          return ""
      
      
      def create_name(config: Dict[str, str], entry_type: EntryType,
                      schema_name: str = "", table_name: str = ""):
          """Creates a Dataplex v2 hierarchy name."""
          if FORBIDDEN_SYMBOL in schema_name:
              schema_name = schema_name.replace(FORBIDDEN_SYMBOL, ALLOWED_SYMBOL)
          if entry_type == EntryType.INSTANCE:
              name_prefix = (
                  f"projects/{config['target_project_id']}/"
                  f"locations/{config['target_location_id']}/"
                  f"entryGroups/{config['target_entry_group_id']}/"
                  f"entries/"
              )
              return name_prefix + config["host_port"].replace(":", "@")
          if entry_type == EntryType.DATABASE:
              instance = create_name(config, EntryType.INSTANCE)
              return f"{instance}/databases/{config['database']}"
          if entry_type == EntryType.DB_SCHEMA:
              database = create_name(config, EntryType.DATABASE)
              return f"{database}/database_schemas/{schema_name}"
          if entry_type == EntryType.TABLE:
              db_schema = create_name(config, EntryType.DB_SCHEMA, schema_name)
              return f"{db_schema}/tables/{table_name}"
          if entry_type == EntryType.VIEW:
              db_schema = create_name(config, EntryType.DB_SCHEMA, schema_name)
              return f"{db_schema}/views/{table_name}"
          return ""
      
      
      def create_parent_name(config: Dict[str, str], entry_type: EntryType,
                             parent_name: str = ""):
          """Generates a Dataplex v2 name of the parent."""
          if entry_type == EntryType.DATABASE:
              return create_name(config, EntryType.INSTANCE)
          if entry_type == EntryType.DB_SCHEMA:
              return create_name(config, EntryType.DATABASE)
          if entry_type == EntryType.TABLE:
              return create_name(config, EntryType.DB_SCHEMA, parent_name)
          return ""
      
      
      def create_entry_aspect_name(config: Dict[str, str], entry_type: EntryType):
          """Generates an entry aspect name."""
          last_segment = entry_type.value.split("/")[-1]
          return f"{config['target_project_id']}.{config['target_location_id']}.{last_segment}"
      

      Uma vez que o ficheiro name_builder.py é usado para o código principal do Python e o código principal do PySpark, recomendamos que escreva os métodos como funções puras, em vez de como membros de uma classe.

    10. Atualize o ficheiro src/top_entry_builder.py com código para preencher as entradas de nível superior com dados.

      """Non-Spark approach for building the entries."""
      import dataclasses
      import json
      from typing import List, Dict
      
      import proto
      from google.cloud import dataplex_v1
      
      from src.constants import EntryType
      from src import name_builder as nb
      
      
      @dataclasses.dataclass(slots=True)
      class ImportItem:
          """A template class for Import API."""
      
          entry: dataplex_v1.Entry = dataclasses.field(default_factory=dataplex_v1.Entry)
          aspect_keys: List[str] = dataclasses.field(default_factory=list)
          update_mask: List[str] = dataclasses.field(default_factory=list)
      
      
      def _dict_factory(data: object):
          """Factory function required for converting Entry dataclass to dict."""
      
          def convert(obj: object):
              if isinstance(obj, proto.Message):
                  return proto.Message.to_dict(obj)
              return obj
      
          return dict((k, convert(v)) for k, v in data)
      
      
      def _create_entry(config: Dict[str, str], entry_type: EntryType):
          """Creates an entry based on a Dataplex library."""
          entry = dataplex_v1.Entry()
          entry.name = nb.create_name(config, entry_type)
          entry.entry_type = entry_type.value.format(
              project=config["target_project_id"], location=config["target_location_id"]
          )
          entry.fully_qualified_name =