news 2026/9/20 17:29:48

Daft 数据连接器全景指南:对象存储、表格式、Catalog 与多模态数据源

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Daft 数据连接器全景指南:对象存储、表格式、Catalog 与多模态数据源

Daft 数据连接器全景指南:对象存储、表格式、Catalog 与多模态数据源

【免费下载链接】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 是一款面向 AI 与多模态负载的高性能数据引擎,本文基于仓库中的 Data Connectors 文档,系统梳理其完整的连接器生态:从 S3、GCS 等对象存储,到 Iceberg、Delta Lake、Hudi、Paimon、Lance 等开放表格式,再到 Catalog/Table 高级接口、SQL 数据库、文件读取与 Kafka 等流式数据源。读完本文,你将掌握如何在 Daft 中为任意数据形态选择正确的读/写入口,并学会用Catalog抽象统一管理跨系统数据。

连接器总览

Daft 的 I/O 能力统一收敛在daft.io模块之下,所有读函数以daft.read_*、写函数以DataFrame.write_*的形式暴露,完整清单见 Daft Connectors API docs。连接器可划分为六大类:

类别典型入口说明
对象存储s3://gs://az://cos://通过 URL 协议原生读写云存储
开放表格式read_icebergread_deltalake支持 ACID 表格式的读写
CatalogsCatalog.from_*数据治理与元数据抽象层
数据库read_sqlwrite_sqlSQL 数据库与数据仓库
文件read_parquetfrom_files各类文件格式的批量读取
其他来源read_kafkaread_mcapread_huggingface流式与多模态数据源

对象存储:原生 URL 协议

Daft 原生支持主流云对象存储,通过 URL 协议直接识别数据位置,无需额外配置即可读写:

ProviderURL Protocols配置类
AWS S3s3://S3Config
Azure Blob Storageaz://abfs://AzureConfig
Google Cloud Storagegs://gcs://GCSConfig
Tencent Cloud COScos://cosn://CosConfig
GooseFSgoosefs://GooseFSConfig

以 S3 为例,数据的 URL 形式为s3://{BUCKET}/{OBJECT_KEY},其中 Bucket 是顶层命名空间,Object Key 是桶内的唯一标识。

凭证配置的三种方式

方式一:依赖环境自动发现。配置 AWS CLI 后 Daft 会自动发现凭证,也可通过环境变量AWS_ACCESS_KEY_IDAWS_SECRET_ACCESS_KEYAWS_SESSION_TOKEN指定。注意:在 Ray 等分布式环境中,Daft 会从各 worker 机器读取凭证,因此每台 worker 都需要正确配置;若希望统一使用 driver 侧凭证,则需要手动指定。

方式二:手动传入S3Config通过daft.set_planning_configIOConfig设为全局默认配置:

from daft.io import IOConfig, S3Config # key_id 为 AWS Access Key ID,access_key 为 AWS Secret Access Key io_config = IOConfig(s3=S3Config(key_id="key_id", session_token="session_token", access_key="access_key")) # 全局设置默认 IOConfig,作用于后续所有 I/O 调用 daft.set_planning_config(default_io_config=io_config) df = daft.read_parquet("s3://my_bucket/my_path/**/*")

方式三:按操作覆盖。每次调用可传入独立的io_config=关键字参数,灵活地为不同调用使用不同的S3Config

df2 = daft.read_csv("s3://my_bucket/my_other_path/**/*", io_config=io_config)

开放表格式:ACID 数据湖读写

Daft 覆盖了业界主流开放表格式的读写,统一暴露为daft.io.read_*DataFrame.write_*接口:

表格式读函数写函数详细文档
Apache Hudiread_hudiHudi
Apache Icebergread_icebergwrite_icebergIceberg
Delta Lakeread_deltalakewrite_deltalakeDelta Lake
Apache Paimonread_paimonwrite_paimonPaimon
Lanceread_lancewrite_lanceLance

从源码结构看,表格式的读取实现分布在daft/io/下的_iceberg.py_delta_lake.py_paimon.py_lance.py等模块,写接口则统一挂在 DataFrame 上,与daft.ioAPI 文档(docs/api/io.md)中的列出的 Input/Output 清单一一对应。

