dlt 加载文件格式(Loader File Format)完全指南:Parquet、CSV、JSONL 与 SQL INSERT 的底层实现
【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy 🛠️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt
本文基于 dlt 官方文档中的 Loader file format 章节展开,系统讲解 dlt 管线中四种加载文件格式(Parquet、CSV、JSONL、SQL INSERT)的配置方式、参数细节与适用场景。读完本文,你将能够按资源或管线级别选择合适的loader_file_format,掌握[data_writer]与[normalize.data_writer]配置段的全部参数,并理解每种格式在 dlt 源码中的写入实现——包括ParquetDataWriter、CsvWriter等核心写入器的工作机制。
一、什么是 Loader File Format,以及如何配置
Loader file format(加载文件格式)决定了管线在 prepare/normalize 阶段如何把提取出的数据准备、打包,并最终写入目标端。它是 dlt 数据链路中"中间表示"的开关:不同目标端对格式的支持不同,同一份数据在不同目标端可能走 Parquet、CSV、JSONL 或 SQL INSERT 语句。
dlt 提供两个层级的配置入口:
- 资源级:
@dlt.resource(file_format=...)。该设置优先于run方法传入的loader_file_format。特殊取值preferred表示"使用目标端偏好的格式"; - 管线级:
pipeline.run(..., loader_file_format=...)。
从源码看,run方法对loader_file_format做了严格校验:非法取值会直接抛出InvalidPipelineException,且传入的格式必须存在于目标端能力(destination capabilities)声明的supported_loader_file_formats列表中。相关实现见 pipeline.run 的参数校验 与 _verify_destination_capabilities——如果指定的格式不受当前目标端支持,dlt 会列出该目标端实际支持的格式集合,帮助快速定位问题。
@def_resource @dlt.resource(file_format="parquet") # 资源级指定 def my_table(): yield [...] p = dlt.pipeline() p.run(my_source(), loader_file_format="jsonl") # 管线级指定资源级的file_format="preferred"会让 dlt 使用目标端偏好的格式(caps.preferred_loader_file_format),这在目标端能力升级或切换目标端时能保持行为最稳定。
二、Parquet:列式存储格式的完整配置
Apache Parquet 是 Apache Hadoop 生态中的开源列式存储格式。dlt 在配置为 Parquet 格式后能够以该格式存储数据,前提是需要pyarrow包——可以作为 dlt 的 extra 一并安装:
pip install "dlt[parquet]"2.1 目标端能力自动配置(Destination autoconfig)
dlt 会基于 destination capabilities 自动配置 Parquet 写入器,避免用户手工处理类型映射:
- decimal / wei 精度:用于挑选正确的 decimal 类型,并设置 precision 和 scale;
- timestamp 精度:用于挑选正确的 timestamp 类型分辨率(秒、微秒或纳秒);
supports_dictionary_encoding:控制常量列(如_dlt_load_id)是否使用字典编码的 Arrow 数组。字典编码对重复值内存友好,但并非所有目标端都支持,默认true。在 ParquetFormatConfiguration 中可以看到该字段的默认值与注释:设为False适合使用 ADBC 驱动且不支持字典类型的目标端(如 MSSQL)。
类型映射的实际执行发生在写入 header 阶段:ParquetDataWriter.write_header 通过columns_to_arrow(columns_schema, self._caps, self.timestamp_timezone)把 dlt 的表 schema 转成 Arrow schema,self._caps即传入的目标端能力上下文——decimal、timestamp 精度的自动挑选就发生在这里。
2.2 Writer settings(写入器参数全解)
dlt 底层使用 pyarrow 的 parquet writer 创建文件。以下是各参数的含义、默认值与注意事项,与 ParquetFormatConfiguration 中的默认值一一对应:
| 参数 | 默认值 | 说明 |
|---|---|---|
flavor | None(pyarrow 默认) | 对 schema 做兼容性清洗以适配不同目标系统,例如"spark" |
version | "2.4" | 决定可用的 Parquet 逻辑类型集合;>=2.6才支持纳秒精度时间戳 |
compression | "snappy" | Parquet 内部压缩编码,可选"snappy"、"gzip"、"brotli"、"zstd"、"lz4"、"none"。务必确认你的数据库能解压所选编码 |
data_page_size | None(pyarrow 默认) | 列块内数据页目标编码大小(字节) |
row_group_size | None | 行组行数,配合大内存缓冲或单一大表时有用,见下文 |
timestamp_timezone | 已弃用 | 请改用 context timezone,所有写入器遵循;未设置时写入器会用 context timezone 标注 tz-aware 列。空字符串有特殊含义,见 2.3 |
coerce_timestamps | None | 时间戳强制转换的分辨率:s/ms/us/ns |
allow_truncated_timestamps | False | 截断时间戳若丢失精度则抛异常 |
write_page_index | False | 是否写入 page index |
use_content_defined_chunking | False | 是否启用 Content-Defined Chunking;要求pyarrow>=21.0.0,低版本会被忽略 |
arrow_concat_promote_options | "none" | 拼接多个 Arrow 表时的类型提升策略:"none"(要求 schema 一致、零拷贝)、"default"(同族提升,如 int32→int64)、"permissive"(跨族提升,如 int64→double) |
这些参数最终在 _create_writer 中组装为pyarrow.parquet.ParquetWriter的关键字参数;值得注意的是源码中显式判断了pyarrow.__version__主版本是否>= 21才传入use_content_defined_chunking,与文档"低版本被忽略"的说明完全一致。
version字段还参与一个关键推断:max_timestamp_precision 方法中,flavor="spark"按 INT96 视为纳秒级,version >= 2.6支持纳秒,否则最高微秒;若设置了coerce_timestamps则取两者较小值。
[data_writer] # 示例取值 flavor="spark" version="2.4" compression="zstd" data_page_size=1048576也可用环境变量配置:
DATA_WRITER__FLAVOR DATA_WRITER__VERSION DATA_WRITER__COMPRESSION DATA_WRITER__DATA_PAGE_SIZE DATA_WRITER__TIMESTAMP_TIMEZONE DATA_WRITER__ARROW_CONCAT_PROMOTE_OPTIONS提示:dlt 默认使用 Parquet 2.4,并把时间戳强制转换为微秒、静默截断纳秒。这一设定提供了与数据库系统(包括加载默认纳秒精度的 pandas DataFrame)最好的互操作性。如需保留纳秒精度时间戳,请设置
version="2.6"。
2.3 作用域:normalize 阶段与按资源/按 source 覆盖
Parquet 写入发生在 normalize 阶段,因此可以只对该阶段生效:
NORMALIZE__DATA_WRITER__FLAVOR=spark设置 normalize 阶段 Parquet 内部编码用NORMALIZE__DATA_WRITER__COMPRESSION=zstd。注意它只作用于 Parquet 的内部 codec,不会改变jsonl、csv等文本文件的 gzip 外层压缩;要禁用文本文件的 gzip 压缩,请使用data_writer.disable_compression。
当 source/resource 产出 Arrow 表 / pandas DataFrame / polars DataFrame 时,可以按 source 粒度覆盖:
SOURCES__<SOURCE_MODULE>__<SOURCE_NAME>__DATA_WRITER__FLAVOR=spark2.4 时间戳与时区
dlt 在所有精度(秒到纳秒)下都会为 timestamp 列附加时区(UTC 调整)信息,并在目标端创建 tz-aware 的 timestamp 列。DuckDB 是例外。列上timezonehint 设为False时保持 naive,由 context timezone 决定 dlt 使用哪个时区。
源码印证了这一点:ParquetDataWriter.init中,timestamp_timezone为None时取get_context_timezone_name(),否则使用传入值——即"未设置时回落到 context timezone"的行为。
禁用时区/UTC 调整的两种方式:
- 把
flavor设为spark。所有时间戳将通过已弃用的int96物理类型生成,不带逻辑类型; - 把
timestamp_timezone设为空字符串(DATA_WRITER__TIMESTAMP_TIMEZONE=""),生成不带 UTC 调整信息的逻辑类型。
据文档所述,按现有知识,Arrow 会把 tz-aware 的 DateTime 转换为 UTC 后存入 Parquet,不再保留时区信息。
2.5 Row group size 与内存缓冲
pyarrow parquet writer 会把每个 item(即每张表或 record batch)写入独立的 row group,这会产生大量小 row group,对某些查询引擎并不友好——例如 DuckDB 按 row group 做并行查询。dlt 的对策是:写入前对表和批次做缓冲与零拷贝拼接(zero-copy concat),通过控制缓冲中保留的最大行数来控制 row group 大小:
[data_writer] buffer_max_items=10e6注意 dlt 是把这些表缓存在内存中的,上例 1000 万行可能消耗大量 RAM。row_group_size配置在 pyarrow writer 下作用有限,只在写入单个超大 pyarrow 表、或内存缓冲非常大时才有用武之地。
当 source 产出 Arrow 表时,normalize 阶段还有专门的 ArrowItemsNormalizer:如果不需要加_dlt_id且文件已是 parquet,它会直接import_items_file零拷贝引入;否则逐 row group 流式重写(pq_stream_with_new_columns),在需要时调用normalize_py_arrow_item重排列、补齐缺失列。arrow_concat_promote_options="none"(默认)时要求 schema 完全一致以实现零拷贝拼接,这正是 2.2 中该参数的意义所在。
三、CSV:最基础但最有兼容性的格式
CSV 是存储表格数据最基础的格式,所有值都是字符串、以分隔符(通常是逗号)分隔。dlt 出于性能和兼容性考虑在特定场景使用它。内部根据数据 item 的形状选用两套实现:
- Python 标准库 CSV writer:resource 产出 Python 对象(dict)时使用;
- PyArrow CSV writer:高速多线程写入器,resource 产出 Arrow 表、pandas DataFrame 或 polars DataFrame 时使用。
3.1 统一的外观约定
dlt 尽量让两个 writer 产出外观一致的文件:
- 分隔符为逗号;
- 引号为
",转义为""; NULL值表现为空字符串或空 token,例如:
text1,text2,text3 A,B,C A,,""最后一行的text2和text3都是 NULL。由于 Pythoncsvwriter 无法写出未加引号的None,最终约定为"";
- 默认使用 UNIX 换行符
"\n"; - 日期以 ISO 8601 表示;
- quoting 风格为"需要时才加引号"。
所有支持csv格式的目标端都接受按上述标准设置写出的文件。
3.2 Write settings(normalize 阶段写入设置)
以下设置在[normalize.data_writer]段配置,控制 dlt 在 normalize 阶段如何写 CSV,对filesystem目标端场景尤其有用。其他目标端都用标准设置做了测试。默认值与 CsvFormatConfiguration 一致:
delimiter:分隔字符(默认',');include_header:是否写入表头行(默认True);lineterminator:行终止字符串(默认"\n",Windows 换行用"\r\n")。仅对 Python CSV writer 生效;PyArrow writer 恒用"\n";encoding:写文件用的编码(默认utf-8)。可用utf-8-sig为旧版 Excel 加 BOM,或latin-1/cp1252适配遗留导入器。两个 writer 都遵循;encoding_errors:无法用encoding表示的字符如何处理(默认strict——加载失败)。取 Python 错误处理器名,如replace(替换为?)或backslashreplace(保留为转义序列);quoting:何时给字段值加引号:quote_needed(默认):只给需要的值加引号(非数值)。Python CSV writer 会给所有非数值加引号;PyArrow CSV writer 的行为没有完全文档化,观察到某些情况下字符串也不加引号;quote_all:所有值都加引号。两个 writer 均支持;quote_minimal:只给含特殊字符(分隔符、引号、行终止符)的字段加引号。仅 Python CSV writer 支持;quote_none:永不加引号。Python CSV writer 在数据含分隔符时使用转义字符;PyArrow CSV writer 遇到特殊字符直接报错。
编码参数的健壮性在源码中也有保障:canonical_encoding / validated_encoding_errors 会先用codecs.lookup校验编码名与错误处理器名,非法值抛出InvalidEncoding/InvalidEncodingErrors;当目标编码不是 UTF-8 时,Utf8TranscodingWrapper 会把 writer 写出的 UTF-8 字节流增量转码为目标编码,且对 PyArrow 这种二进制 writer 同样有效。
[normalize.data_writer] delimiter="|" include_header=false quoting="quote_all" lineterminator="\r\n" encoding="latin-1"或环境变量:
NORMALIZE__DATA_WRITER__DELIMITER=| NORMALIZE__DATA_WRITER__INCLUDE_HEADER=False NORMALIZE__DATA_WRITER__QUOTING=quote_all NORMALIZE__DATA_WRITER__LINETERMINATOR=$"\r\n" NORMALIZE__DATA_WRITER__ENCODING=latin-1注意环境变量里"\r\n"前的$前缀,用于转义换行符。
3.3 Read settings(目标端读取设置)
把 CSV 文件复制进表的目标端(postgres和snowflake)按各自的csv_format配置来读取文件。这些设置不改变 dlt 的写文件方式,它们描述的是目标端将要加载的文件,配置在目标端上:
[destination.postgres.csv_format] delimiter="|" encoding="latin-1"读取时encoding告诉目标端如何解码 CSV 文件(默认utf-8),且有一个只在读取时使用的选项:
on_error_continue:跳过出错的行(仅 Snowflake 支持)。
csv_format同时接受上面写设置中的各选项——当被加载的文件偏离默认(比如不同分隔符、无表头)时设置它们。
写设置与读设置的职责划分(这是最容易被配错的地方):
- dlt 既写又读(标准流程,如加载到 postgres/snowflake):dlt 写出文件、目标端立刻读回,文件只是内部传输格式。两侧保持默认即可;
- 外部系统读取文件(
filesystem目标端为最终落点):文件就是产物。把[normalize.data_writer]调成消费者期望的样子,例如旧版 Excel 用encoding="utf-8-sig"、遗留导入器用cp1252; - 文件不是 dlt 写的(importing external files):在
[destination.<name>.csv_format]里描述文件。注意encoding会原样进入目标端的 COPY 语句,必须是目标端接受的编码名——这些名字并不总是和 Python 的一致(如latin-1vsISO_8859_1)。
如果你坚持组合自定义写编码与数据库目标端,请在目标端csv_format中镜像同样的值——否则目标端按utf-8解码文件,加载要么失败、要么非 ASCII 字符乱码。
3.4 局限性
arrow writer(PyArrow CSV writer)
- binary 列仅支持包含有效 UTF-8 字符的情况;
- json(嵌套/struct)类型不受支持。
csv writer(Python 标准库)
- binary 列仅支持有效 UTF-8 字符(容易扩展更多编码);
- json 列用
json.dumps转储; None值始终带引号。
四、JSONL:行分隔 JSON 文档
JSONL(JSON Lines / JSON Delimited)在一个文件中存储多个 JSON 文档,文档之间以换行分隔。dlt 的 JsonlWriter 是 Python dict 类数据的默认格式之一。
额外数据类型的存储规则:
datetime与date存为 ISO 字符串。tz-aware 时间戳携带数字偏移量,UTC 写作2024-01-15T23:30:00+00:00而非2024-01-15T23:30:00Z。两种写法都能读回,因此旧版 dlt 写出的文件依然可以加载;decimal存为十进制数的文本表示;binary存为 base64 编码字符串;HexBytes存为 hex 编码字符串;json序列化为字符串。
该格式默认压缩。
五、SQL INSERT:直接生成 INSERT...VALUES 语句
该格式生成的文件内容是将在load阶段于目标端执行的INSERT...VALUES语句。实现见 InsertValuesWriter:
- writer 会根据目标端能力上下文中的
insert_values_writer_type选择输出风格:default生成(values),(values)形式,select_union生成SELECT ... UNION ALL形式; - 列名通过
escape_identifier转义,值通过escape_literal转义——即值格式由目标端能力决定; - 该格式要求目标端能力上下文存在(源码中
assert caps is not None),这也是它不能用于文件系统类目标端的原因。
额外数据类型的存储规则:
datetime与date存为 ISO 字符串;decimal存为十进制数文本表示;binary的存储取决于目标端接受的格式;json的存储同样取决于目标端接受的格式。
该格式默认压缩。
六、选型小结
| 格式 | 典型适用场景 | 关键依赖 | 压缩 |
|---|---|---|---|
| Parquet | 列式分析、大数据量、DuckDB/Spark 等引擎并行读取 | pyarrow(pip install "dlt[parquet]") | 由compression参数控制的内部 codec |
| CSV | 与遗留系统/Excel 互操作、filesystem落点 | 无额外依赖(Arrow 表路径需 pyarrow) | 默认 gzip 外层压缩,可用disable_compression关闭 |
| JSONL | 需要保留 dlt 全部数据类型(decimal、binary、HexBytes、嵌套 json)的通用默认格式 | 无 | 默认开启 |
| SQL INSERT | 目标端只能执行 SQL 语句、不支持文件导入 | 目标端能力上下文(escape_literal/escape_identifier) | 默认开启 |
实践建议:优先使用目标端preferred格式,仅在消费者(外部数据库、Excel、遗留导入器)有明确诉求时才在[normalize.data_writer]或资源级file_format上覆盖;凡涉及编码的分流——谁写文件调写设置,谁读文件调目标端csv_format——这一原则同样适用于理解 dlt 其他 loader 格式的边界。
(本文技术细节以当前仓库代码为准;涉及 pyarrow 版本的能力差异,如use_content_defined_chunking要求pyarrow>=21.0.0,请在实际环境中以所安装版本为准。)
【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy 🛠️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考