使用 Workflows 从自定义来源导入元数据

本文档介绍了如何通过在 Workflows 中运行托管式连接流水线,将来自第三方来源的元数据导入 Dataplex Universal Catalog 中。

如需设置托管式连接流水线,您需要为数据源构建连接器。然后在 Workflows 中运行该流水线。该流水线会从您的数据源提取元数据,然后将元数据导入 Dataplex Universal Catalog 中。如有必要,该流水线还会在您的 Google Cloud 项目中创建 Dataplex Universal Catalog 条目组。

如需详细了解托管式连接,请参阅托管式连接概览

准备工作

在导入元数据之前,请完成本部分中的任务。

构建连接器

连接器会从您的数据源提取元数据,并生成可由 Dataplex Universal Catalog 导入的元数据导入文件。该连接器是一个可在Google Cloud Serverless for Apache Spark 上运行的 Artifact Registry 映像。

配置 Google Cloud 资源

  1. Enable the Workflows, Dataproc, Cloud Storage, Dataplex, Secret Manager, Artifact Registry, and Cloud Scheduler 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.

    Enable the APIs

    如果您不打算按时间表运行流水线,则无需启用 Cloud Scheduler API。

  2. 在 Secret Manager 中创建 Secret,以存储第三方数据源的凭证。

  3. 配置虚拟私有云 (VPC) 网络以运行 Serverless for Apache Spark 工作负载。

  4. 创建 Cloud Storage 存储桶以存储元数据导入文件。

  5. 创建以下 Dataplex Universal Catalog 资源:

    1. 为要导入的条目创建自定义切面类型

    2. 为要导入的条目创建自定义条目类型

所需的角色

服务账号代表工作流的身份,并决定了工作流拥有的权限以及它可以访问的 Google Cloud 资源。您需要一个 Workflows 服务账号(用于运行流水线)和一个 Serverless for Apache Spark 服务账号(用于运行连接器)。

您可以使用 Compute Engine 默认服务账号 (PROJECT_NUMBER-compute@developer.gserviceaccount.com),也可以创建自己的服务账号来运行托管式连接流水线。

控制台

  1. 在 Google Cloud 控制台中,前往 IAM 页面。

    转到 IAM

  2. 选择要导入元数据的目标项目。

  3. 点击 授予访问权限,然后输入服务账号的邮箱。

  4. 为服务账号分配以下角色:

    • Logs Writer
    • Dataplex Entry Group Owner
    • Dataplex Metadata Job Owner
    • Dataplex Catalog Editor
    • Dataproc Editor
    • Dataproc Worker
    • Secret Manager Secret Accessor -针对存储数据源凭证的 Secret
    • Storage Object User - 针对 Cloud Storage 存储桶
    • Artifact Registry Reader - 针对包含连接器映像的 Artifact Registry 存储库
    • Service Account User - 如果使用不同的服务账号,请为运行 Workflows 的账号授予此角色,以便其可调用运行 Serverless for Apache Spark 批量作业的服务账号
    • Workflows Invoker - 如果您想安排流水线
  5. 保存更改。