Catalogs 与 Tables:数据治理的高级抽象

Catalog(目录)是组织、治理数据的中心化服务,负责创建表与命名空间、管理事务和访问控制。其核心价值在于屏蔽物理存储细节——你只需关注数据的逻辑结构,无需关心文件格式、分区方案或存储位置。

⚠️ 官方提示:Catalog API 仍处于早期开发阶段,欢迎通过 issue 反馈需求。

Daft 通过CatalogTable两个接口整合了以下目录实现,同时复用既有daft.read_*df.write_*能力支持 Iceberg、Delta Lake 等开放表格式:

Catalog说明
AWS GlueAWS Glue Data Catalog
AWS S3 TablesAmazon S3 Tables catalog
Apache GravitinoApache Gravitino metadata lake
Unity Catalog (Databricks)Databricks Unity Catalog

Catalog 示例:从 PyIceberg 目录出发

下面的示例基于 Daft Sessions 教程创建的 Iceberg Catalog:

import daft from daft import Catalog # iceberg_catalog 来自 'Sessions' 教程 iceberg_catalog = load_catalog(...) # 从 pyiceberg catalog 实例创建 daft catalog catalog = Catalog.from_iceberg(iceberg_catalog) # 验证 catalog """ Catalog('default') """ # 以 DataFrame 形式读取表 catalog.read_table("example.tbl").schema() """ ╭─────────────┬─────────╮ │ column_name ┆ type │ ╞═════════════╪═════════╡ │ x ┆ Boolean │ ├╌╌╌╌╌╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌┤ │ y ┆ Int64 │ ├╌╌╌╌╌╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌┤ │ z ┆ Utf8 │ ╰─────────────┴─────────╯ """ # 给定一个 dataframe... df = daft.from_pylist([{ "x": False, "y": -1, "z": "xyz" }]) # 可以写入表 catalog.write_table("example.tbl", df, mode="append") # 也可以获取表实例 t = catalog.get_table("example.tbl") # 参见 'Working with Tables' t """ Table('tbl') """

Working with Catalogs:目录操作

Catalog接口支持get_tablelist_tables等目录动作。从源码(daft/catalog/init.py)可见,Catalog是一个抽象基类,定义了_create_table_drop_table_get_table_list_tables_list_namespaces等核心抽象方法,并提供一组from_*静态工厂方法:

import daft from daft import Catalog, Table # 从 pyiceberg catalog 对象创建 catalog _ = Catalog.from_iceberg(pyiceberg_catalog) # 从 unity catalog 对象创建 catalog _ = Catalog.from_unity(unity_catalog) # 各种类型都可以注册为表,注意它们是等价的 example_dict = { "x": [ 1, 2, 3 ] } example_df = daft.from_pydict(example_dict) example_table = Table.from_df("temp", example_df) # 从 pydict(名称到表的映射)创建 catalog catalog = Catalog.from_pydict( { "R": example_dict, "S": example_df, "T": example_table, } ) # 列出可用表 # 注意:pattern 语法与 catalog 实现相关。 # 原生/内存与 Postgres catalog 使用 SQL LIKE 语法(%, _), # Iceberg 等 catalog 使用前缀匹配。 catalog.list_tables(pattern=None) """ ['R', 'S', 'T'] """ # 按名称获取表 table_t = catalog.get_table("T") table_t.show() """ ╭───────╮ │ x │ │ --- │ │ Int64 │ ╞═══════╡ │ 1 │ ├╌╌╌╌╌╌╌┤ │ 2 │ ├╌╌╌╌╌╌╌┤ │ 3 │ ╰───────╯ """

更多工厂方法(来自源码 daft/catalog/init.py)

  • Catalog.from_gravitino(endpoint, metalake_name, auth_type="simple", username=None, password=None, token=None):连接 Gravitino metalake,支持simpleoauth2两种认证方式;
  • Catalog.from_s3tables(table_bucket_arn, client=None, session=None):基于 S3 Tables bucket ARN 创建,可传入 boto3 client 或 session(二者只能提供一个),否则底层使用 Iceberg REST client;
  • Catalog.from_glue(name, client=None, session=None):AWS Glue 中 Database → Daft Namespace、Table → Daft Table 的映射;
  • Catalog.from_paimon(catalog, name="paimon"):包装pypaimon的 catalog 对象;
  • Catalog.from_postgres(connection_string, extensions=("vector",)):从 PostgreSQL 连接串创建,自动执行CREATE EXTENSION IF NOT EXISTS(默认创建vector扩展)。

