Daft I/O 全指南:从内存、文件、数据湖到数据目录的读写 API 详解
【免费下载链接】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 与多模态负载的高性能数据引擎,其 I/O 层覆盖了从内存对象、本地/云上文件、开放表格式(Iceberg/Delta Lake/Hudi/Lance)到数据库与外部集成(Kafka、SQL、Hugging Face、WebDataset)的全链路读写能力。本文以 Daft 官方 API 文档 docs/api/io.md 为骨架,结合仓库源码逐层拆解from_*/read_*输入家族、write_*输出家族、用户自定义DataSource/DataSink扩展点以及谓词/投影/limit 下推机制,帮助读者掌握在不同数据源之间高效搬运数据的完整实战方案。
Daft I/O API 全景
Daft 创建 DataFrame 的方式分为两大类:
- 输入(Input):将内存数据、文件、数据目录与外部集成转换为 DataFrame,包括
from_*系列(daft/convert.py 与 daft/io/file_path.py)和read_*系列(daft/io/init.py 统一导出); - 输出(Output):将 DataFrame 写回文件、表格式或外部服务,即
DataFrame.write_*系列(实现于 daft/dataframe/dataframe.py)。
除此之外,Daft 还暴露了两层高级扩展 API:用户自定义 I/O(daft/io/source.py 与 daft/io/sink.py,标记为 experimental)以及Pushdowns 下推机制(daft/io/pushdowns.py),后者是 Daft 在扫描阶段省去不必要 IO 的核心手段。
与各种 Connector 的更多用法可进一步参考 docs/connectors/index.md(对象存储、开放表格式、数据库、文件与目录等)。
输入:从内存数据构建 DataFrame
from_*系列 API 用于把已经存在于进程内存中的数据对象转换为 Daft DataFrame,是快速上手与单元测试的常用入口,全部定义在 daft/convert.py。
from_pydict 与 from_pylist
daft.from_pydict(data: dict[str, InputListType]):以"列名 → 序列"的字典构建 DataFrame。每个 value 必须是等长的 Python list、NumPy array 或 PyArrow array(daft/convert.py)。daft.from_pylist(data: list[dict[str, Any]]):以"行"为单位(列表中的每个 dict 是一行,key 为列名)构建 DataFrame(daft/convert.py)。
import daft # 列式构建 df = daft.from_pydict({"foo": [1, 2]}) df.show() # ╭───────╮ # │ foo │ # │ --- │ # │ Int64 │ # ╞═══════╡ # │ 1 │ # │ 2 │ # ╰───────╯ # 行式构建 df2 = daft.from_pylist([{"foo": 1}, {"foo": 2}])from_pandas 与 from_arrow
daft.from_pandas(data):接受单个 pandas DataFrame 或 pandas DataFrame 列表(daft/convert.py)。daft.from_arrow(data):接受 pyarrow Table、Table 列表/迭代器,或任何实现了 Arrow PyCapsule 接口(即拥有__arrow_c_stream__方法)的对象,例如 pyarrow RecordBatchReader、pandas 2.2+ 的 DataFrame、nanoarrow 数组等(daft/convert.py)。源码注释表明:非pa.Table的 Arrow 流对象走DataFrame._from_arrow_stream路径,而pa.Table优先走 pyarrow 感知路径,以支持扩展类型、Decimal256 等 Rust FFI 流无法表达的类型。
import pyarrow as pa import pandas as pd t = pa.table({"a": [1, 2, 3], "b": ["foo", "bar", "baz"]}) df = daft.from_arrow(t) pd_df = pd.DataFrame({"a": [1, 2, 3], "b": ["foo", "bar", "baz"]}) df = daft.from_pandas(pd_df)from_ray_dataset 与 from_dask_dataframe
daft.from_ray_dataset(ds):从 Ray Dataset 构建 DataFrame(daft/convert.py);daft.from_dask_dataframe(ddf):从 Dask DataFrame 构建,要求该 Dask DataFrame 基于 Dask-on-Ray 创建(daft/convert.py)。
两者均要求 Daft 运行在 RayRunner 之上,即先调用daft.set_runner_ray()再使用:
import ray import daft daft.set_runner_ray() ds = ray.data.from_items([{"a": 1, "b": "foo"}, {"a": 2, "b": "bar"}]) df = daft.from_ray_dataset(ds)from_glob_path:把文件列表变成 DataFrame
daft.from_glob_path(path, io_config=None)根据 glob 模式返回一个包含文件元数据的 DataFrame(daft/io/file_path.py),支持四种通配符:
*:匹配任意多个字符(含 0 个);?:匹配任意单个字符;[...]:匹配方括号内的任意单个字符;**:递归匹配任意多层目录。
返回的 DataFrame 包含三列:path(文件/目录路径)、size(字节大小)、rows(Parquet 对象的行数,其他格式为 None)。当 glob 无匹配时,返回空 DataFrame 而非报错;多个 glob 模式可用列表传入。
df = daft.from_glob_path("/path/to/files/*.jpeg") df = daft.from_glob_path(["/path/to/files/*.jpeg", "/path/to/others/*.jpeg"])from_glob_path也是read_blob的底层基础,其通过LogicalPlanBuilder.from_glob_scan构建扫描计划(见 daft/io/file_path.py)。
输入:读取文件与对象存储
read_*家族负责把磁盘或远端对象存储(s3://、gs://等)上的数据读入 DataFrame。它们几乎都支持通配符路径与目录路径,并统一接受io_config用于配置云存储凭据(如IOConfig(s3=S3Config(region="us-west-2", anonymous=True)))。
read_parquet
daft.read_parquet(path, row_groups=None, infer_schema=True, schema=None, io_config=None, file_path_column=None, hive_partitioning=False, coerce_int96_timestamp_unit=None, ignore_corrupt_files=False, checkpoint=None)(daft/io/_parquet.py):
row_groups:按文件指定要读取的行组列表;仅在读取多个非通配文件时支持,且列表长度必须与 path 数量一致,否则抛ValueError;infer_schema=False时必须同时提供schema,否则报错;schema:当infer_schema=True时作为"schema 提示",用于覆盖推断列的类型或追加推断未发现的列;file_path_column:将源文件路径作为指定名称的列注入结果;hive_partitioning:从文件路径推断 Hive 风格分区并作为列;coerce_int96_timestamp_unit:将 Int96 时间戳统一到ns/us/ms精度;ignore_corrupt_files:静默跳过损坏文件(仅忽略真正的格式错误,如坏魔数、截断 footer、损坏的 row-group 数据;网络与权限错误仍会抛出),被跳过的文件记录在df.skipped_corrupt_files中;checkpoint:传入daft.CheckpointConfig实现跨运行断点续跑(要求 RayRunner)。
df = daft.read_parquet("/path/to/files-*.parquet") from daft.io import S3Config, IOConfig io_config = IOConfig(s3=S3Config(region="us-west-2", anonymous=True)) df = daft.read_parquet("s3://path/to/files-*.parquet", io_config=io_config)实现上,read_parquet将参数打包为ParquetSourceConfig与StorageConfig,再通过get_tabular_files_scan(daft/io/common.py)构造TabularFilesScan逻辑计划。值得注意的细节:在 Ray runner 下默认关闭多线程 IO(multithreaded_io = runners.get_or_create_runner().name != "ray"),以减少每个 Ray worker 的线程池与连接数争用。
read_csv
daft.read_csv(path, infer_schema=True, schema=None, has_headers=True, delimiter=None, double_quote=True, quote=None, escape_char=None, comment=None, allow_variable_columns=False, io_config=None, file_path_column=None, hive_partitioning=False, ignore_corrupt_files=False, checkpoint=None)(daft/io/_csv.py):
delimiter:字段分隔符,默认,;double_quote:是否支持双引号转义,默认 True;quote:包裹含分隔符字段的引号字符,默认";escape_char:转义字符;comment:注释行起始字符,None 表示不支持注释;allow_variable_columns:允许行间列数不一致,True 时少列补 Null、多列忽略多余列;- 其余参数语义与
read_parquet一致(io_config、file_path_column、hive_partitioning、ignore_corrupt_files、checkpoint)。
df = daft.read_csv("/path/to/files-*.csv") df = daft.read_csv("s3://path/to/files-*.csv", io_config=io_config)read_json
daft.read_json(path, infer_schema=True, schema=None, io_config=None, file_path_column=None, hive_partitioning=False, skip_empty_files=False, checkpoint=None)(daft/io/_json.py)读取**行分隔 JSON(JSONL)**文件;skip_empty_files=True可跳过空文件。
read_blob:以原始字节读取任意文件
daft.read_blob(path, *, max_connections=32, on_error="raise", io_config=None)(daft/io/_blob.py)把每个文件读成一行原始字节,语义类似 DuckDB 的read_blob,非常适合图片、音频等非表格化二进制数据。返回三列:path、size(字节数)、content(原始字节)。on_error可选"raise"(立即报错)或"null"(记录错误并回退为 Null)。
其实现是先from_glob_path列出文件,再用daft.functions.url.download表达式下载内容并alias("content")——这正是from_glob_path与表达式系统组合的典型示例。
read_video_frames:视频帧流式读取
daft.read_video_frames(path, image_height, image_width, is_key_frame=None, *, sample_interval_seconds=None, io_config=None)(daft/io/av/init.py)将视频流式读取为 DataFrame of images,需要 PyAV(pip install av)。输出字段包括:path、frame_index、frame_time(秒)、frame_time_base、frame_pts、frame_dts、frame_duration、is_key_frame。
is_key_frame:True 只取关键帧、False 只取非关键帧、None 取全部;sample_interval_seconds:按帧时间近似采样(取时间戳 ≥ 目标时间点的首帧),无有效时间戳的帧被跳过。
df = daft.read_video_frames("/path/to/file.mp4", image_height=480, image_width=640) # 约每秒一帧 df = daft.read_video_frames("/path/to/file.mp4", image_height=480, image_width=640, sample_interval_seconds=1.0)read_warc 与 read_webdataset:网页抓取与多模态数据集
daft.read_warc(path, io_config=None, file_path_column=None, checkpoint=None)(daft/io/_warc.py)读取 WARC 或 gzip 压缩的 WARC 文件(实验特性),返回 DataFrame 含强制元数据列(WARC-Record-ID、WARC-Type、WARC-Date、Content-Length)与可选字段;daft.read_webdataset(path, io_config=None, batch_size=1000)(daft/io/webdataset/_webdataset.py)读取 WebDataset TAR shards:文件名前缀相同的连续 TAR 成员被合并为一行,成员后缀成为列名(WebDataset 约定)。图像、音频、视频等二进制成员以惰性daft.File引用呈现,JSON/text/class 侧车文件则立即解码。返回列含__key__、__url__及各成员后缀列。
WebDataset 注意事项(源码 docstring 明确):不支持压缩 TAR(无法做惰性 range 引用)、不支持稀疏 TAR 成员;schema 从前五个样本推断,跨 shard 不一致会报错而非丢数据。其实现即WebDatasetSource(DataSource)类,实现了get_tasks并按pushdowns.limit收缩batch_size(见 daft/io/webdataset/_webdataset.py)。
read_huggingface
daft.read_huggingface(repo, io_config=None, format=None)(daft/io/huggingface/init.py)读取 Hugging Face 数据集,repo形如username/dataset_name:
format=None(默认)或"parquet":走快速路径read_parquet(f"hf://datasets/{repo}");若 parquet 文件不存在(glob 无匹配)或返回 400(parquet 尚未生成),则回退到datasets库;format="webdataset":等价于read_webdataset(f"hf://datasets/{repo}/**/*.tar")。
read_kafka:流式数据源
daft.read_kafka(topics, bootstrap_servers, group_id, start="earliest", end="latest", ...)基于 librdkafka 风格配置(daft/io/_kafka.py)。源码显示其start/end边界支持:"earliest"/"latest"、整数时间戳毫秒、datetime/ISO 字符串,以及分区偏移映射{partition: offset}或{topic: {partition: offset}}(多 topic 场景)。kafka_client_config可传入额外客户端配置,但bootstrap.servers与group.id受保护,不可被覆盖。
read_sql:从数据库执行查询
daft.read_sql(sql, conn, partition_col=None, num_partitions=None, partition_bound_strategy="min-max", disable_pushdowns_to_sql=False, infer_schema=True, infer_schema_length=10, schema=None)(daft/io/_sql.py):
conn:SQLAlchemy 连接工厂(Callable[[], Connection])或数据库 URL(如"sqlite:///my_database.db");- 分区读取:指定
partition_col后可按列分片并行读取。partition_bound_strategy="min-max"按该列最小/最大值均分区间;"percentile"则用PERCENTILE_DISC求百分位分界(如num_partitions=3时取 33 分位与 66 分位)。指定num_partitions时必须同时指定partition_col; - 执行引擎:优先使用 ConnectorX,除非显式传了 SQLAlchemy 连接工厂或方言不被 ConnectorX 支持;
- 下推:过滤、投影、limit 默认尽可能下推进 SQL,可用
disable_pushdowns_to_sql=True关闭; - 方言:基于 SQLGlot 做方言翻译;
- schema 推断默认扫描 10 行(
infer_schema_length)。
df = daft.read_sql("SELECT * FROM my_table", "sqlite:///my_database.db")输入:开放表格式与数据目录
针对数据湖/向量库场景,Daft 提供 Iceberg、Delta Lake、Hudi、Lance 的专属读取入口,均实现为自定义DataSource(内部通过ScanOperatorHandle.from_data_source接入,见 daft/io/iceberg/_iceberg.py)。
read_iceberg
daft.read_iceberg(table, snapshot_id=None, branch=None, tag=None, io_config=None, checkpoint=None, ignore_corrupt_files=False)(daft/io/iceberg/_iceberg.py):
table:PyIceberg Table 或指向 metadata 文件(s3://bucket/path/to/iceberg/metadata.json)的路径;传路径时内部用StaticTable.from_metadata加载;snapshot_id/branch/tag:三者互斥,用于指定时间旅行目标;resolve_snapshot_id负责解析;- 需要安装 PyIceberg;过滤条件(如
df.where(df["foo"] > 5))会被下推进 Iceberg 扫描。
read_deltalake
daft.read_deltalake(table, version=None, io_config=None, ignore_deletion_vectors=False)(daft/io/delta_lake/_deltalake.py):
table:Delta 表 URI 或 Unity Catalog 的UnityCatalogTable实例;version:int 为版本号,str/datetime 为时间戳版本(RFC 3339 / ISO 8601,datetime 默认按 UTC);ignore_deletion_vectors:跳过 deletion vectors 检查;- 需要
deltalake库。
read_hudi 与 read_lance
daft.read_hudi(table_uri, io_config=None, checkpoint=None)(daft/io/hudi/_hudi.py):读取 Hudi 表 URI;daft.read_lance(uri, io_config=None, version=None, asof=None, ...)(daft/io/lance/_lance.py):读取 LanceDB 表,支持version(int 版本号或 str tag)、asof(加载不晚于给定时间的版本)、block_size(最小 I/O 请求大小提示)、commit_lock、index_cache_size(默认 256)等参数。read_lance通过LazyImport惰性加载daft_lance扩展,避免拖慢import daft。
输出:写回文件、表格式与外部服务
write_*系列是阻塞调用,执行 DataFrame 并返回一个包含写入结果(通常为写入文件路径或统计信息)的新 DataFrame。文件类写入统一支持write_mode:"append"(默认)、"overwrite"、"overwrite-partitions"(仅替换被写分区的数据,需配合partition_cols),文件名为随机 UUID。
write_parquet / write_csv / write_json
df.write_parquet(root_dir, compression="snappy", write_mode="append", write_success_file=False, partition_cols=None, io_config=None, column_compression=None, single_file=False)(daft/dataframe/dataframe.py):
compression:"snappy"、"gzip"、"zstd"、"lz4"、"lz4_raw"、"brotli"、"uncompressed"/"none"(不区分大小写);column_compression:按列覆盖压缩算法,key 为点分列路径(如"user.name"表示嵌套 struct 字段);write_success_file:写_SUCCESS标记文件;single_file=True:合并为单文件,此时root_dir被视为精确文件路径;不能与partition_cols或overwrite-partitions组合,且仅支持 native runner(否则抛 ValueError)。
df = daft.from_pydict({"x": [1, 2, 3], "y": ["a", "b", "c"]}) df.write_parquet("output_dir", write_mode="overwrite") df.write_parquet("output.parquet", single_file=True)df.write_csv(root_dir, write_mode="append", partition_cols=None, io_config=None, delimiter=None, quote=None, escape=None, header=True, date_format=None, timestamp_format=None)(daft/dataframe/dataframe.py):date_format/timestamp_format使用 chrono strftime 格式(如"%Y-%m-%d"、"%+");时区感知时间戳会先转换到目标时区再格式化。
df.write_json(root_dir, write_mode="append", partition_cols=None, io_config=None, ignore_null_fields=False, date_format=None, timestamp_format=None)(daft/dataframe/dataframe.py):ignore_null_fields=True可在写出时忽略 Null 字段。
df.write_json("output_dir", write_mode="overwrite") df.write_json("output_dir", date_format="%d/%m/%Y") # "15/01/2024" df.write_json("output_dir", timestamp_format="%+") # "2024-01-15T10:30:45+00:00"write_deltalake 与 write_iceberg
df.write_deltalake(table, partition_cols=None, mode="append", schema_mode=None, name=None, description=None, configuration=None, custom_metadata=None, dynamo_table_name=None, allow_unsafe_rename=False, io_config=None, checkpoint=None)(daft/dataframe/dataframe.py):mode支持append/overwrite/error/ignore;checkpoint为IdempotentCommit,通过daft.idempotence-key元数据实现幂等提交(仅append,需 RayRunner),崩溃恢复时以相同 key 重试不会产生重复提交;df.write_iceberg(table, mode="append", io_config=None, snapshot_properties=None, checkpoint=None, overwrite_filter=None, validate_overwrite_filter=True)(daft/dataframe/dataframe.py):overwrite_filter支持 Daft 表达式或 Iceberg 谓词字符串(如"dt = '2024-01-01'"),实现静态分区覆盖;提交前会用写出文件的列统计验证"删除范围覆盖写出行",无法证明时拒绝写入(统计只会放宽、不会误收)。
write_lance 与向量场景
df.write_lance(uri, mode="create", io_config=None, schema=None, left_on=None, right_on=None, **kwargs)(daft/dataframe/dataframe.py):
mode:"create"(不存在则创建,存在报错)、"append"、"overwrite"、"merge"(向已存在数据集追加新列);schema:可传 Daft Schema 或 pyarrow Schema,写前强制 cast;mode="merge"时通过left_on/right_on(默认"_rowaddr")对齐行,且 DataFrame 需包含fragment_id列;- 返回元数据 DataFrame:
num_fragments、num_deleted_rows、num_small_files、version。
write_sql / write_clickhouse / write_bigtable / write_huggingface / write_turbopuffer
df.write_sql(table_name, conn, write_mode="append", column_types=None, non_primitive_handling=None)(daft/dataframe/dataframe.py):原始类型列经 pandasto_sql写入;非原始类型列(list、struct、map、tensor、image、embedding 等)按non_primitive_handling归一化——"str"(默认,容器转 JSON 文本)、"bytes"(文本的 UTF-8 字节)、"error"(直接报错)。返回单行 DataFrame:total_written_rows、total_written_bytes;df.write_clickhouse(table, *, host, port=None, user=None, password=None, database=None, client_kwargs=None, write_kwargs=None)(daft/dataframe/dataframe.py):写入 ClickHouse 表,同样返回写入行数/字节统计;df.write_bigtable(project_id, instance_id, table_id, row_key_column, column_family_mappings, client_kwargs=None, write_kwargs=None, serialize_incompatible_types=True)(daft/dataframe/dataframe.py):写入 Google Cloud Bigtable,需指定行键列与列族映射(Bigtable 单元格只接受可转字节的类型);df.write_huggingface(repo, split="train", data_dir="data", revision="main", overwrite=False, commit_message="Upload dataset using Daft", commit_description=None, io_config=None)(daft/dataframe/dataframe.py):将 DataFrame 推送为 HF 数据集;df.write_turbopuffer(namespace, api_key=None, region=None, distance_metric=None, schema=None, id_column=None, vector_column=None, client_kwargs=None, write_kwargs=None)(daft/dataframe/dataframe.py):写入 Turbopuffer 向量命名空间,id列必选,vector列在命名空间有向量索引时必选,其余列成为属性;namespace也可传表达式实现按命名空间分片写入。
用户自定义 I/O:DataSource / DataSink
当内置读写无法覆盖业务场景时,Daft 提供低层扩展 API。这些 API 处于早期演进阶段,官方文档明确标记为 experimental(!!! warning)。
DataSource 与 DataSourceTask
daft.io.source.DataSource(daft/io/source.py)是读取数据的低层接口,职责是"把数据切成可并行处理的任务":
name:调试用源名称;schema:各任务输出 RecordBatch 共享的 schema;get_partition_fields:声明分区字段(磁盘布局,用于逐行注入值);get_clustering_keys:声明执行期分布保证(hash或range),可让优化器跳过 shuffle;supports_count_pushdown:能否吸收 count 聚合下推(为 True 时优化器可用get_tasks从目录元数据直接产出行数而无需扫数据文件);async def get_tasks(self, pushdowns) -> AsyncIterator[DataSourceTask]:在执行期按 pushdowns 产出任务;read():将该 DataSource 作为 DataFrame 读取(内部经ScanOperatorHandle.from_data_source构造 TabularScan)。
DataSourceTask(daft/io/source.py)表示一个可独立处理的数据分区,推荐覆盖async def read() -> AsyncIterator[RecordBatch](旧的get_micro_partitions已标记 deprecated)。特别地,静态工厂DataSourceTask.parquet(...)用原生 Parquet reader 创建扫描任务,是构建 Iceberg/Paimon 等目录连接器时读取 Parquet 文件的推荐方式,可传num_rows、size_bytes(供任务合并启发式)、pushdowns、partition_values、stats与iceberg_delete_files(Iceberg position delete 文件)。
DataSink 与 WriteResult
daft.io.sink.DataSink(daft/io/sink.py)是写外部存储的接口,配合df.write_sink(sink)使用(daft/dataframe/dataframe.py)。写入时序如下:
start():写入开始时调用一次(可初始化资源、打开连接、开启事务);- DataFrame 执行,输出被切分为 micropartitions;
write(micropartitions) -> Iterator[WriteResult]:对每个 micropartition 并行调用;- 所有写结果在单节点汇总;
finalize(write_results) -> MicroPartition:产出最终结果(write_sink返回的 DataFrame 即源于此),其 schema 必须与sink.schema()一致,否则报错。
WriteResult(daft/io/sink.py)是write()返回值的包装,包含result、bytes_written、rows_written三个字段。safe_write将不可序列化的异常包装为带 sink 名称的 RuntimeError,便于分布式环境下定位问题。write_clickhouse、write_bigtable、write_huggingface、write_turbopuffer等正是基于write_sink实现的(见 daft/dataframe/dataframe.py)。
Pushdowns:谓词、投影与 limit 下推
Daft 在扫描阶段支持 predicate(谓词)、projection(投影)与 limit 三类下推,这是减少扫描 IO、提升查询性能的关键机制。
daft.io.pushdowns.Pushdowns(daft/io/pushdowns.py)是一个 frozen dataclass,字段包括:
| 字段 | 类型 | 含义 |
|---|---|---|
filters | Expression \| None | 作用于行的过滤谓词 |
partition_filters | Expression \| None | 作用于分区/文件的过滤谓词(分区裁剪) |
columns | list[str] \| None | 投影的列名列表 |
limit | int \| None | 返回行数上限 |
aggregation | Expression \| None | count 聚合下推 |
Pushdowns在查询规划期被发送给扫描源,提供_from_pypushdowns/_to_pypushdowns与 Rust 侧PyPushdowns互转,以及filter_required_column_names()(返回谓词依赖的列集合,用于确定投影下推的最小列集)。此外,SupportsPushdownFilters.push_filters(filters)接口用于实现过滤下推,返回(pushed_filters, post_filters)——被推入扫描的谓词与仍需扫描后求值的谓词。
daft.io.scan.ScanOperator(daft/io/scan.py)是旧版 Python 扫描基类(官方注释说明正在迁移到daft.io.source.DataSource),通过can_absorb_filter/can_absorb_limit/can_absorb_select/supports_count_pushdown声明各类下推能力,to_scan_tasks(pushdowns)将扫描算子按给定 pushdowns 转换为扫描任务。
具体到各读取器,下推效果可验证如下:
- 文件格式:
read_parquet等将参数编码进ParquetSourceConfig/CsvSourceConfig/JsonSourceConfig后交给get_tabular_files_scan(daft/io/common.py),原生扫描器据此做行组/页级裁剪; - 目录连接器:
read_iceberg的 docstring 明确写出"Filters on this dataframe can now be pushed into the read operation from Iceberg",read_deltalake同理(见 daft/io/delta_lake/_deltalake.py); - SQL 读取:
read_sql会把过滤、投影、limit 翻译进 SQL,disable_pushdowns_to_sql=True可关闭(daft/io/_sql.py); - 自定义源:
DataSource.get_tasks(pushdowns)收到的Pushdowns即可用于自主决定产出哪些任务(如 WebDataset 按pushdowns.columns投影、按pushdowns.limit收缩 batch_size,见 daft/io/webdataset/_webdataset.py)。
实践建议
- 内存构建选
from_pydict/from_pylist:小型数据与测试首选;跨语言互操作选from_arrow(支持 PyCapsule 接口,生态最广)。 - 文件读取统一用 glob:
read_parquet/read_csv/read_json/read_blob都支持*/?/[...]/**通配符与目录路径,配合file_path_column/hive_partitioning可保留来源信息与分区语义。 - 云存储记得传
io_config:匿名桶可用IOConfig(s3=S3Config(region="...", anonymous=True));更多对象存储配置见 docs/connectors/index.md。 - 大数据量用分区 + 下推:指定
partition_cols写分区表、利用where触发谓词下推、用limit触发 limit 下推,从源头减少扫描数据量。 - 分布式运行注意 IO 行为:在 RayRunner 下,
read_parquet/read_warc/read_deltalake等会自动把_multithreaded_io置为 False 以减少资源争用;from_ray_dataset/from_dask_dataframe则必须先daft.set_runner_ray()。 - 跨运行断点续跑:文件读取与 Delta/Iceberg 写入都支持
checkpoint参数(读取用CheckpointConfig,写入用IdempotentCommit),适用于长任务重跑场景(详见 docs/use-case/checkpointing.md)。 - 扩展新数据源:优先考虑在
DataSource/DataSourceTask(读)与DataSink/WriteResult(写)之上实现,注意这些 API 仍处于实验阶段、可能变化。
【免费下载链接】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),仅供参考