ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Daft I/O 全指南:从内存、文件、数据湖到数据目录的读写 API 详解

Daft I/O 全指南:从内存、文件、数据湖到数据目录的读写 API 详解 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/DaftDaft 作为面向 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/Odaft/io/source.py 与 daft/io/sink.py标记为 experimental以及Pushdowns 下推机制daft/io/pushdowns.py后者是 Daft 在扫描阶段省去不必要 IO 的核心手段。与各种 Connector 的更多用法可进一步参考 docs/connectors/index.md对象存储、开放表格式、数据库、文件与目录等。输入从内存数据构建 DataFramefrom_*系列 API 用于把已经存在于进程内存中的数据对象转换为 Daft DataFrame是快速上手与单元测试的常用入口全部定义在 daft/convert.py。from_pydict 与 from_pylistdaft.from_pydict(data: dict[str, InputListType])以列名 → 序列的字典构建 DataFrame。每个 value 必须是等长的 Python list、NumPy array 或 PyArrow arraydaft/convert.py。daft.from_pylist(data: list[dict[str, Any]])以行为单位列表中的每个 dict 是一行key 为列名构建 DataFramedaft/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_arrowdaft.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_dataframedaft.from_ray_dataset(ds)从 Ray Dataset 构建 DataFramedaft/convert.pydaft.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把文件列表变成 DataFramedaft.from_glob_path(path, io_configNone)根据 glob 模式返回一个包含文件元数据的 DataFramedaft/io/file_path.py支持四种通配符*匹配任意多个字符含 0 个?匹配任意单个字符[...]匹配方括号内的任意单个字符**递归匹配任意多层目录。返回的 DataFrame 包含三列path文件/目录路径、size字节大小、rowsParquet 对象的行数其他格式为 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(s3S3Config(regionus-west-2, anonymousTrue))。read_parquetdaft.read_parquet(path, row_groupsNone, infer_schemaTrue, schemaNone, io_configNone, file_path_columnNone, hive_partitioningFalse, coerce_int96_timestamp_unitNone, ignore_corrupt_filesFalse, checkpointNone)daft/io/_parquet.pyrow_groups按文件指定要读取的行组列表仅在读取多个非通配文件时支持且列表长度必须与 path 数量一致否则抛ValueErrorinfer_schemaFalse时必须同时提供schema否则报错schema当infer_schemaTrue时作为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(s3S3Config(regionus-west-2, anonymousTrue)) df daft.read_parquet(s3://path/to/files-*.parquet, io_configio_config)实现上read_parquet将参数打包为ParquetSourceConfig与StorageConfig再通过get_tabular_files_scandaft/io/common.py构造TabularFilesScan逻辑计划。值得注意的细节在 Ray runner 下默认关闭多线程 IOmultithreaded_io runners.get_or_create_runner().name ! ray以减少每个 Ray worker 的线程池与连接数争用。read_csvdaft.read_csv(path, infer_schemaTrue, schemaNone, has_headersTrue, delimiterNone, double_quoteTrue, quoteNone, escape_charNone, commentNone, allow_variable_columnsFalse, io_configNone, file_path_columnNone, hive_partitioningFalse, ignore_corrupt_filesFalse, checkpointNone)daft/io/_csv.pydelimiter字段分隔符默认,double_quote是否支持双引号转义默认 Truequote包裹含分隔符字段的引号字符默认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_configio_config)read_jsondaft.read_json(path, infer_schemaTrue, schemaNone, io_configNone, file_path_columnNone, hive_partitioningFalse, skip_empty_filesFalse, checkpointNone)daft/io/_json.py读取**行分隔 JSONJSONL**文件skip_empty_filesTrue可跳过空文件。read_blob以原始字节读取任意文件daft.read_blob(path, *, max_connections32, on_errorraise, io_configNone)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_frameNone, *, sample_interval_secondsNone, io_configNone)daft/io/av/init.py将视频流式读取为 DataFrame of images需要 PyAVpip install av。输出字段包括path、frame_index、frame_time秒、frame_time_base、frame_pts、frame_dts、frame_duration、is_key_frame。is_key_frameTrue 只取关键帧、False 只取非关键帧、None 取全部sample_interval_seconds按帧时间近似采样取时间戳 ≥ 目标时间点的首帧无有效时间戳的帧被跳过。df daft.read_video_frames(/path/to/file.mp4, image_height480, image_width640) # 约每秒一帧 df daft.read_video_frames(/path/to/file.mp4, image_height480, image_width640, sample_interval_seconds1.0)read_warc 与 read_webdataset网页抓取与多模态数据集daft.read_warc(path, io_configNone, file_path_columnNone, checkpointNone)daft/io/_warc.py读取 WARC 或 gzip 压缩的 WARC 文件实验特性返回 DataFrame 含强制元数据列WARC-Record-ID、WARC-Type、WARC-Date、Content-Length与可选字段daft.read_webdataset(path, io_configNone, batch_size1000)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_huggingfacedaft.read_huggingface(repo, io_configNone, formatNone)daft/io/huggingface/init.py读取 Hugging Face 数据集repo形如username/dataset_nameformatNone默认或parquet走快速路径read_parquet(fhf://datasets/{repo})若 parquet 文件不存在glob 无匹配或返回 400parquet 尚未生成则回退到datasets库formatwebdataset等价于read_webdataset(fhf://datasets/{repo}/**/*.tar)。read_kafka流式数据源daft.read_kafka(topics, bootstrap_servers, group_id, startearliest, endlatest, ...)基于 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_colNone, num_partitionsNone, partition_bound_strategymin-max, disable_pushdowns_to_sqlFalse, infer_schemaTrue, infer_schema_length10, schemaNone)daft/io/_sql.pyconnSQLAlchemy 连接工厂Callable[[], Connection]或数据库 URL如sqlite:///my_database.db分区读取指定partition_col后可按列分片并行读取。partition_bound_strategymin-max按该列最小/最大值均分区间percentile则用PERCENTILE_DISC求百分位分界如num_partitions3时取 33 分位与 66 分位。指定num_partitions时必须同时指定partition_col执行引擎优先使用 ConnectorX除非显式传了 SQLAlchemy 连接工厂或方言不被 ConnectorX 支持下推过滤、投影、limit 默认尽可能下推进 SQL可用disable_pushdowns_to_sqlTrue关闭方言基于 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_icebergdaft.read_iceberg(table, snapshot_idNone, branchNone, tagNone, io_configNone, checkpointNone, ignore_corrupt_filesFalse)daft/io/iceberg/_iceberg.pytablePyIceberg 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_deltalakedaft.read_deltalake(table, versionNone, io_configNone, ignore_deletion_vectorsFalse)daft/io/delta_lake/_deltalake.pytableDelta 表 URI 或 Unity Catalog 的UnityCatalogTable实例versionint 为版本号str/datetime 为时间戳版本RFC 3339 / ISO 8601datetime 默认按 UTCignore_deletion_vectors跳过 deletion vectors 检查需要deltalake库。read_hudi 与 read_lancedaft.read_hudi(table_uri, io_configNone, checkpointNone)daft/io/hudi/_hudi.py读取 Hudi 表 URIdaft.read_lance(uri, io_configNone, versionNone, asofNone, ...)daft/io/lance/_lance.py读取 LanceDB 表支持versionint 版本号或 str tag、asof加载不晚于给定时间的版本、block_size最小 I/O 请求大小提示、commit_lock、index_cache_size默认 256等参数。read_lance通过LazyImport惰性加载daft_lance扩展避免拖慢import daft。输出写回文件、表格式与外部服务write_*系列是阻塞调用执行 DataFrame 并返回一个包含写入结果通常为写入文件路径或统计信息的新 DataFrame。文件类写入统一支持write_modeappend默认、overwrite、overwrite-partitions仅替换被写分区的数据需配合partition_cols文件名为随机 UUID。write_parquet / write_csv / write_jsondf.write_parquet(root_dir, compressionsnappy, write_modeappend, write_success_fileFalse, partition_colsNone, io_configNone, column_compressionNone, single_fileFalse)daft/dataframe/dataframe.pycompressionsnappy、gzip、zstd、lz4、lz4_raw、brotli、uncompressed/none不区分大小写column_compression按列覆盖压缩算法key 为点分列路径如user.name表示嵌套 struct 字段write_success_file写_SUCCESS标记文件single_fileTrue合并为单文件此时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_modeoverwrite) df.write_parquet(output.parquet, single_fileTrue)df.write_csv(root_dir, write_modeappend, partition_colsNone, io_configNone, delimiterNone, quoteNone, escapeNone, headerTrue, date_formatNone, timestamp_formatNone)daft/dataframe/dataframe.pydate_format/timestamp_format使用 chrono strftime 格式如%Y-%m-%d、%时区感知时间戳会先转换到目标时区再格式化。df.write_json(root_dir, write_modeappend, partition_colsNone, io_configNone, ignore_null_fieldsFalse, date_formatNone, timestamp_formatNone)daft/dataframe/dataframe.pyignore_null_fieldsTrue可在写出时忽略 Null 字段。df.write_json(output_dir, write_modeoverwrite) df.write_json(output_dir, date_format%d/%m/%Y) # 15/01/2024 df.write_json(output_dir, timestamp_format%) # 2024-01-15T10:30:4500:00write_deltalake 与 write_icebergdf.write_deltalake(table, partition_colsNone, modeappend, schema_modeNone, nameNone, descriptionNone, configurationNone, custom_metadataNone, dynamo_table_nameNone, allow_unsafe_renameFalse, io_configNone, checkpointNone)daft/dataframe/dataframe.pymode支持append/overwrite/error/ignorecheckpoint为IdempotentCommit通过daft.idempotence-key元数据实现幂等提交仅append需 RayRunner崩溃恢复时以相同 key 重试不会产生重复提交df.write_iceberg(table, modeappend, io_configNone, snapshot_propertiesNone, checkpointNone, overwrite_filterNone, validate_overwrite_filterTrue)daft/dataframe/dataframe.pyoverwrite_filter支持 Daft 表达式或 Iceberg 谓词字符串如dt 2024-01-01实现静态分区覆盖提交前会用写出文件的列统计验证删除范围覆盖写出行无法证明时拒绝写入统计只会放宽、不会误收。write_lance 与向量场景df.write_lance(uri, modecreate, io_configNone, schemaNone, left_onNone, right_onNone, **kwargs)daft/dataframe/dataframe.pymodecreate不存在则创建存在报错、append、overwrite、merge向已存在数据集追加新列schema可传 Daft Schema 或 pyarrow Schema写前强制 castmodemerge时通过left_on/right_on默认_rowaddr对齐行且 DataFrame 需包含fragment_id列返回元数据 DataFramenum_fragments、num_deleted_rows、num_small_files、version。write_sql / write_clickhouse / write_bigtable / write_huggingface / write_turbopufferdf.write_sql(table_name, conn, write_modeappend, column_typesNone, non_primitive_handlingNone)daft/dataframe/dataframe.py原始类型列经 pandasto_sql写入非原始类型列list、struct、map、tensor、image、embedding 等按non_primitive_handling归一化——str默认容器转 JSON 文本、bytes文本的 UTF-8 字节、error直接报错。返回单行 DataFrametotal_written_rows、total_written_bytesdf.write_clickhouse(table, *, host, portNone, userNone, passwordNone, databaseNone, client_kwargsNone, write_kwargsNone)daft/dataframe/dataframe.py写入 ClickHouse 表同样返回写入行数/字节统计df.write_bigtable(project_id, instance_id, table_id, row_key_column, column_family_mappings, client_kwargsNone, write_kwargsNone, serialize_incompatible_typesTrue)daft/dataframe/dataframe.py写入 Google Cloud Bigtable需指定行键列与列族映射Bigtable 单元格只接受可转字节的类型df.write_huggingface(repo, splittrain, data_dirdata, revisionmain, overwriteFalse, commit_messageUpload dataset using Daft, commit_descriptionNone, io_configNone)daft/dataframe/dataframe.py将 DataFrame 推送为 HF 数据集df.write_turbopuffer(namespace, api_keyNone, regionNone, distance_metricNone, schemaNone, id_columnNone, vector_columnNone, client_kwargsNone, write_kwargsNone)daft/dataframe/dataframe.py写入 Turbopuffer 向量命名空间id列必选vector列在命名空间有向量索引时必选其余列成为属性namespace也可传表达式实现按命名空间分片写入。用户自定义 I/ODataSource / DataSink当内置读写无法覆盖业务场景时Daft 提供低层扩展 API。这些 API 处于早期演进阶段官方文档明确标记为 experimental!!! warning。DataSource 与 DataSourceTaskdaft.io.source.DataSourcedaft/io/source.py是读取数据的低层接口职责是把数据切成可并行处理的任务name调试用源名称schema各任务输出 RecordBatch 共享的 schemaget_partition_fields声明分区字段磁盘布局用于逐行注入值get_clustering_keys声明执行期分布保证hash或range可让优化器跳过 shufflesupports_count_pushdown能否吸收 count 聚合下推为 True 时优化器可用get_tasks从目录元数据直接产出行数而无需扫数据文件async def get_tasks(self, pushdowns) - AsyncIterator[DataSourceTask]在执行期按 pushdowns 产出任务read()将该 DataSource 作为 DataFrame 读取内部经ScanOperatorHandle.from_data_source构造 TabularScan。DataSourceTaskdaft/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_filesIceberg position delete 文件。DataSink 与 WriteResultdaft.io.sink.DataSinkdaft/io/sink.py是写外部存储的接口配合df.write_sink(sink)使用daft/dataframe/dataframe.py。写入时序如下start()写入开始时调用一次可初始化资源、打开连接、开启事务DataFrame 执行输出被切分为 micropartitionswrite(micropartitions) - Iterator[WriteResult]对每个 micropartition 并行调用所有写结果在单节点汇总finalize(write_results) - MicroPartition产出最终结果write_sink返回的 DataFrame 即源于此其 schema 必须与sink.schema()一致否则报错。WriteResultdaft/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.Pushdownsdaft/io/pushdowns.py是一个 frozen dataclass字段包括字段类型含义filtersExpression \| None作用于行的过滤谓词partition_filtersExpression \| None作用于分区/文件的过滤谓词分区裁剪columnslist[str] \| None投影的列名列表limitint \| None返回行数上限aggregationExpression \| Nonecount 聚合下推Pushdowns在查询规划期被发送给扫描源提供_from_pypushdowns/_to_pypushdowns与 Rust 侧PyPushdowns互转以及filter_required_column_names()返回谓词依赖的列集合用于确定投影下推的最小列集。此外SupportsPushdownFilters.push_filters(filters)接口用于实现过滤下推返回(pushed_filters, post_filters)——被推入扫描的谓词与仍需扫描后求值的谓词。daft.io.scan.ScanOperatordaft/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_scandaft/io/common.py原生扫描器据此做行组/页级裁剪目录连接器read_iceberg的 docstring 明确写出Filters on this dataframe can now be pushed into the read operation from Icebergread_deltalake同理见 daft/io/delta_lake/_deltalake.pySQL 读取read_sql会把过滤、投影、limit 翻译进 SQLdisable_pushdowns_to_sqlTrue可关闭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 接口生态最广。文件读取统一用 globread_parquet/read_csv/read_json/read_blob都支持*/?/[...]/**通配符与目录路径配合file_path_column/hive_partitioning可保留来源信息与分区语义。云存储记得传io_config匿名桶可用IOConfig(s3S3Config(region..., anonymousTrue))更多对象存储配置见 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),仅供参考
返回列表