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 与 Table
Daft 与 AWS Glue 的对接通过 Daft Catalog 接口完成。理解接入方式前,必须先记住两套术语之间的映射关系,这是使用 Glue 集成时最容易混淆的地方:
| Daft 概念 | Glue 概念 |
|---|---|
| Catalog | Data Catalog(服务级数据目录) |
| Namespace | Database |
| Table | Table |
也就是说,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 实现):
| 参数 | 类型 | 说明 |
|---|---|---|
name | str | Catalog 名称,建议全小写以兼容 Hive 命名习惯 |
region_name | str(可选) | 要连接的 AWS 区域 |
api_version | str(可选) | Glue 服务使用的 API 版本 |
use_ssl | bool(可选) | 是否使用 SSL 连接 |
verify | bool \| str(可选) | 是否校验 SSL 证书,或指向 CA bundle 的路径 |
endpoint_url | str(可选) | 替代端点地址(可指向本地兼容端点) |
aws_access_key_id | str(可选) | 认证用的 Access Key ID |
aws_secret_access_key | str(可选) | 认证用的 Secret Key |
aws_session_token | str(可选) | 临时凭证的会话 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_name="us-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", session=sess)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( name="my_glue_catalog", region_name="us-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")会拆分为DatabaseName+Expression传给 Glue 的GetTables,实现前缀过滤(GlueCatalog._list_tables)。注意:Glue 端不支持省略 namespace 的全局列表,list_tables(None)会抛出GlueCatalog requires the pattern to contain a namespace的ValueError,且结果通过NextToken自动翻页取全量。 - 列出 namespace:
catalog.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_table:
df = catalog.read_table("my_namespace.my_table") catalog.write_table("my_namespace.my_table", df, mode="append") # 或 mode="overwrite"write_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解析 schema;skip.header.line.count == "1"时视为有表头,delimiter参数控制分隔符(默认逗号)。读取时调用daft.io._csv.read_csv并强制infer_schema=False,即完全使用 Glue 元数据中的 schema,不做数据侧推断(GlueCsvTable.read)。 - GlueParquetTable:要求
classification == "parquet",读取时同样infer_schema=False,直接以 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, mode="append"|"overwrite")完成(GlueIcebergTable.read)。 - GlueDeltaTable:要求
table_type == "delta",内部构造一个指向 S3 Location 的UnityCatalogTable,读取走daft.io.delta_lake._deltalake.read_deltalake。需要提示的是,从当前源码结构看,GlueDeltaTable.append/overwrite目前仍会抛出NotImplementedError(GlueDeltaTable),即 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 DataType |
|---|---|
boolean | bool |
byte | int8 |
short | int16 |
integer | int32 |
long/bigint | int64 |
float | float32 |
double | float64 |
decimal | decimal128(precision=38, scale=18) |
string | string |
timestamp | timestamp(us, UTC) |
date | date |
其他类型(包括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 pytest='True'") 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_name="us-west-2") gc._table_impls.append(GlueTestTable) # 注册自定义表实现from_table_info中判断表是否属于自己的常用字段是table["Parameters"](classification、table_type等)以及StorageDescriptor(Location、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, mode="overwrite")写回后再读回校验一致性。
当前局限与注意事项小结
- 该 API 处于早期开发阶段(官方文档开头即有此警示),接口可能随版本变化;
- CSV/Parquet 仅支持读,且不支持 Hive 风格分区;Iceberg 读/写;Delta 读取已实现,写回路径在源码中尚未落地;
- 不支持创建 Glue 表、不支持函数注册(
create_function抛出NotImplementedError); - Glue/Hive 复杂类型(list/map/struct)暂未纳入类型映射;
list_tables必须带 namespace,list_namespaces不支持 pattern 过滤;- 自定义表格式注册依赖内部字段
_table_impls,不属于稳定契约。
整体而言,Daft 的 Glue 集成遵循"元数据解析交给目录、数据读写复用既有连接器"的设计:Glue 元数据被转换为GlueTable实例后,CSV/Parquet 走daft.io._csv/daft.io._parquet,Iceberg 走read_iceberg/write_iceberg,Delta 走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),仅供参考