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_iceberg、read_deltalake等 | 支持 ACID 表格式的读写 |
| Catalogs | Catalog.from_* | 数据治理与元数据抽象层 |
| 数据库 | read_sql、write_sql等 | SQL 数据库与数据仓库 |
| 文件 | read_parquet、from_files等 | 各类文件格式的批量读取 |
| 其他来源 | read_kafka、read_mcap、read_huggingface | 流式与多模态数据源 |
对象存储:原生 URL 协议
Daft 原生支持主流云对象存储,通过 URL 协议直接识别数据位置,无需额外配置即可读写:
| Provider | URL Protocols | 配置类 |
|---|---|---|
| AWS S3 | s3:// | S3Config |
| Azure Blob Storage | az://、abfs:// | AzureConfig |
| Google Cloud Storage | gs://、gcs:// | GCSConfig |
| Tencent Cloud COS | cos://、cosn:// | CosConfig |
| GooseFS | goosefs:// | GooseFSConfig |
以 S3 为例,数据的 URL 形式为s3://{BUCKET}/{OBJECT_KEY},其中 Bucket 是顶层命名空间,Object Key 是桶内的唯一标识。
凭证配置的三种方式
方式一:依赖环境自动发现。配置 AWS CLI 后 Daft 会自动发现凭证,也可通过环境变量AWS_ACCESS_KEY_ID、AWS_SECRET_ACCESS_KEY、AWS_SESSION_TOKEN指定。注意:在 Ray 等分布式环境中,Daft 会从各 worker 机器读取凭证,因此每台 worker 都需要正确配置;若希望统一使用 driver 侧凭证,则需要手动指定。
方式二:手动传入S3Config。通过daft.set_planning_config将IOConfig设为全局默认配置:
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 Hudi | read_hudi | — | Hudi |
| Apache Iceberg | read_iceberg | write_iceberg | Iceberg |
| Delta Lake | read_deltalake | write_deltalake | Delta Lake |
| Apache Paimon | read_paimon | write_paimon | Paimon |
| Lance | read_lance | write_lance | Lance |
从源码结构看,表格式的读取实现分布在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 通过Catalog与Table两个接口整合了以下目录实现,同时复用既有daft.read_*与df.write_*能力支持 Iceberg、Delta Lake 等开放表格式:
| Catalog | 说明 |
|---|---|
| AWS Glue | AWS Glue Data Catalog |
| AWS S3 Tables | Amazon S3 Tables catalog |
| Apache Gravitino | Apache 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_table、list_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,支持simple与oauth2两种认证方式;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")目前可以从pyiceberg与daft.unity的表对象进行读取(源码见 daft/catalog/init.py 中Table抽象类与Table.from_df、Table.from_iceberg等工厂方法)。
Reference:完整 API 文档见 Catalog & Table API docs。三个核心类型:
Catalog:创建与访问表、命名空间的接口;Identifier:对象的路径,如catalog.namespace.table;Table:读写 DataFrame 的接口。
数据库连接器
Daft 支持多种数据库的读写:
| 数据库 | 入口 | 说明 |
|---|---|---|
| Google Cloud Bigtable(实验性) | write_bigtable | 写入 Bigtable,API 可能变化,详见 Bigtable |
| ClickHouse | write_clickhouse | 写入 ClickHouse,详见 ClickHouse |
| PostgreSQL | Catalog.from_postgres | 从 PostgreSQL 数据库创建 catalog,详见 Postgres |
| SQL 数据库 | read_sql、write_sql | 通用 SQL 读写,详见 SQL Databases |
| Turbopuffer | write_turbopuffer | 写入 Turbopuffer namespace,详见 Turbopuffer |
SQL 读写的三大能力
以daft.read_sql()为例,SQL 连接器(详见 SQL Databases)提供:
- 20+ SQL 方言:借助 SQLGlot 在不同方言间转换 SQL 查询;
- 并行与分布式读取:默认多线程 runner 使用本机所有核心,使用 distributed Ray runner 时则跨多机并行;
- 数据跳过优化:只读取匹配
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_col与num_partitions,Daft 会在 SQL 查询后追加WHERE col > ... AND col <= ...子句切分数据,再并行读取每个分区:
df = daft.read_sql( "SELECT * FROM table", "sqlite:///big_table.db", partition_col="col", num_partitions=3, )谓词下推:where、select表达式会被注入 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。
常见文件格式
| 格式 | 读函数 | 写函数 |
|---|---|---|
| CSV | read_csv | write_csv |
| JSON | read_json | write_json |
| Parquet | read_parquet | write_parquet |
| Text | read_text | —(每行一个字符串行) |
| Blob | read_blob | —(每个文件一行原始字节) |
| WARC | read_warc | — |
| WebDataset | read_webdataset | —(多模态 TAR shards,支持惰性媒体读取,见 WebDataset) |
| Video | read_video_frames | —(读取视频帧) |
Text/Blob 的详细用法见 Text Files 与 Blob Files。所有基于文件的读取器还共享一组通用选项(压缩、schema 提示等),完整列表见 Generic File Source Options。
通用文件源选项:忽略损坏文件
ignore_corrupt_files是read_parquet、read_csv、read_iceberg共用的关键选项(read_json、read_warc、read_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 Kafka | read_kafka | 实验性:仅支持有界批量读取,无流式/无界模式,无 offset 提交管理,详见 Apache Kafka |
| MCAP | read_mcap | 实验性,详见 MCAP |
| Hugging Face Datasets | read_huggingface | 读取 HF 数据集,详见 Hugging Face Datasets |
自定义连接器
Daft 不支持的数据源可通过实现自定义 source/sink 接入(详见 Custom Connectors):
| 接口 | 说明 |
|---|---|
DataSink | 从 DataFrame 写数据的接口 |
DataSource | 将数据读入 DataFrame 的接口 |
DataSourceTask | 表示可独立处理的数据分区 |
WriteResult | 由 DataSink 写入的中间结果包装 |
write_sink | 将 DataFrame 写入给定 DataSink |
自定义目录则通过扩展Catalog与Table接口实现,详见 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),仅供参考