Développer un connecteur personnalisé pour l'importation de métadonnées

Ce document fournit un modèle de référence pour vous aider à créer un connecteur personnalisé permettant d'extraire les métadonnées de sources tierces, telles que MySQL, SQL Server et Oracle. Vous pouvez utiliser ce connecteur pour importer des métadonnées dans Knowledge Catalog (anciennement Dataplex Universal Catalog) à l'aide d'un pipeline de connectivité gérée. Un exemple de connecteur Python pour Oracle Database Express Edition (XE) est inclus comme point de départ. Vous pouvez également développer des connecteurs à l'aide de Java, Scala ou R.

Fonctionnement des connecteurs

Un connecteur extrait les métadonnées d'une source de données tierce, les transforme au format ImportItem de Knowledge Catalog et génère des fichiers d'importation de métadonnées qui peuvent être importés par Knowledge Catalog.

Le connecteur fait partie d'un pipeline de connectivité gérée. Un pipeline de connectivité gérée est un workflow orchestré que vous utilisez pour importer des métadonnées Knowledge Catalog. Il exécute le connecteur et effectue d'autres tâches dans le workflow d'importation, comme exécuter un job d'importation de métadonnées et capturer des journaux.

Le pipeline de connectivité gérée exécute le connecteur à l'aide d'un job par lot Managed Service pour Apache Spark. Managed Service pour Apache Spark fournit un environnement d'exécution Spark sans serveur. Bien que vous puissiez créer un connecteur qui n'utilise pas Spark, nous vous recommandons d'utiliser Spark, car il peut améliorer les performances de votre connecteur.

Exigences concernant le connecteur

Le connecteur doit répondre aux exigences suivantes :

  • Le connecteur doit être une image Artifact Registry pouvant être exécutée sur Managed Service pour Apache Spark.
  • Le connecteur doit générer des fichiers de métadonnées dans un format pouvant être importé par un job d'importation de métadonnées Knowledge Catalog (méthode API metadataJobs.create). Pour en savoir plus sur les exigences, consultez la section Fichier d'importation de métadonnées.
  • Le connecteur doit accepter les arguments de ligne de commande suivants pour recevoir des informations du pipeline :

    Argument de ligne de commande Valeur fournie par le pipeline
    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

    Le connecteur utilise ces arguments pour générer des métadonnées dans un groupe d'entrées cible projects/PROJECT_ID/locations/REGION/entryGroups/ENTRY_GROUP_ID et pour écrire dans un bucket Cloud Storage gs://CLOUD_STORAGE_BUCKET_ID/FOLDER_ID. Chaque exécution du pipeline crée un dossier FOLDER_ID dans le bucket CLOUD_STORAGE_BUCKET_ID. Le connecteur doit écrire les fichiers d'importation de métadonnées dans ce dossier.

Les modèles de pipeline sont compatibles avec les connecteurs PySpark. Les modèles supposent que le pilote (mainPythonFileUri) est un fichier local sur l'image du connecteur nommée main.py. Vous pouvez modifier les modèles de pipeline pour d'autres scénarios, tels qu'un connecteur Spark, un autre URI de pilote ou d'autres options.

Voici comment utiliser PySpark pour créer un élément d'importation dans le fichier d'importation de métadonnées.

"""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(