news 2026/9/17 8:03:31

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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 与 Table

Daft 与 AWS Glue 的对接通过 Daft Catalog 接口完成。理解接入方式前,必须先记住两套术语之间的映射关系,这是使用 Glue 集成时最容易混淆的地方:

Daft 概念Glue 概念
CatalogData Catalog(服务级数据目录)
NamespaceDatabase
TableTable

也就是说,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 ID
aws_secret_access_keystr(可选)认证用的 Secret Key
aws_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_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_gluefrom_clientfrom_session构造。

方式三:公开入口 Catalog.from_glue

面向一般用户更推荐走 Catalog.from_glue 这个工厂方法,它要求必须提供clientsession二者之一(同时提供会报错,都不提供则抛出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以“让位”给下一个实现;若所有实现都不匹配,则抛出包含classificationtable_type信息的ValueError。表不存在时抛出统一的NotFoundError

默认的_table_impls列表在 GlueCatalog.new中初始化为GlueCsvTableGlueParquetTableGlueIcebergTableGlueDeltaTable四类内置实现。

围绕 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 namespaceValueError,且结果通过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_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_idbranchtagignore_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
booleanbool
byteint8
shortint16
integerint32
long/bigintint64
floatfloat32
doublefloat64
decimaldecimal128(precision=38, scale=18)
stringstring
timestamptimestamp(us, UTC)
datedate

其他类型(包括listmapstruct等复杂类型)目前会抛出Unsupported Glue typeValueError——源码中的 TODO 注释表明这是有意留待扩展的边界。这意味着如果你的 Glue 表 schema 中包含复杂类型,get_table阶段就会失败,需要先行简化 schema 或等待后续版本支持。

扩展自定义表格式:注册你自己的 GlueTable

对于内置四类实现无法覆盖的表格式,可以自定义GlueTable子类并注册到 Catalog。需要提醒的是:这不是稳定 API,源码与文档均标注其为"补丁式技术"

机制回顾:

  1. 实现GlueTable抽象类,核心是实现类方法from_table_info(cls, catalog, table)——接收 GlueGetTable返回的表元数据字典,元数据匹配时返回实例,不匹配时抛ValueError(这是get_table逐个尝试的约定);同时实现readappendoverwrite
  2. 通过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"]classificationtable_type等)以及StorageDescriptorLocationColumns)。同样的模式可以在 tests/catalog/test_glue.py 中的 GlueTestTable 找到完整参照。

验证与测试:从本地 mock 到真实 AWS 环境

这套集成的行为有两层测试保障,阅读它们是确认边界的最快途径:

  • 本地单测(tests/catalog/test_glue.py):使用 moto 的mock_aws模拟 Glue 服务端,覆盖三种构造方式(load_glue全参数、from_clientfrom_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 8:02:35

微信占用60G?用‘降级’法彻底清理到4.8G的实战指南

打开手机存储空间的时候我愣了一下&#xff1a;微信&#xff0c;60.3G。一部手机总共才256G&#xff0c;一个微信就吃掉了将近四分之一。最离谱的是&#xff0c;这还不是我一个人的问题——群里随手一问&#xff0c;七八个朋友晒出来的截图都在40G到80G之间&#xff0c;有个老哥…

作者头像 李华
网站建设 2026/9/17 8:02:22

Vivado/Vitis 2024.2升级2024.2.1安装器找不到现有安装的解决方法

Vivado/Vitis 2024.2 升级 2024.2.1&#xff1a;安装器找不到现有安装的原因与完整解决办法搞FPGA的兄弟应该都懂&#xff0c;Vivado和Vitis这套工具链的安装和升级&#xff0c;一直是让人又爱又恨的环节。好不容易把 2024.2 的大版本用顺手了&#xff0c;结果 2024.2.1 的更新…

作者头像 李华
网站建设 2026/9/17 7:59:37

Linux下OpenClaw与QQ机器人集成开发指南

1. Linux 环境下 OpenClaw 与 QQ 机器人集成指南在当今自动化与智能化技术快速发展的背景下&#xff0c;将 AI 能力接入即时通讯平台已成为提升工作效率和用户体验的重要方式。作为一名长期从事 Linux 系统管理和 AI 应用开发的工程师&#xff0c;我将分享如何在 Debian 12 系统…

作者头像 李华