gcloud

  1. 向服务账号授予角色。运行以下命令:

    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:SERVICE_ACCOUNT_ID" \
        --role=roles/logging.logWriter
    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:SERVICE_ACCOUNT_ID" \
        --role=roles/dataplex.entryGroupOwner
    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:SERVICE_ACCOUNT_ID" \
        --role=roles/dataplex.metadataJobOwner
    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:SERVICE_ACCOUNT_ID" \
        --role=roles/dataplex.catalogEditor
    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:SERVICE_ACCOUNT_ID" \
        --role=roles/dataproc.editor
    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:SERVICE_ACCOUNT_ID" \
        --role=roles/dataproc.worker
    

    替换以下内容:

    • PROJECT_ID:用于导入元数据的目标 Google Cloud项目名称。
    • SERVICE_ACCOUNT_ID:服务账号,例如 my-service-account@my-project.iam.gserviceaccount.com
  2. 在资源级别为该服务账号授予以下角色:

    gcloud secrets add-iam-policy-binding SECRET_ID \
        --member="serviceAccount:SERVICE_ACCOUNT_ID" \
        --role=roles/secretmanager.secretaccessor
    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:SERVICE_ACCOUNT_ID" \
        --role=roles/storage.objectUser \
        --condition=resource.name.startsWith('projects/_/buckets/BUCKET_ID')
    gcloud artifacts repositories add-iam-policy-binding REPOSITORY \
        --location=REPOSITORY_LOCATION \
        --member=SERVICE_ACCOUNT_ID} \
        --role=roles/artifactregistry.reader
    

    替换以下内容:

    • SECRET_ID:用于存储数据源凭证的 Secret 的 ID。其格式为 projects/PROJECT_ID/secrets/SECRET_ID
    • BUCKET_ID:Cloud Storage 存储桶的名称。
    • REPOSITORY:包含连接器映像的 Artifact Registry 制品库。
    • REPOSITORY_LOCATION:存储库托管所在的 Google Cloud位置。
  3. 为运行 Workflows 的服务账号授予运行 Serverless for Apache Spark 批量作业的服务账号的 roles/iam.serviceAccountUser 角色。即使您为 Workflows 和 Serverless for Apache Spark 使用相同的服务账号,也必须授予此角色。

    gcloud iam service-accounts add-iam-policy-binding \
        serviceAccount:SERVICE_ACCOUNT_ID \
        --member='SERVICE_ACCOUNT_ID' \
        --role='roles/iam.serviceAccountUser'
    

    如果您使用的是不同的服务账号,--member 标志的值就是运行 Serverless for Apache Spark 批量作业的服务账号。

  4. 如果您想安排流水线,请为服务账号授予以下角色:

    gcloud projects add-iam-policy-binding PROJECT_ID \
     --member="SERVICE_ACCOUNT_ID" \
     --role=roles/workflows.invoker
    

导入元数据

如需导入元数据,请创建并执行一个用于运行托管式连接性流水线的工作流。(可选)您还可以为流水线创建一个运行时间表。

控制台

  1. 创建工作流。提供以下信息:

    • 服务账号:您在本文档的所需角色部分中配置的服务账号。
    • 加密:选择 Google-managed encryption key

    • 定义工作流:提供以下定义文件:

      main:
        params: [args]
        steps:
          - init:
              assign:
              - WORKFLOW_ID: ${"metadataworkflow-" + sys.get_env("GOOGLE_CLOUD_WORKFLOW_EXECUTION_ID")}
              - NETWORK_URI: ${default(map.get(args, "NETWORK_URI"), "")}
              - SUBNETWORK_URI: ${default(map.get(args, "SUBNETWORK_URI"), "")}
              - NETWORK_TAGS: ${default(map.get(args, "NETWORK_TAGS"), [])}
      
          - check_networking:
              switch:
                - condition: ${NETWORK_URI != "" and SUBNETWORK_URI != ""}
                  raise: "Error: cannot set both network_uri and subnetwork_uri. Please select one."
                - condition: ${NETWORK_URI == "" and SUBNETWORK_URI == ""}
                  steps:
                   - submit_extract_job_with_default_network_uri:
                        assign:
                          - NETWORK_TYPE: "networkUri"
                          - NETWORKING: ${"projects/" + args.TARGET_PROJECT_ID + "/global/networks/default"}  
                - condition: ${NETWORK_URI != ""}
                  steps:
                    - submit_extract_job_with_network_uri:
                        assign:
                          - NETWORKING: ${NETWORK_URI}
                          - NETWORK_TYPE: "networkUri"
                - condition: ${SUBNETWORK_URI != ""}
                  steps:
                    - submit_extract_job_with_subnetwork_uri:
                        assign:
                          - NETWORKING: ${SUBNETWORK_URI}
                          - NETWORK_TYPE: "subnetworkUri"
              next: check_create_target_entry_group
      
          - check_create_target_entry_group:
              switch:
                - condition: ${args.CREATE_TARGET_ENTRY_GROUP == true}
                  next: create_target_entry_group
                - condition: ${args.CREATE_TARGET_ENTRY_GROUP == false}
                  next: prepare_pyspark_job_body
      
          - create_target_entry_group:
              call: http.post
              args:
                url: ${"https://dataplex.googleapis.com/v1/projects/" + args.TARGET_PROJECT_ID + "/locations/" + args.CLOUD_REGION + "/entryGroups?entry_group_id=" + args.TARGET_ENTRY_GROUP_ID}
                auth:
                  type: OAuth2
                  scopes: "https://www.googleapis.com/auth/cloud-platform"
              next: prepare_pyspark_job_body
      
          -