Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 35 additions & 0 deletions src/include/storage/ducklake_stats_copy.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
//===----------------------------------------------------------------------===//
// DuckDB
//
// storage/ducklake_stats_copy.hpp
//
//
//===----------------------------------------------------------------------===//

#pragma once

#include "duckdb/function/copy_function.hpp"

namespace duckdb {

//! Wrap a sink-based COPY function so per-file column statistics (min/max,
//! null_count, num_values, has_nan) are computed from the chunks as they
//! stream through the sink — the same point where the parquet writer computes
//! its statistics — and reported through copy_to_get_written_statistics /
//! CopyFunctionReturnType::WRITTEN_FILE_STATISTICS.
//!
//! For COPY functions that do not report written statistics themselves (the
//! vortex COPY returns only a file list). Costs one comparison pass over each
//! chunk at write time; no re-read of the written file. Single-file sink-based
//! writes only: the wrapped function must not use the batch or rotation APIs.
struct DuckLakeStatsCopy {
//! Whether `fn` needs wrapping to report written statistics
static bool NeedsWrapping(const CopyFunction &fn);
//! The wrapping COPY function forwarding to `inner`
static CopyFunction WrapFunction(const CopyFunction &inner);
//! Wrap `inner`'s bound data; `names`/`types` describe the sunk chunks
static unique_ptr<FunctionData> WrapBindData(const CopyFunction &inner, unique_ptr<FunctionData> inner_bind,
vector<string> names, vector<LogicalType> types);
};

} // namespace duckdb
1 change: 1 addition & 0 deletions src/storage/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ add_library(
ducklake_partition_data.cpp
ducklake_secret.cpp
ducklake_stats.cpp
ducklake_stats_copy.cpp
ducklake_table_entry.cpp
ducklake_initializer.cpp
ducklake_autoload_helper.cpp
Expand Down
74 changes: 29 additions & 45 deletions src/storage/ducklake_insert.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
#include "duckdb/planner/expression/bound_cast_expression.hpp"
#include "duckdb/planner/expression/bound_reference_expression.hpp"
#include "duckdb/common/multi_file/multi_file_reader.hpp"
#include "storage/ducklake_stats_copy.hpp"
#include "duckdb/main/extension_helper.hpp"
#include "duckdb/function/function_binder.hpp"
#include "duckdb/planner/operator/logical_projection.hpp"
Expand Down Expand Up @@ -108,46 +109,21 @@ DuckLakeColumnStats DuckLakeInsert::ParseColumnStats(const LogicalType &type, co

void DuckLakeInsert::AddWrittenFiles(ClientContext &context, DuckLakeInsertGlobalState &global_state, DataChunk &chunk,
const string &encryption_key, optional_idx partition_id, bool set_snapshot_id) {
if (global_state.file_format != "parquet") {
// non-parquet writers return CHANGED_ROWS_AND_FILE_LIST: {count, files[]}. There are no per-file
// stats, so derive the size from the filesystem and map the whole row count to the single file.
// (Combinations that need per-file stats - partitioning, flush, NOT NULL - are rejected in
// GetCopyOptions before anything is written.)
auto &fs = FileSystem::GetFileSystem(context);
for (idx_t r = 0; r < chunk.size(); r++) {
auto files_value = chunk.GetValue(1, r);
if (files_value.IsNull()) {
continue;
}
auto &file_list = ListValue::GetChildren(files_value);
if (file_list.empty()) {
continue;
}
if (file_list.size() != 1) {
throw NotImplementedException("Writing '%s' data files is only supported as a single file per write",
global_state.file_format);
}
DuckLakeDataFile data_file;
data_file.file_name = file_list[0].GetValue<string>();
data_file.file_format = global_state.file_format;
data_file.row_count = chunk.GetValue(0, r).GetValue<idx_t>();
auto handle = fs.OpenFile(data_file.file_name, FileFlags::FILE_FLAGS_READ);
data_file.file_size_bytes = NumericCast<idx_t>(fs.GetFileSize(*handle));
data_file.encryption_key = encryption_key;
if (partition_id.IsValid()) {
data_file.partition_id = partition_id.GetIndex();
}
global_state.written_files.push_back(std::move(data_file));
}
return;
}
// Both parquet and wrapped non-parquet (vortex) COPYs return
// WRITTEN_FILE_STATISTICS ({filename, row_count, file_size, footer_size,
// column_stats, partition_keys}); the non-parquet path is single-file and
// non-partitioned but consumes the same chunk shape. (Non-parquet functions
// report no footer size - see the NULL tolerance below.)
for (idx_t r = 0; r < chunk.size(); r++) {
DuckLakeDataFile data_file;
data_file.file_name = chunk.GetValue(0, r).GetValue<string>();
data_file.file_format = global_state.file_format;
data_file.row_count = chunk.GetValue(1, r).GetValue<idx_t>();
data_file.file_size_bytes = chunk.GetValue(2, r).GetValue<idx_t>();
data_file.footer_size = chunk.GetValue(3, r).GetValue<idx_t>();
auto footer_size = chunk.GetValue(3, r);
if (!footer_size.IsNull()) {
data_file.footer_size = footer_size.GetValue<idx_t>();
}
data_file.encryption_key = encryption_key;
if (partition_id.IsValid()) {
data_file.partition_id = partition_id.GetIndex();
Expand Down Expand Up @@ -562,11 +538,6 @@ DuckLakeCopyOptions DuckLakeInsert::GetCopyOptions(ClientContext &context, DuckL
"set data_inlining_row_limit to 0 for such tables",
data_file_format);
}
if (copy_input.has_not_null_columns) {
// NOT NULL is enforced by inspecting written null-count statistics
throw NotImplementedException("NOT NULL columns are not yet supported for the '%s' data file format",
data_file_format);
}
if (copy_input.is_compaction) {
// compaction's directory + rotation output does not compose with the single-file path used
// for non-parquet writes yet
Expand Down Expand Up @@ -674,7 +645,17 @@ DuckLakeCopyOptions DuckLakeInsert::GetCopyOptions(ClientContext &context, DuckL

auto function_data = copy_fun.function.copy_to_bind(context, bind_input, names_to_write, casted_types);

DuckLakeCopyOptions result(std::move(info), copy_fun.function);
// COPY functions without a written-statistics hook (vortex) are wrapped so
// the per-file column statistics are computed from the chunks as they
// stream through the sink and reported like parquet's.
auto copy_function = copy_fun.function;
if (!is_parquet && DuckLakeStatsCopy::NeedsWrapping(copy_function)) {
function_data =
DuckLakeStatsCopy::WrapBindData(copy_function, std::move(function_data), names_to_write, casted_types);
copy_function = DuckLakeStatsCopy::WrapFunction(copy_function);
}

DuckLakeCopyOptions result(std::move(info), std::move(copy_function));
result.bind_data = std::move(function_data);

result.use_tmp_file = false;
Expand All @@ -700,11 +681,14 @@ DuckLakeCopyOptions DuckLakeInsert::GetCopyOptions(ClientContext &context, DuckL
if (is_parquet) {
result.return_type = CopyFunctionReturnType::WRITTEN_FILE_STATISTICS;
} else {
// non-parquet COPY functions (vortex) can't report per-file statistics; take the written file
// list instead and derive size/row-count ourselves. They also don't implement rotate_next_file,
// so we can't use DuckDB's directory+rotation naming - generate a single full file path here
// (the plain single-file COPY path) so the writer receives a file, not the table directory.
result.return_type = CopyFunctionReturnType::CHANGED_ROWS_AND_FILE_LIST;
// Non-parquet COPY functions (vortex) don't report per-file statistics
// themselves; the DuckLakeStatsCopy wrapper (applied below) computes
// them from the sunk chunks and reports WRITTEN_FILE_STATISTICS like
// parquet. They also don't implement rotate_next_file, so we can't use
// DuckDB's directory+rotation naming - generate a single full file path
// here (the plain single-file COPY path) so the writer receives a file,
// not the table directory.
result.return_type = CopyFunctionReturnType::WRITTEN_FILE_STATISTICS;
result.rotate = false;
result.per_thread_output = false;
result.partition_output = false;
Expand Down
Loading
Loading