
数据工程数据编排ETL任务调度批处理流处理数据集成后端【免费下载链接】mage-ai Build, run, and manage data pipelines for integrating and transforming data.项目地址https://gitcode.com/gh_mirrors/ma/mage-ai点击查看免费下载输出文章mage-ai Azure Blob Storage 数据源接入指南配置、发现原理与数据同步实战在 mage-ai 的数据集成框架mage_integrations中Azure Blob Storage 是官方内置的数据源之一。它允许你直接以容器内的 CSV / Parquet 文件为数据源经模式发现discover、目录生成catalog与全表同步sync等标准流程将对象存储中的文件数据接入下游管道。本文将围绕 azure_blob_storage 数据源的官方说明展开先给出完整的配置参数与连接串获取方法再深入其 源码实现讲解模式发现、类型推断、同步与连接测试的底层机制让你既能照配置立刻跑通也能理解它为什么这么工作。一、配置参数详解在 Mage UI 中新建 Data Integration Pipeline 并选择Azure Blob Storage作为 Source 时需要提供以下三项凭据。它们的完整定义、说明与示例值如下Key说明示例值connection_string存储账户Storage account的连接字符串包含账户端点与鉴权信息。BlobEndpointhttps://xxx.blob.core.windows.net/;yyyyyysigtestsigcontainer_name数据所在的 Blob 存储容器Container名称。testcontainerprefix文件所在的路径前缀。注意不要在此值中包含容器名或表名。users/ds/20221225这三项参数与仓库内的模板文件完全对应见 templates/config.json{ connection_string: , container_name: , prefix: null }其中prefix默认为null即不传前缀时默认扫描整个容器。1. connection_string连接字符串连接字符串Connection String是整个接入的鉴权核心。它通常形如BlobEndpointhttps://storage-account.blob.core.windows.net/;SharedAccessSignature...sig...它同时携带端点地址BlobEndpoint与鉴权凭据因此源码中可以直接用它构造BlobServiceClient而无需单独配置账号密钥或令牌。从源码结构看AzureBlobStorage类正是将配置中的connection_string直接传给官方 SDK 完成客户端构建见init.py。2. container_name容器名称Blob 存储以容器为顶层命名空间数据文件都存放在容器内部。容器名的典型取值范围为 363 个小写字母、数字与连字符本数据源对其不做额外转换直接以字符串透传给 Azure SDK。3. prefix路径前缀prefix用于圈定待同步的文件范围。例如仓库文档示例中的users/ds/20221225它表示容器内以该路径开头的所有对象即只同步这个目录含子目录下的文件配置中不能把容器名testcontainer或表名写进prefix否则会拼出不存在的路径源码中prefix被原样传给container_client.list_blobs(prefix)用于列取符合条件的 Blob 对象。二、如何获取连接字符串原文档给出了两种获取connection_string的官方途径方式一从存储账户的访问密钥Access Keys获取在 Azure 门户进入你的 Storage account → “Access keys” 页面可以查看主/次连接字符串直接复制使用。这是最常见、权限最完整的方式适合该连接串只在内网环境中使用的场景。方式二生成共享访问签名SAS后拼装连接字符串如果你希望把访问范围限制在存储账户下的某个资源上例如只授予某个容器或某个 Blob 的读写权限可以在 Azure 门户生成 SASShared Access Signature令牌再以“连接字符串 SAS 令牌”的形式拼出受限的连接字符串。这种方式适合需要最小化授权、或需要把凭据下发给第三方/外部任务使用的场景。在 Mage 中连接测试Test connection按钮会实际调用源码里的test_connection方法先判断容器是否存在再尝试列取prefix下的第一个 Blob任一环节失败都会直接报错因此在保存前可以立即验证连接串与路径是否正确。三、源码级原理数据源如何工作AzureBlobStorage继承自 Source 基类并实现了discover、load_data、test_connection三个核心方法。它的完整实现位于 mage_integrations/mage_integrations/sources/azure_blob_storage/init.py。3.1 底层客户端基于 Azure 官方 SDKdef build_client(self): return BlobServiceClient.from_connection_string(self.config[connection_string])build_client使用azure.storage.blob.BlobServiceClient.from_connection_string从连接串直接构造服务客户端随后通过get_container_client(self.container_name)拿到容器客户端再调用list_blobs(self.prefix)枚举该前缀下的所有 Blob 对象。也就是说整个数据源只依赖三项配置即可完成从鉴权到列对象的全过程。3.2 discover文件级模式发现与流Stream生成discover方法源码 L35-L106负责枚举文件并生成 Catalog遍历prefix下的所有 Blob跳过size 0的空对象以文件名的目录层级生成stream_id将路径按/切分后把除最后一个文件名之外的部分用_连接例如users/ds/20221225/data.csv会被映射为流users_ds_20221225读取文件内容为 DataFrame对每个非空列做类型推断生成 JSON Schema以FULL_TABLE复制方式、UPDATE冲突处理方式构建CatalogEntry只取第一个文件生成流后即break即以第一个文件的列结构作为该流的 Schema 基线。其中类型推断逻辑值得一提源码 L53-L79对每一列先取非空值调用 infer_dtypes基于 pandas 的infer_dtype推断若推断结果为mixed列内类型混杂则统计各 Python 类型的出现次数取出现最多者作为列类型list→array、dict→object其余一律按string处理推断出的列类型会映射为 JSON Schema 中的type: [null, col_type]即所有列都允许为 NULL避免因个别空值导致同步失败。最终每个流都会携带标准元数据get_standard_metadata其中的forced-replication-method固定为FULL_TABLE并把unique_conflict_method设为UPDATE对应 constants.py 中的REPLICATION_METHOD_FULL_TABLE与UNIQUE_CONFLICT_METHOD_UPDATE。3.3 load_data全表式数据读取load_data方法源码 L108-L121是同步阶段的实际数据来源再次列取prefix下的所有 Blob跳过空文件对每个非空文件调用__build_df读取为 DataFrame并以df.to_dict(records)逐文件产出记录批次。__build_df是文件格式分派的关键源码 L131-L146blob_data blob_client.download_blob() buffer io.BytesIO(blob_data.readall()) if .parquet in key: df pd.read_parquet(buffer) elif .csv in key: df pd.read_csv(buffer)文件下载到内存BytesIO缓冲区后按扩展名解析当前支持.parquet与.csv两种格式文件不存在时返回空 DataFrame保证流程不中断由于discover与load_data都按list_blobs(prefix)的结果迭代且复制方式为FULL_TABLE因此每次运行都会把前缀下的全部文件当作整表重新同步。3.4 整体处理流程discover → catalog → sync当以命令行方式运行该数据源时会执行main(AzureBlobStorage)见文件末尾进入 Source.process() 的完整流程test_connection模式仅验证连接并退出discover模式输出 CatalogJSON Schema 元数据count_records模式统计可同步记录数普通同步模式基于 Catalog 中选中的流依次写出SCHEMA消息、读取load_data的每一批记录、写出RECORD消息与STATE消息。基类还支持--config、--state、--catalog、--query、--test_connection、--load_sample_data等标准 CLI 参数见 sources/utils.py 中的 parse_args因此在 Mage 之外也可以直接用命令行跑通该数据源。四、在 Mage 中使用该数据源的完整步骤结合上述原理在 mage-ai 中使用 Azure Blob Storage 数据源的实操路径如下准备存储在 Azure 门户创建 Storage account 与容器上传 CSV / Parquet 文件并按业务日期/分区组织目录如users/ds/20221225/。获取连接串按“如何获取连接字符串”一节中的任一方式取得connection_string建议生产环境使用 SAS 受限连接串。配置 Source在 Mage 的 Data Integration Pipeline 中选择 Azure Blob Storage依次填写connection_string、container_name、prefix三项配置字段含义与示例见上文参数表。测试连接点击连接测试源码会校验容器存在性并尝试列出前缀下的首个 Blob若失败请检查连接串权限、容器名拼写与前缀路径。Discover 与配置 Schema运行模式发现后Mage 会展示自动推断出的流与列类型由于key_properties为空表级主键需由你在 Catalog 配置中自行指定。运行同步以全表复制方式运行管道load_data将逐文件读取并写入目标端。五、注意事项与适用边界复制方式固定为全量源码将流的replication_method固定为FULL_TABLE不支持基于bookmark的增量复制每次运行都会重读前缀下的全部文件。从源码结构看若需要增量应在管道层做增量调度或按日期前缀拆分多个流。文件格式限制当前仅支持.csv与.parquet其他格式如 JSON Lines、Excel不会被解析会被视为空 DataFrame 而跳过。类型推断以首个文件为基线discover仅取第一个非空文件构建 Schemabreak逻辑若同一前缀下不同文件列结构不一致可能出现 Schema 与后续文件不匹配的问题建议同一prefix下保持统一的表结构。空文件会被跳过size 0的 Blob 在发现与同步阶段都会被过滤避免产生空流或空批次。权限最小化建议若连接串使用 SAS需确保其授权范围覆盖目标容器与prefix前缀下的列list与读read权限否则test_connection会失败。六、参考文件索引官方数据源说明配置参数与连接串获取方法。源码实现discover/load_data/test_connection/__build_df完整实现。配置模板三个配置项的默认模板。Source 基类process/sync整体处理流程与 CLI 入口main。参数解析parse_args与get_standard_metadata。类型映射常量COLUMN_TYPE_*与REPLICATION_METHOD_FULL_TABLE、UNIQUE_CONFLICT_METHOD_UPDATE的定义。类型推断工具infer_dtypes/convert_data_type的实现。通过以上配置与源码层面的解析你应该可以独立完成 Azure Blob Storage 数据源的接入、排错与二次定制。赞分享数据工程数据编排ETL任务调度批处理流处理数据集成后端【免费下载链接】mage-ai Build, run, and manage data pipelines for integrating and transforming data.项目地址https://gitcode.com/gh_mirrors/ma/mage-ai点击查看免费下载相关推荐5关做完你的VPet桌面宠物MOD新手向开源桌宠教程5关做完你的VPet桌面宠物MOD新手向开源桌宠教程 VPet 是一个开源虚拟桌宠模拟器能把宠物嵌进任何 WPF 应用里。你画的图能在桌面上走动、被摸、会桌面应用游戏开发DataHub Azure Blob Storageabs数据源接入指南从 Blob 到 Dataset 的元数据治理实践DataHub Azure Blob Storageabs数据源接入指南从 Blob 到 Dataset 的元数据治理实践 Azure Blob Stor数据目录数据治理数据血缘后端前端数据工程数据集成HYG-Database许可证变更从CC BY-SA 2.5到4.0的重要变化HYG Database许可证变更从CC BY SA 2.5到4.0的重要变化 HYG Database是一个重要的恒星数据库项目近期其许可证从CC BY数据工程数据编排ETL任务调度批处理流处理数据集成后端前端上一篇超越传统OCRpaddlepaddle/devanagari_PP-OCRv5_mobile_rec_safetensors移动端部署最佳实践下一篇Apollo Link REST 终极性能优化指南缓存策略和请求合并的10个最佳实践 创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考