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 Storage(GCS)连接器是 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 集成源(
S3Source),Path Specs 的详细语义与 S3 连接器保持一致。
从源码看,这一"代理"结构非常直观:GCSSource在初始化时通过create_equivalent_s3_source构建一个内部的S3Source实例,并将自身所有工作单元委托给它(gcs_source.py):
create_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.py),region 固定为auto;get_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 | 可选启用(通过配置开启) |
当前该源的官方支持状态为BETA(SupportStatus.BETA)。
二、前置条件
运行摄取之前,需要确保:
- 网络连通性:DataHub 摄取环境能够访问 GCS 的 S3 互操作端点
https://storage.googleapis.com; - 有效认证凭据:根据所选认证方式准备好对应凭据(详见下文);
- 读取权限:服务账号/外部身份需具备该模块所需元数据 API 的读取权限,最典型的是
Storage Object Viewer(roles/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 Credentials(ADC) | DataHub 运行在 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.py):
source: 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.py),HMAC 模式下连接器直接构造标准的AwsConnectionConfig:aws_endpoint_url指向https://storage.googleapis.com,aws_access_key_id/aws_secret_access_key填入 HMAC 密钥对,aws_region固定为auto。这正是"GCS 的 S3 互操作端点 + HMAC 密钥即 S3 风格的访问凭据"的实现体现。
配置校验(重要):GCSSourceConfig.validate_credential(gcs_source.py)保证:
auth_type为hmac时,credential必填,否则直接报错credential is required when auth_type is 'hmac';auth_type非hmac时,credential必须为空,避免凭据误配。
3.2 GKE Workload Identity(推荐:DataHub 在 GKE 内)
当 DataHub 运行在 GKE 集群内部且已启用 Workload Identity 时,这是最简方案:无需任何凭据文件或 Secret,Pod 的 Kubernetes Service Account 会被自动使用。配置步骤:
在 GKE 集群上启用 Workload Identity;
创建一个具备
Storage Object Viewer角色的 Google Service Account(GSA);将 Kubernetes Service Account(KSA)绑定到该 GSA:
gcloud 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-account=GSA_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_credentials(gcs_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_id="gcs-oauth"、aws_secret_access_key="not-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 Federation(WIF)。配置步骤:
- 在 Google Cloud 中创建 Workload Identity Pool 与 Provider;
- 授予外部身份模拟某个具备
Storage Object Viewer角色的 GCP 服务账号的权限; - 从 Google Cloud Console 下载或生成 WIF 凭据配置文件(JSON 格式);
- 通过以下三种互斥方式之一将配置交给连接器:
| 配方字段 | 说明 |
|---|---|
gcp_wif_configuration | WIF 配置 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。该 Mixin(GCPWIFConfig)被设计为可被 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配方示例二:内联 dict
source: 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_credentials(gcs_source.py),最终调用load_wif_credentials(common/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_options(gcs_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 │ └── year=2022 │ └── month=2 │ ├── 1.parquet │ └── 2.parquet └── returns └── year=2021 └── month=2 └── 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 │ └── year=2022 │ └── month=2 │ ├── 1.parquet │ └── 2.parquet └── tmp_orders └── year=2021 └── month=2 └── 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 └── year=2022 └── month=2 ├── 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.py):
autodetect_partitions(默认true):当文件夹形如year=2024的{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.py):
- CSV
- TSV
- JSONL
- JSON
- Parquet
- Apache Avro
Schema 推断行为如下:
- Parquet 与 Avro:Schema 按文件原样提取(文件自带 schema);
- CSV、TSV、JSONL:Schema 通过推断获得,默认读取前 100 行,可通过配方中的
max_rows参数控制; - JSON:Schema 基于整个文件推断(因为难以只抽取文件前几个对象),可能影响性能;项目正在研究基于迭代器(iterator)的 JSON 解析器以避免读入整个 JSON 对象。
从配置模型看(gcs_source.py):max_rows默认 100(最小值 1),同时还会约束 JSON/JSONL 推断——顶层数组最多读取这么多条记录,单个 JSON 对象内的数组也会被截断到该数量,因此"仅在更靠后位置出现的字段不会被报告";number_of_files_to_sample默认 100,控制用于 Schema 推断的文件采样数(当 path spec 的sample_files为false时被忽略)。
六、数据 Profiling
GCS 连接器支持数据 Profiling,启用后会提取:
- 每个数据集的行数与列数;
- 每列(在启用时):null 计数与比例、distinct 计数与比例、最小值/最大值/均值/中位数/标准差及若干分位数、直方图或唯一值频率。
实现要点(见 gcs_post.md):
- Profiling 是纯 Python 实现,构建于
pyarrow与 Apache DataSketches 之上,不需要 Spark、Hadoop 或 JVM; - distinct 计数与分位数/直方图为近似值(DataSketches);
- GCS 文件通过 S3 互操作端点、使用与摄取相同的凭据读取,因此启用 Profiling 无需额外设置;
- 启用 Profiling 会拖慢摄取运行速度。
配置入口为GCSSourceConfig中的profiling(DataLakeProfilerConfig)与profile_patterns(默认放行所有表的AllowDenyPattern,gcs_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.com(S3 互操作端点),并注意跨区域出口流量成本; - 作用域/过滤:检查
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.md):
| GCS 源概念 | DataHub 概念 | 备注 |
|---|---|---|
"Google Cloud Storage" | Data Platform | |
| GCS 对象 / 包含对象的文件夹 | Dataset | |
| GCS bucket | Container | 子类型GCS bucket |
| GCS folder | Container | 子类型Folder |
此外,该集成还支持有状态删除检测(stateful deletion detection):桶、文件夹、数据集等实体的状态由摄取运行状态跟踪,元数据被移除时可被检测并清理(GCSSourceConfig继承自StatefulIngestionConfigBase,并透传stateful_ingestion配置)。
延伸阅读:Path Specs 的 S3 侧完整语义可查阅 S3 Data Lake 连接器文档 s3 对应章节;PathSpec 模型与遍历实现位于 data_lake_common/path_spec.py;GCS 连接器实现本体与配置校验见 gcs_source.py。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考