ARTICLE DETAIL

资讯详情

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

Daft 集成 AWS Glue:通过 GlueCatalog 读写 Glue 表的完整指南

Daft 集成 AWS Glue:通过 GlueCatalog 读写 Glue 表的完整指南 Daft 集成 AWS Glue通过 GlueCatalog 读写 Glue 表的完整指南【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft本篇指南讲解 Daft 如何通过其 Catalog 接口接入 AWS Glue Data Catalog包括如何创建并配置GlueCatalog实例、内置支持的四种表格式CSV、Parquet、Iceberg、Delta Lake及其读写能力边界、Glue/Hive 类型到 Daft 类型的实际转换规则以及如何注册自定义GlueTable实现来扩展不支持的表格式。读完本文你可以直接在自己的 AWS 环境中用 Daft DataFrame 读取 Glue 元数据表、将结果写回 Iceberg 表并在需要时按源码协议补齐自定义表格式。术语映射先理清 Catalog、Namespace 与 TableDaft 与 AWS Glue 的对接通过 Daft Catalog 接口完成。理解接入方式前必须先记住两套术语之间的映射关系这是使用 Glue 集成时最容易混淆的地方Daft 概念Glue 概念CatalogData Catalog服务级数据目录NamespaceDatabaseTableTable也就是说load_glue(name...)中的name只是给这个 Catalog 实例起的名字文档中建议遵循 Hive 兼容习惯使用全小写它并不直接对应某个 Glue Database真正的 Database 在 Daft 中表现为 Namespace表标识符采用database_name.table_name两段式写法。这一约束在源码中得到了印证GlueCatalog._get_table 与_drop_table都会校验标识符长度必须为 2否则抛出Expected identifier with form database_name.table_name的ValueError对应测试见 tests/catalog/test_glue.py。创建 GlueCatalog 实例的三种方式方式一load_glue 直接传 boto3 客户端配置最直接的入口是daft.catalog.__glue.load_glue其完整签名与参数说明如下源自 load_glue 的 docstring 实现参数类型说明namestrCatalog 名称建议全小写以兼容 Hive 命名习惯region_namestr可选要连接的 AWS 区域api_versionstr可选Glue 服务使用的 API 版本use_sslbool可选是否使用 SSL 连接verifybool \| str可选是否校验 SSL 证书或指向 CA bundle 的路径endpoint_urlstr可选替代端点地址可指向本地兼容端点aws_access_key_idstr可选认证用的 Access Key IDaws_secret_access_keystr可选认证用的 Secret Keyaws_session_tokenstr可选临时凭证的会话 Token函数内部仅将非None的参数收集进 options 字典再调用boto3.client(glue, **options)创建底层客户端——因此省略所有 AWS 参数时会走 boto3 的标准凭证链环境变量、IAM 角色、~/.aws/credentials等。方式二复用已有 boto3/botocore 客户端或会话如果项目中已经存在配置好的 AWS 客户端或会话可以通过静态方法直接传入避免重复构造import boto3 from daft.catalog import Catalog from daft.catalog.__glue import GlueCatalog # 传入现成的 boto3 glue client client boto3.client(glue, region_nameus-west-2) catalog GlueCatalog.from_client(my_glue_catalog, client) # 或传入 boto3.Session / botocore.session.Session sess boto3.session.Session() catalog GlueCatalog.from_session(my_glue_catalog, sessionsess)from_session同时兼容boto3.Session调用session.client(glue)和更底层的botocore.session.Session调用session.create_client(glue)类型不匹配会抛出TypeError见 GlueCatalog.from_session。注意GlueCatalog()直接调用构造函数是不被支持的会抛出ValueError只能通过load_glue、from_client或from_session构造。方式三公开入口 Catalog.from_glue面向一般用户更推荐走 Catalog.from_glue 这个工厂方法它要求必须提供client或session二者之一同时提供会报错都不提供则抛出Must provide either a client or session.。该方法还会在依赖缺失时给出明确的安装提示pip install -U daft[aws]对应的依赖定义在 pyproject.toml 中awsextra 包含boto3当前仓库锁定1.44.0与mypy-boto3-glue。读取表从元数据到 DataFrame加载 Catalog 后的核心读表流程非常短from daft.catalog.__glue import load_glue # 加载 Glue catalog 实例 catalog load_glue( namemy_glue_catalog, region_nameus-west-2 ) # 加载一张 Glue 表 tbl catalog.get_table(my_namespace.my_table) # 作为 DataFrame 读取 df tbl.read() df.show()get_table的底层解析逻辑值得注意GlueCatalog._get_table 先调用 Glue 的GetTableAPI 拿到原始Table元数据然后依次尝试_table_impls列表中每个GlueTable实现的from_table_info方法——每个实现负责判断元数据是否匹配自己的格式不匹配时抛ValueError以“让位”给下一个实现若所有实现都不匹配则抛出包含classification与table_type信息的ValueError。表不存在时抛出统一的NotFoundError。默认的_table_impls列表在 GlueCatalog.new中初始化为GlueCsvTable、GlueParquetTable、GlueIcebergTable、GlueDeltaTable四类内置实现。围绕 Catalog 的其他常用操作及边界行为如下均有 tests/catalog/test_glue.py 中基于 motomock_aws的测试覆盖列出表catalog.list_tables(my_namespace)列出整个 namespace 的表catalog.list_tables(my_namespace.table_prefix)会拆分为DatabaseNameExpression传给 Glue 的GetTables实现前缀过滤GlueCatalog._list_tables。注意Glue 端不支持省略 namespace 的全局列表list_tables(None)会抛出GlueCatalog requires the pattern to contain a namespace的ValueError且结果通过NextToken自动翻页取全量。列出 namespacecatalog.list_namespaces()会翻完所有GetDatabases分页由于 Glue 的get_databases尚无 pattern 过滤能力传入 pattern 会直接抛出ValueError源码注释说明是等待 AWS 提供官方 scheme 后再支持见 GlueCatalog._list_namespaces。namespace 生命周期create_namespace/create_namespace_if_not_exists/drop_namespace均已实现分别映射到 Glue 的create_database/delete_database重复创建抛出ValueError(... already exists)删除不存在的 namespace 抛出NotFoundError。删除表drop_table(db.table)映射到 GlueDeleteTable。Catalog 级别的读写捷径继承自 Catalog 基类后还可以在 Catalog 层面直接读写无需手动get_tabledf catalog.read_table(my_namespace.my_table) catalog.write_table(my_namespace.my_table, df, modeappend) # 或 modeoverwritewrite_table内部调用Table.write(df, mode...)按 mode 分派到具体实现的append/overwrite见 Table.write。表格式支持矩阵与实现细节Glue 数据目录本身只存元数据真正可读取的内容取决于表的底层格式。文档给出的支持矩阵如下表格式支持情况CSV读Parquet读Iceberg读、写Delta Lake读、写同时需要记住三条明确的限制CSV/Parquet 尚不支持 Hive 风格分区读取尚不支持在 Glue 中创建表GlueCatalog._create_table目前直接抛出NotImplementedError(Table creation not yet implemented)见 源码。各格式的判定与实现要点均可在 daft/catalog/__glue.py 中核对GlueCsvTable要求Parameters.classification CSV从StorageDescriptor.Location取 S3 路径从Columns解析 schemaskip.header.line.count 1时视为有表头delimiter参数控制分隔符默认逗号。读取时调用daft.io._csv.read_csv并强制infer_schemaFalse即完全使用 Glue 元数据中的 schema不做数据侧推断GlueCsvTable.read。GlueParquetTable要求classification parquet读取时同样infer_schemaFalse直接以 Glue 声明的 schema 打开 Parquet 文件GlueParquetTable。GlueIcebergTable要求table_type ICEBERG。实现上会借助 pyiceberg 的 Glue Catalog复用当前GlueCatalog持有的同一个 boto3 客户端gc.glue catalog._client把 Glue 表元数据转换为 PyIceberg Table再走 Daft 的 Iceberg 读写路径GlueIcebergTable._create_pyiceberg_table。读取时支持snapshot_id、branch、tag、ignore_corrupt_files四个选项由Table._validate_options严格校验传入其他选项会报错写入则通过df.write_iceberg(table, modeappend|overwrite)完成GlueIcebergTable.read。GlueDeltaTable要求table_type delta内部构造一个指向 S3 Location 的UnityCatalogTable读取走daft.io.delta_lake._deltalake.read_deltalake。需要提示的是从当前源码结构看GlueDeltaTable.append/overwrite目前仍会抛出NotImplementedErrorGlueDeltaTable即 Delta 路径在当前仓库状态下实际可用的主要是读取写回能力以文档支持矩阵为准、可能仍在迭代中。更完整的 Iceberg、Delta Lake 独立连接器用法可分别参考 Iceberg 连接器文档 与 Delta Lake 连接器文档。类型系统Glue/Hive 类型的实际转换规则Glue 的目录类型系统基于 Apache Hive 类型系统但 Data Catalog不校验写入 type 字段的取值实际解析行为取决于底层表格式。CSV 表使用 Glue CSV 分类器类型Parquet 表对应 Arrow 类型系统Iceberg / Delta Lake 则各有各的类型体系。对 CSV/Parquet 这类直接走 GlueColumns元数据的路径Daft 使用 _convert_glue_type 做显式映射支持的类型与转换结果如下该映射表有完整的单元测试 test_convert_glue_schema 佐证Glue/Hive 类型Daft DataTypebooleanboolbyteint8shortint16integerint32long/bigintint64floatfloat32doublefloat64decimaldecimal128(precision38, scale18)stringstringtimestamptimestamp(us, UTC)datedate其他类型包括list、map、struct等复杂类型目前会抛出Unsupported Glue type的ValueError——源码中的 TODO 注释表明这是有意留待扩展的边界。这意味着如果你的 Glue 表 schema 中包含复杂类型get_table阶段就会失败需要先行简化 schema 或等待后续版本支持。扩展自定义表格式注册你自己的 GlueTable对于内置四类实现无法覆盖的表格式可以自定义GlueTable子类并注册到 Catalog。需要提醒的是这不是稳定 API源码与文档均标注其为补丁式技术。机制回顾实现GlueTable抽象类核心是实现类方法from_table_info(cls, catalog, table)——接收 GlueGetTable返回的表元数据字典元数据匹配时返回实例不匹配时抛ValueError这是get_table逐个尝试的约定同时实现read、append、overwrite通过catalog._table_impls.append(YourTable)追加注册注意是追加不会覆盖内置的四种实现。最小示例from typing import Any, Literal from daft.catalog import Catalog from daft.catalog.__glue import GlueCatalog, GlueTable, load_glue from daft.dataframe import DataFrame class GlueTestTable(GlueTable): GlueTestTable 演示如何注册自定义表实现。 classmethod def from_table_info(cls, catalog: GlueCatalog, table: dict[str, Any]) - GlueTable: if bool(table[Parameters].get(pytest)): return cls(catalog, table) raise ValueError(Expected Parameter pytestTrue) def read(self, **options) - DataFrame: raise NotImplementedError def append(self, df: DataFrame, **options) - None: raise NotImplementedError def overwrite(self, df: DataFrame, **options) - None: raise NotImplementedError gc load_glue(my_glue_catalog, region_nameus-west-2) gc._table_impls.append(GlueTestTable) # 注册自定义表实现from_table_info中判断表是否属于自己的常用字段是table[Parameters]classification、table_type等以及StorageDescriptorLocation、Columns。同样的模式可以在 tests/catalog/test_glue.py 中的 GlueTestTable 找到完整参照。验证与测试从本地 mock 到真实 AWS 环境这套集成的行为有两层测试保障阅读它们是确认边界的最快途径本地单测tests/catalog/test_glue.py使用 moto 的mock_aws模拟 Glue 服务端覆盖三种构造方式load_glue全参数、from_client、from_session含 boto3/botocore 两种会话、namespace/表的生命周期与错误路径、四种内置表的元数据匹配与拒绝逻辑例如把classification: parquet的元数据喂给GlueCsvTable会被正确拒绝以及上文类型转换表的逐类型断言。无需真实 AWS 凭证即可在本地运行。集成测试tests/integration/iceberg/test_glue_iceberg.py标注integration且需要--credentials标志运行验证真实 AWS 上的 Glue Iceberg 表读写闭环——包括GlueCatalog.from_session建连、table.read()读取并交给daft.sql做 join、以及 pandas → Daft →table.write(df, modeoverwrite)写回后再读回校验一致性。当前局限与注意事项小结该 API 处于早期开发阶段官方文档开头即有此警示接口可能随版本变化CSV/Parquet 仅支持读且不支持 Hive 风格分区Iceberg 读/写Delta 读取已实现写回路径在源码中尚未落地不支持创建 Glue 表、不支持函数注册create_function抛出NotImplementedErrorGlue/Hive 复杂类型list/map/struct暂未纳入类型映射list_tables必须带 namespacelist_namespaces不支持 pattern 过滤自定义表格式注册依赖内部字段_table_impls不属于稳定契约。整体而言Daft 的 Glue 集成遵循元数据解析交给目录、数据读写复用既有连接器的设计Glue 元数据被转换为GlueTable实例后CSV/Parquet 走daft.io._csv/daft.io._parquetIceberg 走read_iceberg/write_icebergDelta 走read_deltalake从而让 Glue 成为 Daft DataFrame 生态中一个透明的表来源与写入目标。【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表