Mengembangkan konektor kustom untuk impor metadata

Dokumen ini menyediakan template referensi untuk membuat konektor kustom guna mengekstrak metadata dari sumber pihak ketiga, seperti MySQL, SQL Server, dan Oracle. Anda dapat menggunakan konektor ini untuk mengimpor metadata ke Knowledge Catalog (sebelumnya Dataplex Universal Catalog) melalui pipeline konektivitas terkelola. Contoh konektor Python untuk Oracle Database Express Edition (XE) disertakan sebagai titik awal. Anda juga dapat mengembangkan konektor menggunakan Java, Scala, atau R.

Cara kerja konektor

Konektor mengekstrak metadata dari sumber data pihak ketiga, mengubah metadata ke format ImportItem Knowledge Catalog, dan membuat file impor metadata yang dapat diimpor oleh Knowledge Catalog.

Konektor adalah bagian dari pipeline konektivitas terkelola. Pipeline konektivitas terkelola adalah alur kerja yang diorkestrasi yang Anda gunakan untuk mengimpor metadata Knowledge Catalog. Pipeline konektivitas terkelola menjalankan konektor dan melakukan tugas lain dalam alur kerja impor, seperti menjalankan tugas impor metadata dan mencatat log.

Pipeline konektivitas terkelola menjalankan konektor menggunakan tugas batch Managed Service untuk Apache Spark. Managed Service untuk Apache Spark menyediakan lingkungan eksekusi Spark serverless. Meskipun Anda dapat membuat konektor yang tidak menggunakan Spark, sebaiknya gunakan Spark karena dapat meningkatkan performa konektor Anda.

Persyaratan konektor

Konektor memiliki persyaratan berikut:

  • Konektor harus berupa image Artifact Registry yang dapat dijalankan di Managed Service untuk Apache Spark.
  • Konektor harus membuat file metadata dalam format yang dapat diimpor oleh tugas impor metadata Knowledge Catalog (metode API metadataJobs.create). Untuk mengetahui persyaratan mendetail, lihat File impor metadata.
  • Konektor harus menerima argumen command line berikut untuk menerima informasi dari pipeline:

    Argumen command line Nilai yang diberikan oleh 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

    Konektor menggunakan argumen ini untuk membuat metadata dalam grup entri target projects/PROJECT_ID/locations/REGION/entryGroups/ENTRY_GROUP_ID, dan untuk menulis ke bucket Cloud Storage gs://CLOUD_STORAGE_BUCKET_ID/FOLDER_ID. Setiap eksekusi pipeline akan membuat folder baru FOLDER_ID di bucket CLOUD_STORAGE_BUCKET_ID. Konektor harus menulis file impor metadata ke folder ini.

Template pipeline mendukung konektor PySpark. Template mengasumsikan bahwa driver (mainPythonFileUri) adalah file lokal pada image konektor bernama main.py. Anda dapat mengubah template pipeline untuk skenario lain, seperti konektor Spark, URI driver yang berbeda, atau opsi lainnya.

Berikut cara menggunakan PySpark untuk membuat item impor dalam file impor metadata.

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

Sebelum memulai

Panduan ini mengasumsikan bahwa Anda sudah terbiasa dengan Python dan PySpark.

Tinjau informasi berikut:

Lakukan hal berikut. Buat semua resource di lokasi Google Cloud yang sama.

  1. Buat atau pilih Google Cloud project.

    Peran yang diperlukan untuk memilih atau membuat project

    • Pilih project: Memilih project tidak memerlukan peran IAM tertentu—Anda dapat memilih project mana pun yang telah diberi peran.
    • Membuat project: Untuk membuat project, Anda memerlukan peran Project Creator (roles/resourcemanager.projectCreator), yang berisi izin resourcemanager.projects.create. Pelajari cara memberikan peran.
    • Buat Google Cloud project:

      gcloud projects create PROJECT_ID

      Ganti PROJECT_ID dengan nama untuk Google Cloud project yang Anda buat.

    • Pilih project Google Cloud yang Anda buat:

      gcloud config set project PROJECT_ID

      Ganti PROJECT_ID dengan nama project Google Cloud Anda.

  2. Verifikasi bahwa penagihan diaktifkan untuk project Google Cloud Anda.

  3. Aktifkan Dataplex, Dataproc, Workflows, dan Artifact Registry API:

    Peran yang diperlukan untuk mengaktifkan API

    Untuk mengaktifkan API, Anda memerlukan izin serviceusage.services.enable. Jika Anda membuat project, kemungkinan Anda sudah memiliki izin ini melalui peran Pemilik (roles/owner). Jika tidak, Anda bisa mendapatkan izin ini melalui peran Admin Penggunaan Layanan (roles/serviceusage.serviceUsageAdmin). Pelajari cara memberikan peran.

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

  5. Jika Anda menggunakan penyedia identitas (IdP) eksternal, Anda harus login ke gcloud CLI dengan identitas gabungan Anda terlebih dahulu.

  6. Untuk melakukan inisialisasi gcloud CLI, jalankan perintah berikut:

    gcloud init
  7. Memberikan peran ke akun pengguna Anda. Jalankan perintah berikut satu kali untuk setiap peran IAM berikut: 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

    Ganti kode berikut:

    • PROJECT_ID: Project ID Anda.
    • USER_IDENTIFIER: ID untuk akun pengguna Anda. Misalnya, myemail@example.com.
    • ROLE: Peran IAM yang Anda berikan ke akun pengguna Anda.
  8. Menyiapkan autentikasi:

    1. Pastikan Anda memiliki peran IAM Create Service Accounts (roles/iam.serviceAccountCreator) dan peran IAM Project IAM Admin (roles/resourcemanager.projectIamAdmin). Pelajari cara memberikan peran.
    2. Buat akun layanan:

      gcloud iam service-accounts create SERVICE_ACCOUNT_NAME

      Ganti SERVICE_ACCOUNT_NAME dengan nama untuk akun layanan.

    3. Memberikan peran IAM roles/owner ke akun layanan.

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

      Ganti kode berikut:

      • SERVICE_ACCOUNT_NAME: nama dari akun layanan.
      • PROJECT_ID: project ID dimana Anda membuat akun layanan
  9. Buat bucket Cloud Storage untuk menyimpan file impor metadata.

  10. Buat resource metadata berikut dalam project yang sama.

    Untuk nilai contoh, lihat bagian Contoh resource metadata untuk sumber Oracle dalam dokumen ini.

    1. Buat grup entri.
    2. Buat jenis aspek kustom untuk entri yang ingin Anda impor. Gunakan konvensi penamaan SOURCE-ENTITY_TO_IMPORT.

      Misalnya, untuk database Oracle, buat jenis aspek bernama oracle-database.

      Secara opsional, Anda dapat membuat jenis aspek tambahan untuk menyimpan informasi lainnya.

    3. Buat jenis entri kustom untuk resource yang ingin Anda impor, dan tetapkan jenis aspek yang relevan kepadanya. Gunakan konvensi penamaan SOURCE-ENTITY_TO_IMPORT.

      Misalnya, untuk database Oracle, buat jenis entri bernama oracle-database. Tautkan ke jenis aspek yang diberi nama oracle-database.

  11. Pastikan sumber pihak ketiga Anda dapat diakses dari project Google Cloud Anda. Untuk mengetahui informasi selengkapnya, lihat Konfigurasi jaringan Managed Service untuk Apache Spark.

Membuat konektor Python dasar

Contoh konektor Python dasar membuat entri tingkat teratas untuk sumber data Oracle menggunakan class library klien Knowledge Catalog. Kemudian, Anda memberikan nilai untuk kolom entri.

Konektor membuat file impor metadata dengan entri berikut:

  • Entri instance, dengan jenis entri projects/PROJECT_ID/locations/LOCATION/entryTypes/oracle-instance. Entri ini mewakili sistem Oracle Database XE.
  • Entri database, yang merepresentasikan database di dalam sistem Oracle Database XE.

Untuk membuat konektor Python dasar, lakukan hal berikut:

  1. Buat clone repositori cloud-dataplex.

  2. Siapkan lingkungan lokal. Sebaiknya gunakan lingkungan virtual.

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

    Gunakan versi Python aktif atau pemeliharaan. Python versi 3.7 dan yang lebih baru didukung.

  3. Buat project Python.

  4. Persyaratan penginstalan:

    pip install -r requirements.txt
    

    Persyaratan berikut diinstal:

    google-cloud-dataplex==2.2.2
    google-cloud-storage
    google-cloud-secret-manager
    
  5. Tambahkan file pipeline main.py di root project.

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

    Saat men-deploy kode ke Managed Service untuk Apache Spark, file main.py berfungsi sebagai titik entri untuk eksekusi. Sebaiknya minimalkan jumlah informasi yang disimpan dalam file main.py; gunakan file ini untuk memanggil fungsi dan class yang ditentukan dalam konektor Anda, seperti class src/bootstrap.py.

  6. Buat folder src untuk menyimpan sebagian besar logika konektor Anda.

  7. Perbarui file src/cmd_reader.py dengan class Python untuk menerima argumen command line. Anda dapat menggunakan modul argeparse untuk melakukannya.

    """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])
    

    Di lingkungan produksi, sebaiknya Anda menyimpan sandi di Secret Manager.

  8. Perbarui file src/constants.py dengan kode untuk membuat konstanta.

    """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. Perbarui file src/name_builder.py dengan metode untuk membuat resource metadata yang Anda inginkan agar konektor dibuat untuk resource Oracle Anda. Gunakan konvensi yang dijelaskan di bagian Contoh resource metadata untuk sumber Oracle dalam dokumen ini.

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