数据工程数据编排ETL任务调度批处理流处理数据集成后端【免费下载链接】mage-ai Build, run, and manage data pipelines for integrating and transforming data.项目地址https://gitcode.com/gh_mirrors/ma/mage-ai点击查看免费下载本指南以开源仓库mage_integrations中 BigQuery Source 模块为核心系统讲解在 Mage AI 数据集成管道中接入 Google BigQuery 所需的全部配置项、服务账号认证方式、最小权限要求以及底层同步机制。读完本文你将能够独立完成 BigQuery Source 的凭证配置、批量读取调优并理解数据发现discover与分页拉取load在源码层面的执行原理。概览BigQuery Source 在 Mage 中的定位Mage AI 的数据集成体系mage_integrations把“从各类数据源读取数据”的能力抽象为Source把“写入各类目标”的能力抽象为Destination。BigQuery Source 位于 sources/bigquery它继承自 SQL 型 Source 基类通过 Google Cloud 官方 Python SDK 连接 BigQuery将数据集Dataset下的表暴露为可供下游管道消费的数据流Stream。该 Source 的核心能力包括表/元数据发现discover自动读取指定 dataset 下所有表的列信息生成符合 Singer 规范的 Catalog。全量同步full table replication按批次batch拉取数据支持增量字段bookmark过滤。列名规范化自动为列名添加 BigQuery 反引号转义。配置要求三类必需凭证项在 Mage UI 中创建该 Source 时必须填写以下凭证信息Key说明示例值path_to_credentials_json_fileGoogle 服务账号凭证 JSON 文件的路径。如果 Mage 运行在 GCP 上可以将此值留空Mage 会改用实例服务账号instance service account进行认证。/path/to/service_account_credentials.jsondataset你要读取数据的 BigQuery 数据集名称。example_datasetcredentials_info直接以内联结构化字段提供 Google 服务账号凭证是path_to_credentials_json_file的替代方案。结构见下文配置模板可在 sources/bigquery/templates/config.json 中找到默认仅包含两个占位字段{ path_to_credentials_json_file: path_to_credentials_json_file, dataset: }两种认证方式的取舍从源码 connections/bigquery/init.py 可以看到两种认证方式的优先级与底层实现if self.credentials_info is not None: if isinstance(self.credentials_info, dict): self.credentials_info service_account.Credentials.from_service_account_info( self.credentials_info, ) elif self.path_to_credentials_json_file is not None: self.credentials_info service_account.Credentials.from_service_account_file( self.path_to_credentials_json_file, ) self.client Client(credentialsself.credentials_info, locationself.location)当配置了credentials_info且为 dict 类型时优先使用from_service_account_info从内存中的字段构建凭证否则若提供了path_to_credentials_json_file则通过from_service_account_file从本地文件读取两者都未提供时credentials_info为Nonegoogle.cloud.bigquery.Client会回退到 Application Default CredentialsADC这正是 README 所说“Mage 运行在 GCP 上可使用实例服务账号认证”的底层机制。建议在非 GCP 环境如本地开发、自建服务器优先使用path_to_credentials_json_file便于密钥管理与轮换credentials_info适合需要把密钥直接写入管道配置如加密后的环境变量的场景。切勿将明文私钥提交到版本库。credentials_info的完整结构credentials_info对应一个 Google 服务账号 JSON 中的所有字段。该结构在源码中被定义为CredentialsInfoTypeTypedDict见 connections/utils/google.py字段类型说明auth_provider_x509_cert_urlstr认证提供方证书 URL服务账号 JSON 固定字段auth_uristrOAuth 认证 URI固定为https://accounts.google.com/o/oauth2/authclient_emailstr服务账号邮箱如xxxproject.iam.gserviceaccount.comclient_idstr服务账号客户端 IDclient_x509_cert_urlstr客户端证书 URLprivate_keystrRSA 私钥内容注意需保留原始换行符private_key_idstr私钥 IDproject_idstrGCP 项目 IDtoken_uristrToken 获取 URI固定为https://oauth2.googleapis.com/tokentypestr固定为service_account以 YAML 形式示意如下auth_provider_x509_cert_url: str auth_uri: str client_email: str client_id: str client_x509_cert_url: str private_key: str private_key_id: str project_id: str token_uri: str type: str在 GCP Console 中创建服务账号密钥JSON 类型后下载的文件内容即与上述字段一一对应可直接复制粘贴进credentials_info。可选配置批量读取大小Key说明示例值batch_fetch_limit每批拉取的行数默认 50000。如果你的实例内存更大可以调大批次大小以获得更少的网络往返。50000该配置在源码层面有明确的默认值与读取链常量定义见 sources/constants.pyBATCH_FETCH_LIMIT 50000、SUBBATCH_FETCH_LIMIT 10000、键名batch_fetch_limit/subbatch_fetch_limit。实际生效逻辑位于 SQL Source 基类 sources/sql/base.pyproperty def fetch_limit(self): config self.config or dict() return ( config.get(SUBBATCH_FETCH_LIMIT_KEY) or config.get(BATCH_FETCH_LIMIT_KEY) or BATCH_FETCH_LIMIT )即优先级为subbatch_fetch_limit→batch_fetch_limit→ 默认 50000。若需控制单批内存占用例如下游处理单元较小可通过subbatch_fetch_limit进一步细分批次。所需权限最小权限集使用 BigQuery Source 的服务账号至少需要以下 IAM 权限bigquery.jobs.create bigquery.readsessions.create bigquery.readsessions.getData此外该账号还需要对配置中指定的 dataset 拥有BigQuery Data Viewer角色。bigquery.jobs.create允许发起查询作业对应同步时执行的SELECT语句bigquery.readsessions.create与bigquery.readsessions.getData对应 BigQuery Storage Read API 的读取会话能力批量读取数据时使用。实操建议在 GCP IAM 中创建一个专用服务账号仅授予上述权限并绑定到目标 dataset 的roles/bigquery.dataViewer然后将该账号的密钥文件路径填入path_to_credentials_json_file。源码视角BigQuery Source 的核心实现1. 数据集前缀与连接构建BigQuery Source 类 将dataset配置直接用于两处property def table_prefix(self) - str: dataset self.config[dataset] return f{dataset}. def build_connection(self) - BigQueryConnection: return BigQueryConnection( credentials_infoself.config.get(credentials_info), path_to_credentials_json_fileself.config.get(path_to_credentials_json_file), )table_prefix让所有生成的 SQL 以dataset.表名形式限定命名空间这正是 README 要求dataset必填的原因build_connection把两类凭证配置透传给 BigQuery 连接层见上文认证逻辑。2. 元数据发现Discoverbuild_discover_queryinit.py通过查询INFORMATION_SCHEMA.COLUMNS获取表名、列名、数据类型、可空性等信息并支持用streams参数仅发现指定表SELECT table_name , column_default , NULL AS column_key , column_name , data_type , is_nullable FROM {dataset}.INFORMATION_SCHEMA.COLUMNS查询结果随后由基类的discover方法sources/sql/base.py转换为标准 Singer Catalog每张表成为一个 stream默认复制方式为REPLICATION_METHOD_FULL_TABLE全量同步并将 BigQuery 的日期/时间/JSON 等类型映射为标准的字符串/对象等 JSON Schema 类型。3. 列名转义与时间戳处理BigQuery 列名可能包含特殊字符因此该 Source 重写了列名处理逻辑init.pydef update_column_names(self, columns: List[str]) - List[str]: return [f{column} for column in columns] def wrap_column_in_quotes(self, column: str) - str: if not in column: return f{column} return column同时它针对COLUMN_FORMAT_DATETIME日期时间格式列重写了column_type_mapping在基于 bookmark 值比较时不做时间类型转换从而保证增量同步边界值的精确匹配def column_type_mapping(self, column_type: str, column_format: str None) - str: if COLUMN_FORMAT_DATETIME column_format: # Not cast datetime value type when comparing bookmark values return None return super().column_type_mapping(column_type, column_format)4. 分页拉取Load数据读取由基类load_datasources/sql/base.py驱动按fetch_limit与offset循环执行LIMIT n OFFSET m查询直到取完所有行每次批次间暂停 1 秒避免对 BigQuery 造成突发负载。生成的 SQL 会自动对 bookmark 列、主键列、唯一约束列做ORDER BY保证分页稳定性在有 bookmark 时生成WHERE 列 值或条件唯一约束列用在count_records模式下改写为COUNT(*) AS number_of_records用于同步前的记录数统计。连接层通过google.cloud.bigquery.dbapi执行 SQL见 connections/bigquery/init.py。测试与验证仓库中已有 BigQuery 连接与目标侧的自动化测试tests/destinations/bigquery/test_bigquery.py验证了path_to_credentials_json_file配置到连接层的透传、SQL 建表命令生成如CREATE TABLE test.test_table (\ID STRING)等行为。Source 侧的连通性同样可以复用该认证链路test_connectionsources/sql/base.py会实际建立一次连接来校验凭证与网络可用性。配置速查清单完成 BigQuery Source 接入时对照以下清单逐项确认已创建具备bigquery.jobs.create、bigquery.readsessions.create、bigquery.readsessions.getData权限且绑定目标 datasetBigQuery Data Viewer角色的服务账号已准备服务账号 JSON 文件路径path_to_credentials_json_file或完整的credentials_info字段dataset填写正确且与服务账号可访问的数据集一致如需更快同步可适当调大batch_fetch_limit默认 50000并留意实例内存在 Mage UI 中保存配置并执行连接测试确认 discover 能列出目标表后再创建管道。补充说明BigQuery 同样支持作为 Destination 写入数据其配置与认证方式与 Source 一脉相承相关说明见 destinations/bigquery/README.md可作为构建“BigQuery → BigQuery”或“多源 → BigQuery”管道的参考。赞分享数据工程数据编排ETL任务调度批处理流处理数据集成后端【免费下载链接】mage-ai Build, run, and manage data pipelines for integrating and transforming data.项目地址https://gitcode.com/gh_mirrors/ma/mage-ai点击查看免费下载相关推荐Yox超轻量MVVM框架入门指南 - 快速构建现代Web应用Yox超轻量MVVM框架入门指南 快速构建现代Web应用 Yox 是一款专为现代Web应用设计的超轻量级MVVM框架它借鉴了Vue.js的优秀设计理念数据工程数据编排ETL任务调度批处理流处理数据集成后端前端终极指南如何将Kubeconform作为Go模块集成到你的应用中终极指南如何将Kubeconform作为Go模块集成到你的应用中 Kubeconform是一款快速的Kubernetes清单验证工具支持自定义资源。本文将详数据工程数据编排ETL任务调度批处理流处理数据集成后端前端Datejs核心原理深度剖析从源码看日期处理机制Datejs核心原理深度剖析从源码看日期处理机制 Datejs作为一款强大的JavaScript日期时间库通过对原生Date对象的扩展和创新的解析引擎为开数据工程数据编排ETL任务调度批处理流处理数据集成后端前端上一篇Cocos 粒子系统实战指南做出好看又不掉帧的粒子效果下一篇探索Kubernetes故障排查从新手到英雄创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
