Parquet 内容定义分块
Parquet 内容定义分块
内容定义分块(Content-Defined Chunking,CDC)是 Parquet 写入器的一项实验性功能,它使数据页边界取决于列值,而不是固定的行数或字节数。这样一来,当使用相同设置写入彼此关联紧密的数据集版本时,未发生变化的区域更有可能产生完全相同的页。
当生成的文件通过可寻址内容或块去重系统进行存储或传输时,CDC 就能发挥作用。这类系统可以复用完全相同的页,而不必再次存储或传输它们。例如,在数据集开头附近的一处小插入可能只改变一个页,而其后的页边界会重新收敛到与上一版本一致。
CDC 本身并不进行数据去重,也不提供页存储。在传统文件系统或对象存储上,每个 Parquet 文件仍然会被完整存储。输出的仍是普通的 Parquet 文件,不需要读取端提供任何 CDC 相关的支持。
何时启用 CDC
当以下条件全部满足时,可以考虑启用 CDC:
- 你会经常写出同一数据集的相似版本。
- 你的存储或传输层能够检测并复用重复的字节范围。
- 相比提升单个文件的写入并行度,减少存储或网络传输开销更重要。
对于普通的 Parquet 输出,默认应保持 CDC 关闭,除非你已在存储或传输这些文件的系统中实际测量到了收益。CDC 默认处于关闭状态。
启用 CDC 后,DataFusion 会为每个输出文件使用顺序 Arrow 写入器,因为分块器的状态必须在行组之间保持。这可能会降低写入吞吐量,相比 DataFusion 的并行写入路径而言。不过写入不同的输出文件仍然可以并发进行。
CDC 对每个输出文件独立运作。当 COPY 的目标为目录时,DataFusion 会以轮询方式将输入的 RecordBatch 分配到各个并行输出文件中;datafusion.execution.minimum_parallel_output_files 默认值为 4。如果数据集版本之间的分批或文件分配发生变化,未改动的行可能会在文件之间移动,从而降低去重效果。要获得最佳结果,请保持输入顺序和输出文件布局稳定。单文件输出请使用文件名作为目标;需要多个文件时,请按稳定键进行分区。参见配置设置。
通过 SQL 启用 CDC
通过 Parquet 格式选项,可以为单次 COPY 操作设置 CDC。本例中的文件名目标会选择单文件输出:
COPY (
SELECT
value AS id,
CONCAT('event-', CAST(value AS VARCHAR)) AS event
FROM generate_series(1, 100000)
) TO 'cdc-output.parquet'
STORED AS PARQUET
OPTIONS (
'format.content_defined_chunking.enabled' 'true'
);默认分块参数是不错的起点。下面的示例在一次写入中显式指定了这些默认值,并未修改它们的取值:
COPY source_table TO 'cdc-output.parquet'
STORED AS PARQUET
OPTIONS (
'format.content_defined_chunking.enabled' 'true',
'format.content_defined_chunking.min_chunk_size' '262144',
'format.content_defined_chunking.max_chunk_size' '1048576',
'format.content_defined_chunking.norm_level' '0'
);请在使用具有代表性的数据进行测量之后,再修改这些值。
也可以在会话中为后续的 Parquet 写入操作启用 CDC:
SET datafusion.execution.parquet.content_defined_chunking.enabled = true;对应的环境变量为 DATAFUSION_EXECUTION_PARQUET_CONTENT_DEFINED_CHUNKING_ENABLED。有关设置会话选项的所有方式,请参阅配置设置。
通过 Rust API 启用 CDC
将 TableParquetOptions 传入 DataFrame::write_parquet:
use datafusion::config::{ParquetCdcOptions, TableParquetOptions};
use datafusion::dataframe::DataFrameWriteOptions;
use datafusion::error::Result;
use datafusion::prelude::SessionContext;
#[tokio::main]
async fn main() -> Result<()> {
let ctx = SessionContext::new();
let df = ctx
.sql("SELECT value AS id FROM generate_series(1, 100000)")
.await?;
let mut parquet_options = TableParquetOptions::default();
parquet_options.global.content_defined_chunking = ParquetCdcOptions::enabled();
df.write_parquet(
"cdc-output.parquet",
DataFrameWriteOptions::new().with_single_file_output(true),
Some(parquet_options),
)
.await?;
Ok(())
}直接设置 ParquetCdcOptions 的字段即可使用非默认的分块大小或归一化设置。
调优
| 选项 | 默认值 | 作用 |
|---|---|---|
min_chunk_size | 256 KiB | 滚动哈希可以选择边界的最小逻辑大小。 |
max_chunk_size | 1 MiB | 写入方强制创建边界的最大逻辑大小。该值必须大于 min_chunk_size。 |
norm_level | 0 | 控制边界选择的激进程度。较高的值可以提升去重效果,但会产生更多小页;建议取值范围为 -3 到 3。 |
分块大小是在编码和压缩之前,根据列的逻辑数据进行度量的。嵌套数据的定义层级(definition level)和重复层级(repetition level)也会计入大小。
在比较数据集的不同版本时,应使用相同的 CDC、编码、压缩和模式(schema)设置。即使逻辑数据没有变化,修改写入方设置也可能改变页字节内容,从而降低去重效果。在修改默认值之前,请先使用有代表性的数据测量去重率、输出大小、网络传输量和写入时间。
评论
登录后参与评论
KnowForge