Daft read_text 深度指南:用 DataFrame 高效读取本地与对象存储中的文本文件

发布时间:2026/9/17 15:37:14
Daft read_text 深度指南:用 DataFrame 高效读取本地与对象存储中的文本文件 Daft read_text 深度指南用 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本篇指南围绕 Daft 的daft.read_text()接口展开介绍如何将行式文本文件日志、配置文件、纯文本数据等读入 DataFrame并覆盖整文件读取、空行处理、Hive 分区、通配符等全部读取选项。读完本文你不仅掌握每个参数的用法与默认值还能结合仓库源码理解 Daft 文本读取器在 Rust 侧的流式读取机制、缓冲/分块策略以及对非 UTF-8 编码文件的实际限制。基本用法从本地到对象存储read_text是 daft/io/_text.py 中定义的公开 API标注了PublicAPI它把按行组织的文本文件转换为 DataFrame默认情况下文件的每一行成为 DataFrame 的一行。本地文件示例import daft df daft.read_text(/path/to/file.txt) df.show()如果文件在对象存储上通过io_config传入对应的云配置。S3 场景import daft from daft.io import S3Config, IOConfig io_config IOConfig(s3S3Config(region_nameus-west-2, anonymousTrue)) df daft.read_text(s3://my-bucket/logs/*.txt, io_configio_config) df.show()GCS 场景import daft from daft.io import GCSConfig, IOConfig io_config IOConfig(gcsGCSConfig(anonymousTrue)) df daft.read_text(gs://my-bucket/logs/*.txt, io_configio_config) df.show()从源码可以看出两个细节若未显式传入io_configread_text会回退到 Daft 上下文的默认 IO 配置daft/io/_text.pycontext.get_context().daft_planning_config.default_io_config即全局配置过的凭证与端点同样生效传入空路径列表会直接抛出ValueError: Cannot read DataFrame from empty list of text filepathsdaft/io/_text.py该行为在 tests/io/test_text.py 中有对应断言。输出 Schemaread_text返回的 DataFrame 具有固定 schema——这一点在 Python 入口中写得很明确schema {text: DataType.string()}daft/io/_text.py不做 schema 推断infer_schemaFalse列名类型说明textstring输入文件的行内容当whole_textTrue时为一整个文件的内容在此基础上file_path_column和hive_partitioning会在 schema 中追加对应列见下文。测试用例 tests/io/test_text.py 验证了基础读取的 schema 即为单列textpa.string()。读取选项全解整文件模式whole_text默认whole_textFalse文件按行切分。设为True后每个文件成为 DataFrame 的一行# Each file becomes a single row df daft.read_text(/path/to/files/*.txt, whole_textTrue) df.show()适用于需要整体处理文件内容的场景例如文档处理、向量化嵌入或内容本身不应被换行符切断的情况。Rust 侧的实现印证了这一机制src/daft-text/src/read.rs 中stream_text检测到whole_text后走read_into_whole_text_stream分支——用read_to_string一次性读入全部内容产出单个 1 行的RecordBatch而按行模式则走read_into_line_chunk_stream逐行流式产出。注意whole_textTrue时_chunk_size参数不生效文档字符串中已注明。空行处理skip_blank_lines默认值为True去掉首尾空白后为空的行会被跳过。要保留这些行df daft.read_text(/path/to/file.txt, skip_blank_linesFalse)tests/io/test_text.py 精确刻画了这个语义对内容line1\n\nline2\n \n\t\nskip_blank_linesFalse得到[line1, , line2, , \t]纯空白行被原样保留包括 和\t而skip_blank_linesTrue只得到[line1, line2]。当whole_textTrue时该选项的语义变为跳过整个内容为空白的文件——源码中对应if convert_options.skip_blank_lines content.trim().is_empty() { return; }src/daft-text/src/read.rs。文件编码encoding与当前实现的限制官方用法示例df daft.read_text(/path/to/file.txt, encodinglatin-1)默认编码为UTF-8。这里必须指出一个与示例相关的重要事实当前仓库的 Rust 读取器仅支持 UTF-8。src/daft-text/src/read.rs 中stream_text会校验编码参数非utf-8/utf8忽略大小写的值会直接抛出ValueError: Unsupported text encoding: {encoding}. Only UTF-8 is currently supported.。也就是说传入encodinglatin-1在现阶段运行时会失败如需读取非 UTF-8 文件建议先在 ETL 前置步骤中完成转码。这是使用encoding参数前必须知晓的适用前提。追加源文件路径列file_path_column给 DataFrame 增加一列记录每行来自哪个源文件df daft.read_text(/path/to/files/*.txt, file_path_columnsource_file) df.show()该参数与hive_partitioning一起在get_tabular_files_scan中接入扫描任务daft/io/_text.py。tests/io/test_text.py 验证了多文件场景下每一行都携带正确的源路径a1/a2指向a.txtb1指向b.txt路径值保留 glob 展开后的完整文件路径。Hive 分区推断hive_partitioning从文件路径中解析keyvalue目录段作为列# For paths like /data/year2024/month01/file.txt df daft.read_text(/data/**/*.txt, hive_partitioningTrue) df.show() # Includes year and month columnstests/io/test_text.py 用keya/keyb两个目录验证了该能力开启后 schema 变为textkey可选的路径列且分区值按路径正确回填到每一行。对于日志归档、爬虫抓取结果这类按日期/来源分目录组织的数据这一选项可以避免额外解析路径列。通配符模式read_text的path参数支持 glob 模式匹配多个文件模式说明*匹配任意数量的字符?匹配任意单个字符[...]匹配方括号内的任意单个字符**递归匹配目录# All .txt files in a directory df daft.read_text(/logs/*.txt) # All .log files recursively df daft.read_text(/logs/**/*.log) # Files matching a pattern df daft.read_text(/logs/server-[0-9].txt)除字符串外path也接受路径列表str | list[str]daft/io/_text.py测试中[str(file_a), str(file_b)]的列表传参用法tests/io/test_text.py可直接复制使用。压缩文件按扩展名自动解码这是一个文档未提及、但由源码与测试共同确认的能力read_text支持读取压缩文本文件。Rust 侧在打开 reader 后会按 URI 扩展名尝试匹配压缩编解码器src/daft-text/src/read.rsCompressionCodec::from_uri(uri)命中则包裹解码层。tests/io/test_text.py 验证了读取compressed.txt.gz时能正常拿到解压缩后的行[l1, l2]Rust 单测src/daft-text/src/read.rs同样覆盖了 gzip 场景下按行分块的正确性。因此带.gz等压缩扩展名的文本归档可以直接喂给read_text无需手动解压。流式读取原理缓冲与分块从源码结构看文本读取是一条流式管道这一设计决定了它在处理超大日志文件时的内存行为打开 reader本地文件经BufReader包装远程流S3/GCS 对象经StreamReader包装两者再统一挂接可选的压缩解码层src/daft-text/src/read.rs按行流式消费行模式下通过reader.lines()逐行拉取行先累积进chunk向量达到chunk_size就产出一次RecordBatchsrc/daft-text/src/read.rs默认值buffer_size默认 8 MiB8 * 1024 * 1024字节chunk_size默认 65536 行src/daft-text/src/read.rs。这两个默认值可以通过read_text的前缀下划线参数_buffer_size、_chunk_size覆盖daft/io/_text.py它们最终写入TextSourceConfigsrc/daft-scan/src/file_format_config.rs。下划线前缀表明这是面向调优的低层参数日常使用可以完全忽略chunk_size在whole_textTrue时无意义。配置结构在 Rust 侧的完整定义见 src/daft-scan/src/file_format_config.rs其中TextSourceConfig默认值为encodingutf-8、skip_blank_linestrue、whole_textfalse与 Python 入口签名默认值一一对应。实战场景处理日志文件读取一批日志并过滤错误行同时保留来源文件以便溯源import daft from daft import col # Read log files and filter for errors df daft.read_text(/var/log/app/*.log, file_path_columnlog_file) errors df.where(col(text).contains(ERROR)) errors.show()日志文件通常体量大、行数多行式流式读取 谓词下推的组合在此场景下是自然的处理路径若日志按天滚动存放到date2024-01-01/之类的目录结构中可叠加hive_partitioningTrue直接得到日期列。文档处理把整个文档作为一行读入供后续嵌入或全文分析使用import daft # Read entire documents df daft.read_text(/documents/*.txt, whole_textTrue, file_path_columndoc_path) # Process with embeddings or other analysis df.show()whole_textTruefile_path_column的组合使每行对应一篇完整文档及其路径可以直接衔接 Daft 的 AI 函数如嵌入生成做批处理。行为边界与验证依据小结空文件读取空.txt不报错返回 0 行且 schema 仍为text: stringtests/io/test_text.py无结尾换行hello\nworld末行无\n同样被正确读出为两行tests/io/test_text.py编码限制当前 Rust 实现仅接受 UTF-8其他编码会抛ValueErrorsrc/daft-text/src/read.rsRust 侧集成测试S3/MinIO 上的文本读取有专门集成测试 tests/integration/io/test_text_s3_minio.py与本地测试共同覆盖远程读取路径。完整的用例与断言可直接参考 tests/io/test_text.py需要复现或回归验证时以该文件为准。【免费下载链接】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),仅供参考