此外,未限定的daft.read_table("my_table")使用默认 catalog 与 namespace,支持"my_namespace.my_table""my_catalog.my_namespace.my_table"两种限定形式;临时表(如daft.create_temp_table创建的)在解析未限定名称时优先于 catalog 表。

Working with Tables:从目录到 DataFrame 的桥梁

Table接口是 catalog 与 dataframe 之间的桥梁:既能把表读成 DataFrame,也能把 DataFrame 写入表。你可以脱离 catalog 单独使用工厂方法创建 Table——底层实际复用的仍是daft.read_*df.write_*API,Table的价值在于对表格式本身提供一层间接抽象,作为 catalog 读写的统一入口

from daft import Table from pyiceberg.table import StaticTable # 假设你有一个 pyiceberg 表 pyiceberg_table = StaticTable("metadata.json") # 包装为 daft table 以使用 daft 的表 API table = Table.from_iceberg(pyiceberg_table) # 读取 DataFrame,等价于 daft.read_iceberg(pyiceberg_table) df = table.read() # 也可以从 dataframe 创建临时表 daft.create_temp_table("my_temp_table", daft.from_pydict({ ... })) # 临时表会像其他表一样被解析 df = daft.read_table("my_temp_table")

目前可以从pyicebergdaft.unity的表对象进行读取(源码见 daft/catalog/init.py 中Table抽象类与Table.from_dfTable.from_iceberg等工厂方法)。

Reference:完整 API 文档见 Catalog & Table API docs。三个核心类型:

  • Catalog:创建与访问表、命名空间的接口;
  • Identifier:对象的路径,如catalog.namespace.table
  • Table:读写 DataFrame 的接口。

数据库连接器

Daft 支持多种数据库的读写:

数据库入口说明
Google Cloud Bigtable(实验性)write_bigtable写入 Bigtable,API 可能变化,详见 Bigtable
ClickHousewrite_clickhouse写入 ClickHouse,详见 ClickHouse
PostgreSQLCatalog.from_postgres从 PostgreSQL 数据库创建 catalog,详见 Postgres
SQL 数据库read_sqlwrite_sql通用 SQL 读写,详见 SQL Databases
Turbopufferwrite_turbopuffer写入 Turbopuffer namespace,详见 Turbopuffer

SQL 读写的三大能力

daft.read_sql()为例,SQL 连接器(详见 SQL Databases)提供:

  1. 20+ SQL 方言:借助 SQLGlot 在不同方言间转换 SQL 查询;
  2. 并行与分布式读取:默认多线程 runner 使用本机所有核心,使用 distributed Ray runner 时则跨多机并行;
  3. 数据跳过优化:只读取匹配df.select()df.limit()df.where()表达式的数据,经常跳过整个分区/列。

安装 SQL 支持:

pip install -U "daft[sql]"

底层通过 ConnectorX(Rust 实现、直接读入 Arrow Table 实现零拷贝)读取;若数据库不被 ConnectorX 支持,则回退到 SQLAlchemy。也可以通过连接工厂传入 SQLAlchemy engine 附加参数:

import daft from sqlalchemy import create_engine def create_connection(): return create_engine("sqlite:///example.db", echo=True).connect() df = daft.read_sql("SELECT * FROM books", create_connection)

并行读取:传入partition_colnum_partitions,Daft 会在 SQL 查询后追加WHERE col > ... AND col <= ...子句切分数据,再并行读取每个分区:

df = daft.read_sql( "SELECT * FROM table", "sqlite:///big_table.db", partition_col="col", num_partitions=3, )

