尧图精选

DataHub GCS 数据源连接器全解析:三种认证方式、Path Specs 配置与 GCS 数据湖摄取实战

🕒 发布时间:2026/9/18 13:43:34 📁 来源:尧图网络
DataHub GCS 数据源连接器全解析三种认证方式、Path Specs 配置与 GCS 数据湖摄取实战【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub导读本文基于当前仓库 metadata-ingestion/docs/sources/gcs/gcs_pre.md 与其配套文档 gcs_post.md、示例配方 gcs_recipe.yml 展开。Google Cloud StorageGCS连接器是 DataHub 面向生产环境的元数据摄取模块它借助 GCS 与 S3 的互操作能力在底层复用 DataHub S3 Data Lake 集成源将 GCS 中的单个文件或文件夹映射为 DataHub 中的 Dataset。读完本文你将掌握 GCS 连接器的 HMAC、GKE Workload Identity、Workload Identity Federation 三种认证方式的选型与配置、Path Specs 规则的精确定义与实战写法以及 Schema 推断、数据 Profiling、支持的文件类型等能力边界并能据此编写可直接运行的摄取配方Recipe。一、连接器概述GCS 摄取是如何实现的gcs模块源码位于 metadata-ingestion/src/datahub/ingestion/source/gcs/gcs_source.py将 Google Cloud Storage 数据集摄取进 DataHub它面向生产摄取工作流设计并具备独立的模块级能力。该连接器的核心设计理念是复用而非重造它允许将单个文件或一组文件夹中的文件映射为 DataHub 中的一个 Dataset指定构成一个数据集的文件分组通过摄取配方Recipe中的path_specs配置完成该源利用了GCS 与 S3 的互操作性Interoperability of GCS with S3即 GCS 提供的兼容 S3 的 XML API 端点在底层直接使用DataHub 的 S3 Data Lake 集成源S3SourcePath Specs 的详细语义与 S3 连接器保持一致。从源码看这一代理结构非常直观GCSSource在初始化时通过create_equivalent_s3_source构建一个内部的S3Source实例并将自身所有工作单元委托给它gcs_source.pycreate_equivalent_s3_path_specs()把gs://前缀的 Path Spec 转换为s3://前缀同时完整保留file_types、table_name、autodetect_partitions、sample_files、exclude、traversal_method、emit_folders_only等全部字段create_equivalent_s3_config()根据认证类型构造 S3 兼容连接配置端点固定为https://storage.googleapis.com常量GCS_ENDPOINT_URL定义在 gcs_utils.pyregion 固定为autoget_workunits_internal()直接透传self.s3_source.get_workunits_internal()的结果。注意GCS 连接器还通过 data_lake_common/object_store.py 中的对象存储适配器create_object_store_adapter(gcs)对 S3 源施加 GCS 定制化例如桶与文件夹被映射为带GCS bucket、Folder子类型的 Container详见 gcs/README.md 的概念映射表。在能力声明上gcs_source.py连接器支持能力状态容器Container含 GCS bucket / Folder 子类型默认启用Schema 元数据默认启用数据 Profiling可选启用通过配置开启当前该源的官方支持状态为BETASupportStatus.BETA。二、前置条件运行摄取之前需要确保网络连通性DataHub 摄取环境能够访问 GCS 的 S3 互操作端点https://storage.googleapis.com有效认证凭据根据所选认证方式准备好对应凭据详见下文读取权限服务账号/外部身份需具备该模块所需元数据 API 的读取权限最典型的是Storage Object Viewerroles/storage.objectViewer角色以允许列举对象、读取对象与对象标签等 S3 互操作操作ListBuckets、ListObjectsV2、GetObject、HeadObject等见 gcs_source.py。三、三种认证方式详解GCS 连接器支持三种认证方式枚举GCSAuthType定义于 gcs_source.py通过配方中的auth_type字段选择auth_type取值认证机制适用场景hmac默认长寿命 HMAC 密钥access id secret简单部署、GCP 之外的 Service Account需要显式配置凭据workload_identity无密钥使用 Application Default CredentialsADCDataHub 运行在 GKE 且已启用 Workload Identity无需任何凭据配置workload_identity_federation无密钥、基于令牌通过 GCP STS 端点将外部身份令牌兑换为短时 GCP 凭据DataHub 运行在 GCP 之外AWS、Azure、本地机房且希望免分发服务账号密钥文件3.1 HMAC 认证默认HMAC 是默认认证方式也是唯一需要显式提供长期密钥的方式。配置步骤创建一个具备Storage Object Viewer角色的服务账号Service Account确认满足生成 HMAC 密钥的前置要求GCS 管理 HMAC 密钥的相关约束为该服务账号创建一对 HMAC 密钥Access ID 与 Secret。在配方中凭据通过credential字段提供结构HMACKey定义于 gcs_utils.pysource: type: gcs config: path_specs: - include: gs://gcs-ingestion-bucket/parquet_example/{table}/year{partition[0]}/*.parquet credential: hmac_access_id: hmac access id hmac_access_secret: hmac access secret从源码实现看gcs_source.pyHMAC 模式下连接器直接构造标准的AwsConnectionConfigaws_endpoint_url指向https://storage.googleapis.comaws_access_key_id/aws_secret_access_key填入 HMAC 密钥对aws_region固定为auto。这正是GCS 的 S3 互操作端点 HMAC 密钥即 S3 风格的访问凭据的实现体现。配置校验重要GCSSourceConfig.validate_credentialgcs_source.py保证auth_type为hmac时credential必填否则直接报错credential is required when auth_type is hmacauth_type非hmac时credential必须为空避免凭据误配。3.2 GKE Workload Identity推荐DataHub 在 GKE 内当 DataHub 运行在 GKE 集群内部且已启用 Workload Identity 时这是最简方案无需任何凭据文件或 SecretPod 的 Kubernetes Service Account 会被自动使用。配置步骤在 GKE 集群上启用 Workload Identity创建一个具备Storage Object Viewer角色的 Google Service AccountGSA将 Kubernetes Service AccountKSA绑定到该 GSAgcloud iam service-accounts add-iam-policy-binding GSA_EMAIL \ --role roles/iam.workloadIdentityUser \ --member serviceAccount:PROJECT_ID.svc.id.goog[NAMESPACE/KSA_NAME] kubectl annotate serviceaccount KSA_NAME \ --namespace NAMESPACE \ iam.gke.io/gcp-service-accountGSA_EMAIL在配方中设置auth_type: workload_identity不需要credential或任何 WIF 配置字段source: type: gcs config: auth_type: workload_identity path_specs: - include: gs://gcs-ingestion-bucket/parquet_example/{table}/year{partition[0]}/*.parquet源码层原理此模式走_setup_adc_credentialsgcs_source.py调用google.auth.default(scopes[https://www.googleapis.com/auth/cloud-platform])加载 ADC。若加载失败会抛出带明确指引的ValueError提示检查 GKE Workload Identity 是否启用、Pod Service Account 是否绑定。由于 boto3 默认使用 SigV4 签名而 GCS XML API 接受 Bearer 令牌连接器通过GCSOAuthAwsConnectionConfig使用虚拟的 AWS 风格密钥aws_access_key_idgcs-oauth、aws_secret_access_keynot-used让 boto3 建立会话再注册before-send.s3.*事件处理器在每个请求发出前把Authorization头替换为Bearer token并附加x-goog-project-id头见 gcs_source.py。若 ADC 未返回 project ID日志会提示通过GCLOUD_PROJECT或GOOGLE_CLOUD_PROJECT环境变量设置。3.3 Workload Identity Federation推荐DataHub 在 GCP 之外当 DataHub运行在 GCP 之外AWS、Azure、本地机房等且希望免密钥认证、不必分发服务账号密钥文件时使用 Workload Identity FederationWIF。配置步骤在 Google Cloud 中创建 Workload Identity Pool 与 Provider授予外部身份模拟某个具备Storage Object Viewer角色的 GCP 服务账号的权限从 Google Cloud Console 下载或生成 WIF 凭据配置文件JSON 格式通过以下三种互斥方式之一将配置交给连接器配方字段说明gcp_wif_configurationWIF 配置 JSON文件路径gcp_wif_configuration_json内联配置dict也兼容直接内联 JSON 字符串gcp_wif_configuration_json_string配置内容以JSON 字符串形式提供便于从 Secret 管理器注入三种选项的定义与互斥校验实现在 common/gcp_wif_config.py同时指定多个选项会直接报错gcp_wif_configuration_json传入字符串时会被自动json.loads解析为 dict向后兼容gcp_wif_configuration_json_string会在配置阶段校验必须是合法 JSON。该 MixinGCPWIFConfig被设计为可被 BigQuery、Dataplex、VertexAI 等其他 GCP 源复用。配方示例一配置文件路径source: type: gcs config: auth_type: workload_identity_federation gcp_wif_configuration: /path/to/gcp_wif_configuration.json path_specs: - include: gs://gcs-ingestion-bucket/parquet_example/{table}/year{partition[0]}/*.parquet配方示例二内联 dictsource: type: gcs config: auth_type: workload_identity_federation gcp_wif_configuration_json: type: external_account audience: //iam.googleapis.com/projects/PROJECT_NUMBER/locations/global/workloadIdentityPools/POOL_ID/providers/PROVIDER_ID subject_token_type: urn:ietf:params:oauth:token-type:jwt token_url: https://sts.googleapis.com/v1/token credential_source: file: /var/run/secrets/tokens/gcp-ksa/token service_account_impersonation_url: https://iamcredentials.googleapis.com/v1/projects/-/serviceAccounts/SERVICE_ACCOUNT_EMAIL:generateAccessToken path_specs: - include: gs://gcs-ingestion-bucket/parquet_example/{table}/year{partition[0]}/*.parquet配方示例三JSON 字符串从配置文件复制粘贴source: type: gcs config: auth_type: workload_identity_federation gcp_wif_configuration_json_string: | { type: external_account, audience: //iam.googleapis.com/projects/PROJECT_NUMBER/locations/global/workloadIdentityPools/POOL_ID/providers/PROVIDER_ID, subject_token_type: urn:ietf:params:oauth:token-type:jwt, token_url: https://sts.googleapis.gov/v1/token, credential_source: { file: /var/run/secrets/tokens/gcp-ksa/token }, service_account_impersonation_url: https://iamcredentials.googleapis.com/v1/projects/-/serviceAccounts/SERVICE_ACCOUNT_EMAIL:generateAccessToken } path_specs: - include: gs://gcs-ingestion-bucket/parquet_example/{table}/year{partition[0]}/*.parquet源码层原理WIF 模式走_setup_wif_credentialsgcs_source.py最终调用load_wif_credentialscommon/gcp_wif_config.py先解析出 WIF 配置 dict再交给google.auth.load_credentials_from_dict构建凭据由于 WIF 走服务账号模拟impersonation连接器会为凭据附加cloud-platformscope否则 IAM 会返回 400 Scope required令牌在首次 API 调用时惰性刷新。与workload_identity一样最终都经由GCSOAuthAwsConnectionConfig以 Bearer 令牌方式访问 S3 互操作端点。配置校验重要validate_gcp_wif_configuration_optionsgcs_source.py保证auth_type为workload_identity_federation时三个 WIF 字段至少提供一个否则报错auth_type为workload_identity时三个 WIF 字段必须全部为空凭据自动来自 ADC。四、Path Specs如何把 GCS 文件/文件夹映射为数据集Path Specs 是本连接器最核心的配置。它通过 data_lake_common/path_spec.py 中的PathSpec模型定义GCSSourceConfig要求path_specs非空且每个path_spec.include必须以gs://开头校验逻辑见 gcs_source.py。4.1 示例一数据集 单个文件桶结构test-gs-bucket ├── employees.csv └── food_items.csv配置path_specs: - include: gs://test-gs-bucket/*.csv此时每个匹配的文件各自成为一个 Dataset。4.2 示例二带分区的数据集桶结构test-gs-bucket ├── orders │ └── year2022 │ └── month2 │ ├── 1.parquet │ └── 2.parquet └── returns └── year2021 └── month2 └── 1.parquet配置path_specs: - include: gs://test-gs-bucket/{table}/{partition_key[0]}{partition[0]}/{partition_key[1]}{partition[1]}/*.parquet这里{table}对应orders/returns文件夹即 Dataset 的划分单位{partition_key[i]}对应分区名year、month{partition[i]}对应分区值2022、2。4.3 示例三分区 exclude 排除桶结构test-gs-bucket ├── orders │ └── year2022 │ └── month2 │ ├── 1.parquet │ └── 2.parquet └── tmp_orders └── year2021 └── month2 └── 1.parquet配置path_specs: - include: gs://test-gs-bucket/{table}/{partition_key[0]}{partition[0]}/{partition_key[1]}{partition[1]}/*.parquet exclude: - **/tmp_orders/**exclude使用 glob 模式支持**用于在扫描时剔除不需要的路径。4.4 示例四混合性质的多个数据集桶结构test-gs-bucket ├── customers │ ├── part1.json │ ├── part2.json │ ├── part3.json │ └── part4.json ├── employees.csv ├── food_items.csv ├── tmp_10101000.csv └── orders └── year2022 └── month2 ├── 1.parquet ├── 2.parquet └── 3.parquet配置多个path_specs条目并存path_specs: - include: gs://test-gs-bucket/*.csv exclude: - **/tmp_10101000.csv - include: gs://test-gs-bucket/{table}/*.json - include: gs://test-gs-bucket/{table}/{partition_key[0]}{partition[0]}/{partition_key[1]}{partition[1]}/*.parquet4.5 合法的path_specs.include格式汇总gs://my-bucket/foo/tests/bar.avro # 单文件表 gs://my-bucket/foo/tests/*.* # 多个文件级表 gs://my-bucket/foo/tests/{table}/*.avro # 无分区的表 gs://my-bucket/foo/tests/{table}/*/*.avro # 分区未指定的表 gs://my-bucket/foo/tests/{table}/*.* # 未指定分区与数据类型 gs://my-bucket/{dept}/tests/{table}/*.avro # 指定用于显示名的关键字 gs://my-bucket/{dept}/tests/{table}/{partition_key[0]}{partition[0]}/{partition_key[1]}{partition[1]}/*.avro # 指定分区键与值格式 gs://my-bucket/{dept}/tests/{table}/{partition[0]}/{partition[1]}/{partition[2]}/*.avro # 仅指定分区值格式 gs://my-bucket/{dept}/tests/{table}/{partition[0]}/{partition[1]}/{partition[2]}/*.* # 全部扩展名 gs://my-bucket/*/{table}/{partition[0]}/{partition[1]}/{partition[2]}/*.* # 表位于桶下 2 层 gs://my-bucket/*/*/{table}/{partition[0]}/{partition[1]}/{partition[2]}/*.* # 表位于桶下 3 层4.6 合法的path_specs.exclude格式汇总**/tests/**gs://my-bucket/hr/***_/tests/_.csv即如**/tests/*.csv之类的通配写法gs://my-bucket/foo/*/my_table/**4.7 重要规则与注意事项{table}代表将为之创建 Dataset 的文件夹include路径必须以*.*或*.[ext]结尾以表示叶层若提供*.[ext]则只扫描指定类型的文件/*/代表单层文件夹{partition[i]}代表分区值{partition_key[i]}代表分区名抽取时用索引 i 将分区键与分区值配对include中必须指定所有文件夹层级只有exclude可以使用**式匹配exclude路径中不能包含命名变量{}文件夹名不能包含{、}、*、/字符{folder}是内部工作保留字不要在命名变量中使用。此外还有两个与 Path Specs 相关的PathSpec高级参数完整定义见 path_spec.pyautodetect_partitions默认true当文件夹形如year2024的{partition_key}{partition_value}格式时自动检测分区键/值traversal_method默认MAX文件夹遍历方式可选ALL遍历全部、MIN_MAX按最小/最大值遍历、MAX只遍历最大值文件夹sample_files默认true是否只采样少量文件推断 Schema会关闭文件计数与大小计算显著影响性能allow_double_stars默认false是否允许include中出现**开启会影响性能include_hidden_folders默认false是否包含以.或_开头的隐藏文件夹tables_filter_pattern默认放行全部用正则精确过滤{table}部分的表实现细粒度包含/排除。⚠️成本警告path_specs.include中请尽量指定足够长的固定前缀不含/*/这将显著减少扫描时间与成本对 Google Cloud Storage 尤其重要。⚠️出口流量警告如果从 Google Cloud Storage 摄取数据集建议在与源同区域的服务器上运行摄取以避免高昂的出口egress费用。如果你需要更复杂的文件名解析逻辑可以引入 {transformer} 实现。五、支持的文件类型与 Schema 推断GCS 连接器支持的文件类型常量SUPPORTED_FILE_TYPES见 path_spec.pyCSVTSVJSONLJSONParquetApache AvroSchema 推断行为如下Parquet 与 AvroSchema 按文件原样提取文件自带 schemaCSV、TSV、JSONLSchema 通过推断获得默认读取前 100 行可通过配方中的max_rows参数控制JSONSchema 基于整个文件推断因为难以只抽取文件前几个对象可能影响性能项目正在研究基于迭代器iterator的 JSON 解析器以避免读入整个 JSON 对象。从配置模型看gcs_source.pymax_rows默认 100最小值 1同时还会约束 JSON/JSONL 推断——顶层数组最多读取这么多条记录单个 JSON 对象内的数组也会被截断到该数量因此仅在更靠后位置出现的字段不会被报告number_of_files_to_sample默认 100控制用于 Schema 推断的文件采样数当 path spec 的sample_files为false时被忽略。六、数据 ProfilingGCS 连接器支持数据 Profiling启用后会提取每个数据集的行数与列数每列在启用时null 计数与比例、distinct 计数与比例、最小值/最大值/均值/中位数/标准差及若干分位数、直方图或唯一值频率。实现要点见 gcs_post.mdProfiling 是纯 Python 实现构建于pyarrow与 Apache DataSketches 之上不需要 Spark、Hadoop 或 JVMdistinct 计数与分位数/直方图为近似值DataSketchesGCS 文件通过 S3 互操作端点、使用与摄取相同的凭据读取因此启用 Profiling 无需额外设置启用 Profiling 会拖慢摄取运行速度。配置入口为GCSSourceConfig中的profilingDataLakeProfilerConfig与profile_patterns默认放行所有表的AllowDenyPatterngcs_source.py二者都会被透传到内部的 S3 DataLake 配置。七、完整配方Recipe参考完整的五种配方示例可直接参考 gcs_recipe.yml覆盖HMAC 默认认证、GKE Workload Identity、WIF 配置文件、WIF 内联 dict、WIF JSON 字符串均以上文第 3 节的 YAML 为准。摄取命令与其他 DataHub 源一致datahub ingest -c recipe.yml八、故障排查如果摄取失败请按以下顺序排查凭据按所选认证方式核对 HMAC 密钥、ADC 或 WIF 配置是否正确注意各auth_type下凭据字段的互斥校验规则权限确认服务账号/外部身份具备Storage Object Viewer及对象列举、读取等元数据 API 权限连通性确认摄取环境可访问https://storage.googleapis.comS3 互操作端点并注意跨区域出口流量成本作用域/过滤检查path_specs的include/exclude与tables_filter_pattern是否误过滤了目标对象日志查看摄取日志中的源相关错误如 WIF 加载失败、ADC 加载失败时源码会抛出带修复指引的ValueError据此调整配置。结合仓库源码中的配置校验器绝大多数凭据类错误都会在配置加载阶段以清晰的报错信息提前暴露例如credential is required when auth_type is hmac、All path_spec.include should start with gs://、Cannot specify multiple WIF configuration options等可作为快速定位的依据。九、概念映射速查GCS 源概念与 DataHub 概念的对应关系详见 gcs/README.mdGCS 源概念DataHub 概念备注Google Cloud StorageData PlatformGCS 对象 / 包含对象的文件夹DatasetGCS bucketContainer子类型GCS bucketGCS folderContainer子类型Folder此外该集成还支持有状态删除检测stateful deletion detection桶、文件夹、数据集等实体的状态由摄取运行状态跟踪元数据被移除时可被检测并清理GCSSourceConfig继承自StatefulIngestionConfigBase并透传stateful_ingestion配置。延伸阅读Path Specs 的 S3 侧完整语义可查阅 S3 Data Lake 连接器文档 s3 对应章节PathSpec 模型与遍历实现位于 data_lake_common/path_spec.pyGCS 连接器实现本体与配置校验见 gcs_source.py。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联 返回资讯列表 →