ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

mage-ai Google Cloud Storage 数据源接入指南:配置、鉴权与文件读取实现解析

mage-ai Google Cloud Storage 数据源接入指南:配置、鉴权与文件读取实现解析 数据工程数据编排ETL任务调度批处理流处理数据集成后端【免费下载链接】mage-ai Build, run, and manage data pipelines for integrating and transforming data.项目地址https://gitcode.com/gh_mirrors/ma/mage-ai点击查看免费下载本篇指南聚焦 mage-ai 开源数据管道平台中的Google Cloud StorageGCS数据源位于仓库 mage_integrations/mage_integrations/sources/google_cloud_storage完整说明该数据源的配置项、服务账号鉴权方式与启动步骤并结合仓库源码剖析其对象发现discover、数据加载load与连接测试test的底层实现帮助你正确搭建从 GCS 存储桶读取 CSV/Parquet 文件的数据集成管道。一、数据源定位与核心能力mage-ai 中的数据集成Data Integration框架遵循 Singer 规范实现了一套 Source 抽象Google Cloud Storage 数据源是该框架下的一个具体实现类定义于 mage_integrations/mage_integrations/sources/google_cloud_storage/init.py 中的GoogleCloudStorage(Source)类。该数据源解决的核心问题是按前缀prefix扫描指定 GCS 存储桶bucket内的对象blob过滤出目标文件类型自动推断每列的数据类型并生成 Catalog数据目录随后按流stream读取数据。它支持两种文件格式parquet与csv读取时使用 Pandas 完成数据帧构建。从代码结构看其工作链路为数据源层Source→ 连接层mage_integrations/mage_integrations/connections/google_cloud_storage/init.py 中的GoogleCloudStorage(Connection)→ Google Cloud Python 官方客户端google.cloud.storage.Client。二、配置项说明配置该数据源时需要提供以下凭证与参数表格内容与数据源官方文档一致并补充了仓库模板中的默认值Key说明示例值path_to_credentials_json_fileGoogle 服务账号凭据 JSON 文件路径。若 Mage 运行在 GCP 上可将此项置空Mage 将使用实例服务账号instance service account完成鉴权/path/to/service_account_credentials.jsonbucket数据要保存/读取的 Google Cloud Storage 存储桶名称my_bucketfile_typeGoogle Cloud Storage 文件的类型。支持的文件类型值parquet、csvparquetcredentials_info服务账号凭据的另一种指定方式内联凭据信息对象结构见下文见下方结构说明prefix用于过滤存储桶中对象的前缀字符串仅考虑名称以该前缀开头的对象sales_data_仓库中对应的配置模板位于 mage_integrations/mage_integrations/sources/google_cloud_storage/templates/config.json模板给出的默认结构如下可作为新建数据源时的参考起点{ bucket: , file_type: parquet, path_to_credentials_json_file: /path/to/your/service/account/key.json, credentials_info: null, prefix: }可见file_type默认值为parquet而credentials_info默认为null——即默认走path_to_credentials_json_file文件路径方式。credentials_info的结构credentials_info用于直接内联服务账号凭据例如从密钥管理系统中注入其结构如下type: str project_id: str private_key_id: str private_key: str client_email: str client_id: str auth_uri: str token_uri: str auth_provider_x509_cert_url: str client_x509_cert_url: str universe_domain: str该结构与 Google Cloud 服务账号 JSON 凭据文件的字段一一对应。连接层使用google.oauth2.service_account提供的Credentials.from_service_account_info(...)将其解析为可用的凭据对象详见下文鉴权实现一节。仓库中定义的类型声明 mage_integrations/mage_integrations/connections/utils/google.py 使用TypedDict声明了这些字段其中包含auth_provider_x509_cert_url、auth_uri、client_email、client_id、client_x509_cert_url、private_key、private_key_id、project_id、token_uri、type等字段。两种鉴权方式的选择根据连接层实现 mage_integrations/mage_integrations/connections/google_cloud_storage/init.py鉴权逻辑按优先级处理内联凭据优先若credentials_info不为None则调用service_account.Credentials.from_service_account_info(self.credentials_info)从字段字典直接构建凭据凭据文件次之若未提供credentials_info但提供了path_to_credentials_json_file则调用service_account.Credentials.from_service_account_file(self.path_to_credentials_json_file)读取文件构建凭据实例服务账号兜底当两者均为空时Credentials参数为None此时 Google Cloud 客户端会自动使用运行环境如 GCP 虚拟机/容器的元数据服务获取实例服务账号Application Default Credentials这正是官方文档中如果 Mage 运行在 GCP 上可留空的实现依据。数据源层的build_client()方法位于__init__.py会把上述两个配置项一并传给连接层最终返回google.cloud.storage.Client客户端对象。三、Get Started快速开始按照以下步骤完成数据源的初始化配置启用 Google Cloud Storage API在 Google Cloud Console 中启用 Cloud Storage API可通过控制台的 API 库找到Cloud Storage服务并启用。创建服务账号前往 Google Cloud Console 的 IAM 与管理 → 服务账号页面创建专属服务账号。下载凭据文件为服务账号生成 JSON 格式的密钥文件并下载将其放入你的 Mage 项目中例如项目根目录下的service_account.json随后在数据源配置中将path_to_credentials_json_file指向该文件。授予权限为便于与存储桶交互建议为该服务账号授予Storage Admin角色若遵循最小权限原则也可参考数据集成场景仅授予Storage Object Viewer读取对象与Storage Object Creator写入对象等更细粒度角色。完成上述配置后即可在 Mage 平台中通过数据集成管道创建该数据源并配置目标端进行数据同步。四、源码级实现剖析数据流如何工作1. 连接测试test_connectiontest_connection()方法mage_integrations/mage_integrations/sources/google_cloud_storage/init.py用于验证配置是否有效def test_connection(self) - None: client self.build_client() if not client.get_bucket(self.bucket).exists(): raise Exception(fBucket {self.bucket} does not exist.) client.list_blobs(self.bucket)它首先通过get_bucket(self.bucket).exists()检查存储桶是否存在不存在则抛出异常随后调用list_blobs验证列表权限确保服务账号具备读取该桶的访问能力。这一方法在数据源配置保存或管道测试连接时被调用可快速反馈鉴权与权限问题。2. 流发现discoverdiscover()方法负责扫描存储桶并生成数据目录Catalog核心逻辑为for blob in client.list_blobs(self.bucket, prefixself.prefix): if blob.size 0: continue key blob.name if not key.endswith(f.{self.file_type}): continue parts key.split(.) stream_id _.join(parts[:-1]) ...可以看出以prefix作为前缀过滤条件列出所有对象空对象size 0会被跳过对象名必须以.file_type即.parquet或.csv结尾否则被过滤掉流 IDstream_id由对象文件名去掉扩展名后以_连接而成例如sales_data_2024.parquet会被解析为流sales_data_2024。对于每个通过过滤的对象__build_df()读取文件内容随后对每个非空列调用infer_dtypes推断类型遇到mixed混合类型时通过Counter统计各实际 Python 类型出现次数并取占比最大的类型映射为COLUMN_TYPE_ARRAY列表、COLUMN_TYPE_OBJECT字典或COLUMN_TYPE_STRING字符串其余类型经convert_data_type转换为标准类型。最终生成每个列的类型为[null, col_type]的 JSON Schema并用get_standard_metadata生成标准元数据其中key_properties[]无主键replication_methodREPLICATION_METHOD_FULL_TABLE全表复制unique_conflict_methodUNIQUE_CONFLICT_METHOD_UPDATE唯一冲突时更新即该数据源默认采用全量同步策略同步时对已存在的记录按更新方式处理冲突。3. 数据加载load_dataload_data()是实际数据读取的入口for blob in client.list_blobs(self.bucket, prefixself.prefix): if blob.size 0: continue key blob.name stream_id _.join(key.split(.)[:-1]) if stream_id in self.selected_streams: df self.__build_df(blob.name) yield df.to_dict(records)它同样按前缀扫描对象、跳过空对象然后仅对在 Catalog 中被用户勾选selected_streams的流构建 DataFrame并以df.to_dict(records)的形式逐条产出记录供下游目标端消费。4. 文件读取细节__build_dfdef __build_df(self, key: str) - pd.DataFrame: client self.build_client() bucket client.get_bucket(self.bucket) blob bucket.get_blob(key) data blob.download_as_bytes() buffer io.BytesIO(data) if .parquet in key: df pd.read_parquet(buffer) elif .csv in key: encoding from_bytes(data).best().encoding df pd.read_csv(buffer, encodingencoding) return df文件以字节流方式下载后按类型解析Parquet 文件直接pd.read_parquet(buffer)读取保留列式存储的 schemaCSV 文件使用charset_normalizer的from_bytes(data).best().encoding自动探测文件编码再以探测出的编码执行pd.read_csv从而兼容 UTF-8、GBK 等常见编码的 CSV。值得注意的是文件格式判定使用.parquet in key与.csv in key的包含匹配与 discover 阶段的endswith严格后缀匹配略有不同实际使用时建议保持对象命名规范让文件名以明确的扩展名结尾。五、与其他 GCS 使用方式的边界需要说明的是本数据源Source解决的是数据集成管道中从 GCS 作为输入源同步到目标端的场景。mage-ai 项目中还存在另一类 GCS 接入方式直接面向 Python 块data loader / data exporter的 IO 类 mage_ai/io/google_cloud_storage.py它通过项目中的io_config.yaml配置如GOOGLE_SERVICE_ACC_KEY_FILEPATH、GOOGLE_SERVICE_ACC_KEY等键实现块级读写并支持GOOGLE_APPLICATION_CREDENTIALS环境变量兜底。两种方式的定位不同本指南的数据源面向数据集成管道关注 Catalog 发现、全量复制策略与多对象批量同步配置项为本文第二节所列五项IO 类面向普通数据加载器/导出器块按对象键精确读写单个文件。实际项目中可按需选用若你的目标是搭建可复用的 GCS → 其他数据存储的持续同步管道则优先采用本文所述数据源。六、配置要点与排障建议鉴权二选一path_to_credentials_json_file与credentials_info任选其一即可若两者都为空仅在运行环境为 GCP具备实例服务账号时可正常工作本地环境会因缺少 Application Default Credentials 而报鉴权错误。权限最小化验证先确认服务账号对目标桶具备storage.objects.list与storage.objects.get权限对应Storage Object Viewer否则test_connection阶段即会失败。前缀语义prefix是对象名前缀匹配而非目录语义sales_data_会匹配所有以该字符串开头的对象结合file_type过滤可精准圈定同步范围。流命名影响对象文件名直接决定流 ID多个对象可能解析出相同流 ID同步时会被视为同一流的多个分片规划对象命名时应保持一致性。文件编码CSV 文件建议保持统一编码虽然数据源内置了编码自动探测但编码探测存在一定不确定性若同步结果出现乱码可优先检查源文件编码。配置完成后即可在 Mage 平台的数据集成页面创建管道将 GCS 数据源与任意已支持的目标端连接实现定时或手动触发的数据同步任务。赞分享数据工程数据编排ETL任务调度批处理流处理数据集成后端【免费下载链接】mage-ai Build, run, and manage data pipelines for integrating and transforming data.项目地址https://gitcode.com/gh_mirrors/ma/mage-ai点击查看免费下载相关推荐mage-ai Google Analytics 数据源接入指南配置、认证与 GA4 报告数据抽取mage ai Google Analytics 数据源接入指南配置、认证与 GA4 报告数据抽取 导读 本文面向在 mage ai 数据集成框架中接入 Go数据工程数据编排ETL任务调度批处理流处理数据集成后端前端Mage-ai 集成 Zendesk 数据源配置、OAuth 鉴权与增量同步实现解析Mage ai 集成 Zendesk 数据源配置、OAuth 鉴权与增量同步实现解析 Mage ai 通过 mage_integrations 中的数据集成框数据工程数据编排ETL任务调度批处理流处理数据集成后端前端Mage 数据集成Google Cloud StorageGCS目标端配置与写入原理详解Mage 数据集成Google Cloud StorageGCS目标端配置与写入原理详解 导读 本文讲解 Mage 数据集成框架中 Google Clou数据工程数据编排ETL任务调度批处理流处理数据集成后端前端上一篇HackRF One硬件架构完整解析从微控制器到射频前端的终极设计指南下一篇learnyounode练习解析每个模块的学习重点与难点突破创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表