谓词下推whereselect表达式会被注入 SQL 查询本身,通过df.explain(show_all=True)可验证。例如对 BigQuery Google Trends 数据集的过滤与投影会被改写为SELECT refresh_date, term, rank ... WHERE rank = 1 AND refresh_date >= CAST('2024-04-01' AS DATE) ...。官方在 12GB RAM 的 Google Colab 上实测:开启下推时运行 8.87s、峰值内存 315.97 MiB;关闭下推时超过 2 分钟并最终 OOM。相比手写改写 SQL,Daft 对大量表达式组合的自动下推更可靠。

文件连接器

文件引用:from_files

daft.from_files()从 glob 模式创建惰性文件引用的 DataFrame——与立即加载内容的读函数不同,它生成可在需要时读取的File对象:

import daft # 本地文件 df = daft.from_files("/path/to/files/*.jpeg") # 远程文件(S3 / GCS 同理) df = daft.from_files("s3://my-bucket/images/*.png") df = daft.from_files("gs://my-bucket/images/*.png")

输出为单列 DataFrame:file列(类型File),表示可按需读取的惰性文件引用。支持标准 glob 通配符:*(任意字符)、?(单个字符)、[...](括号内字符)、**(递归匹配目录),且可传入多个模式列表:

# 递归搜索 df = daft.from_files("/images/**/*.png") # 多个模式 df = daft.from_files(["/images/*.jpeg", "/photos/*.jpeg"])

通过col("file")可以访问文件属性(file_path()file_size()),甚至在下游读取内容前先按属性过滤;未匹配到任何文件时返回空 DataFrame 而不是报错。与daft.from_glob_path()的差异:后者返回path字符串列(仅需路径时使用),前者返回File对象列(需读取内容或访问属性时使用)。图片处理管线推荐from_glob_path+.download()+decode_image()组合,详见 Files。

常见文件格式

格式读函数写函数
CSVread_csvwrite_csv
JSONread_jsonwrite_json
Parquetread_parquetwrite_parquet
Textread_text—(每行一个字符串行)
Blobread_blob—(每个文件一行原始字节)
WARCread_warc
WebDatasetread_webdataset—(多模态 TAR shards,支持惰性媒体读取,见 WebDataset)
Videoread_video_frames—(读取视频帧)

Text/Blob 的详细用法见 Text Files 与 Blob Files。所有基于文件的读取器还共享一组通用选项(压缩、schema 提示等),完整列表见 Generic File Source Options。

通用文件源选项:忽略损坏文件

ignore_corrupt_filesread_parquetread_csvread_iceberg共用的关键选项(read_jsonread_warcread_text不支持)。默认情况下读到损坏文件(损坏、截断、或在列出与读取之间被删除)会抛错并中止查询;开启后这些文件会被静默跳过,查询继续:

import daft df = daft.read_parquet("s3://my-bucket/data/**/*.parquet", ignore_corrupt_files=True) df.collect() # 程序化访问被跳过的文件:(path, reason, partial) 三元组 skipped = df.skipped_corrupt_files for path, reason, partial in skipped: tag = " (partial)" if partial else "" print(f"Skipped{tag}: {path}\n Reason: {reason}")

判定边界:格式非法(Parquet magic bytes 错误、截断的 footer、行/列数不匹配)、数据损坏(不可读的 row group、非法 CSV 编码、行内字段数错误)、文件缺失(列表与读取之间被并发 compaction 或分区覆盖删除)会被跳过;而网络错误(连接重置、读超时、限流)与权限错误(拒绝访问、凭证不足)不会被跳过——这些属于可重试的瞬时问题,静默忽略权限错误会掩盖需要人工介入的配置错误。

可观测性:Daft 会对每个被跳过的文件输出WARNING级日志(含路径与原因);df.skipped_corrupt_files.collect()后可用(.count_rows()等其他执行方法不会填充该属性),可接入告警或死信队列:

df = daft.read_parquet("s3://my-bucket/nightly/**/*.parquet", ignore_corrupt_files=True) df.write_parquet("s3://my-bucket/processed/") skipped = df.skipped_corrupt_files if skipped: for path, reason, partial in skipped: dead_letter_queue.put({"path": path, "reason": reason, "partial": partial, "run": TODAY})

注意两个限制:Iceberg 的delete files损坏不受该选项保护(删除文件损坏通常意味着更严重的 catalog 不一致);Parquet 开启该选项后 count 下推被禁用,df.count()需读取全部 row-group 数据而非走元数据优化,大数据集上可能变慢。详见 Generic File Source Options。

其他数据源

数据源入口状态与说明
Apache Kafkaread_kafka实验性:仅支持有界批量读取,无流式/无界模式,无 offset 提交管理,详见 Apache Kafka
MCAPread_mcap实验性,详见 MCAP
Hugging Face Datasetsread_huggingface读取 HF 数据集,详见 Hugging Face Datasets

自定义连接器

Daft 不支持的数据源可通过实现自定义 source/sink 接入(详见 Custom Connectors):

接口说明
DataSink从 DataFrame 写数据的接口
DataSource将数据读入 DataFrame 的接口
DataSourceTask表示可独立处理的数据分区
WriteResult由 DataSink 写入的中间结果包装
write_sink将 DataFrame 写入给定 DataSink

自定义目录则通过扩展CatalogTable接口实现,详见 Custom Catalogs。

内存数据源

不经过外部存储,直接从内存 Python 对象创建 DataFrame:

函数说明
from_arrow从 PyArrow Table 或 RecordBatch 创建
from_dask_dataframe从 Dask DataFrame 创建
from_pandas从 Pandas DataFrame 创建
from_pydict从 Python 字典创建
from_pylist从 Python 列表创建
from_ray_dataset从 Ray Dataset 创建

小结

Daft 的连接器体系以「URL 协议识别对象存储 →read_*/write_*统一读写表格式与数据库 →Catalog/Table抽象数据治理 →DataSource/DataSink扩展自定义系统 →from_*接入内存数据」为完整闭环。实际选型建议:云上数据直接用s3://等协议路径;需要 ACID 与元数据管理时选 Iceberg/Delta Lake 等表格式并结合对应Catalog;批量文件处理优先read_parquet/read_csv并善用ignore_corrupt_files;异构 SQL 数据源统一走read_sql的方言转换与并行分区读取。更多 API 细节可查阅 Daft Connectors API docs 与 Catalog & Table API docs。

【免费下载链接】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/19 5:07:13

Yue2模型实操指南:AR-NAR混合Transformer部署教程

1. 项目概述&#xff1a;一个被误读的“YuE”——从热搜词迷雾中打捞真实技术信号最近在多个技术社区和搜索平台看到“YuE”“YuE2”频繁出现在Python、Hugging Face相关话题的热搜榜前列&#xff0c;甚至和AR–NAR Mixture-of-Transformers、FontDiffuser、TEI&#xff08;Tex…

作者头像 李华
网站建设 2026/9/20 5:10:26

PSO算法优化光伏MPPT:突破局部遮阴挑战

1. 新能源控制器中的MPPT挑战与机遇光伏发电系统在实际运行中常常面临局部遮阴的困扰&#xff0c;就像一片树林中总有几棵树会投下阴影。当光伏阵列部分区域被遮挡时&#xff0c;其功率-电压(P-V)特性曲线会出现多个峰值点&#xff0c;这给传统的最大功率点跟踪(MPPT)算法带来了…

作者头像 李华
网站建设 2026/9/20 7:16:02

AI时代企业核心生产力重构与开源策略升级

1. 项目概述&#xff1a;AI时代的生产力边界重构最近三年&#xff0c;我观察到企业技术战略正在经历一场静默革命。传统的外包与开源边界在生成式AI冲击下变得模糊不清——某电商平台将80%的客服代码交给AI生成&#xff0c;某金融公司用开源大模型替代了原本外包的智能投顾模块…

作者头像 李华
网站建设 2026/9/20 9:35:18

Notepad-- 文件对比功能实战:多版本差异比对的完整指南

Notepad-- 文件对比功能实战&#xff1a;多版本差异比对的完整指南 【免费下载链接】notepad-- 一个支持windows/linux/mac的文本编辑器&#xff0c;目标是做中国人自己的编辑器&#xff0c;来自中国。 项目地址: https://gitcode.com/GitHub_Trending/no/notepad-- 备份…

作者头像 李华