From 8c47c26d33dfc5d4b755d107f02478f67b7aa7f0 Mon Sep 17 00:00:00 2001 From: Mosha Pasumansky Date: Fri, 4 Sep 2026 10:43:21 -0700 Subject: [PATCH] fix(vortex): compute per-file column statistics in-stream for vortex writes The vortex COPY reports only a file list, so vortex data files were registered without per-file column statistics and the global table stats were never updated. Readers trust those stats, which turned the stale bounds into wrong results on mixed parquet+vortex tables: - filters on values outside the stale bounds folded to empty (WHERE id = 300 -> no rows, count(*) WHERE id >= 250 -> 0) - DuckDB's sort-key compression packed ORDER BY columns into the stale range's width, returning values truncated mod 256 for ORDER BY..LIMIT Fix: DuckLakeStatsCopy wraps any sink-based COPY function that lacks a copy_to_get_written_statistics hook. The wrapper accumulates per-column min/max (arg-min/max raw scan per chunk, one Value materialization per chunk), null/value counts, and NaN detection for FLOAT/DOUBLE (NaN is excluded from bounds and flagged, parquet-style) as the chunks stream through the sink - the same point where the parquet writer computes its statistics - and reports them through WRITTEN_FILE_STATISTICS. No extra pass over the written file; no extension change needed. The insert path then consumes one unified WRITTEN_FILE_STATISTICS chunk shape for every format (the non-parquet file-list special case is gone; footer_size tolerates NULL since only parquet reports one). Vortex files therefore get real ducklake_file_column_stats rows (file pruning now works on them), the global bounds stay correct, and NOT NULL constraints are enforced from the written null counts - the previous NOT NULL rejection for non-parquet formats is lifted. Bounds are recorded for types whose physical order matches logical order and whose values round-trip through the stats VARCHAR form; other columns (nested, blob, uuid, interval, oversized strings) record counts only. This also unblocks requesting WRITTEN_FILE_STATISTICS for vortex, which previously crashed on the missing extension hook (feat/rtdl-vortex-written-stats); if the duckdb-vortex extension later implements the hook natively, the wrapper steps aside automatically (NeedsWrapping). Co-Authored-By: Claude Fable 5 --- src/include/storage/ducklake_stats_copy.hpp | 35 ++ src/storage/CMakeLists.txt | 1 + src/storage/ducklake_insert.cpp | 74 ++-- src/storage/ducklake_stats_copy.cpp | 411 ++++++++++++++++++ .../sql/vortex/vortex_mixed_format_stats.test | 98 +++++ test/sql/vortex/vortex_write.test | 10 +- .../vortex/vortex_written_stats_types.test | 357 +++++++++++++++ 7 files changed, 938 insertions(+), 48 deletions(-) create mode 100644 src/include/storage/ducklake_stats_copy.hpp create mode 100644 src/storage/ducklake_stats_copy.cpp create mode 100644 test/sql/vortex/vortex_mixed_format_stats.test create mode 100644 test/sql/vortex/vortex_written_stats_types.test diff --git a/src/include/storage/ducklake_stats_copy.hpp b/src/include/storage/ducklake_stats_copy.hpp new file mode 100644 index 00000000..934f9edd --- /dev/null +++ b/src/include/storage/ducklake_stats_copy.hpp @@ -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 WrapBindData(const CopyFunction &inner, unique_ptr inner_bind, + vector names, vector types); +}; + +} // namespace duckdb diff --git a/src/storage/CMakeLists.txt b/src/storage/CMakeLists.txt index c31596ca..e78741e8 100644 --- a/src/storage/CMakeLists.txt +++ b/src/storage/CMakeLists.txt @@ -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 diff --git a/src/storage/ducklake_insert.cpp b/src/storage/ducklake_insert.cpp index 08bfdce9..0386a8f9 100644 --- a/src/storage/ducklake_insert.cpp +++ b/src/storage/ducklake_insert.cpp @@ -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" @@ -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(); - data_file.file_format = global_state.file_format; - data_file.row_count = chunk.GetValue(0, r).GetValue(); - auto handle = fs.OpenFile(data_file.file_name, FileFlags::FILE_FLAGS_READ); - data_file.file_size_bytes = NumericCast(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(); data_file.file_format = global_state.file_format; data_file.row_count = chunk.GetValue(1, r).GetValue(); data_file.file_size_bytes = chunk.GetValue(2, r).GetValue(); - data_file.footer_size = chunk.GetValue(3, r).GetValue(); + auto footer_size = chunk.GetValue(3, r); + if (!footer_size.IsNull()) { + data_file.footer_size = footer_size.GetValue(); + } data_file.encryption_key = encryption_key; if (partition_id.IsValid()) { data_file.partition_id = partition_id.GetIndex(); @@ -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 @@ -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; @@ -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; diff --git a/src/storage/ducklake_stats_copy.cpp b/src/storage/ducklake_stats_copy.cpp new file mode 100644 index 00000000..92134b6e --- /dev/null +++ b/src/storage/ducklake_stats_copy.cpp @@ -0,0 +1,411 @@ +#include "storage/ducklake_stats_copy.hpp" + +#include "common/ducklake_util.hpp" +#include "duckdb/common/operator/comparison_operators.hpp" +#include "duckdb/common/types/value.hpp" +#include "duckdb/common/vector_operations/vector_operations.hpp" +#include "duckdb/common/file_system.hpp" + +#include + +namespace duckdb { + +namespace { + +//! Bounds are only recorded for types whose physical-value order matches their +//! logical order and whose values round-trip through the stats VARCHAR +//! representation. Everything else still gets null/value counts. +bool TypeSupportsMinMaxStats(const LogicalType &type) { + switch (type.id()) { + case LogicalTypeId::BOOLEAN: + case LogicalTypeId::TINYINT: + case LogicalTypeId::SMALLINT: + case LogicalTypeId::INTEGER: + case LogicalTypeId::BIGINT: + case LogicalTypeId::HUGEINT: + case LogicalTypeId::UTINYINT: + case LogicalTypeId::USMALLINT: + case LogicalTypeId::UINTEGER: + case LogicalTypeId::UBIGINT: + case LogicalTypeId::UHUGEINT: + case LogicalTypeId::FLOAT: + case LogicalTypeId::DOUBLE: + case LogicalTypeId::DECIMAL: + case LogicalTypeId::DATE: + case LogicalTypeId::TIME: + case LogicalTypeId::TIMESTAMP: + case LogicalTypeId::TIMESTAMP_SEC: + case LogicalTypeId::TIMESTAMP_MS: + case LogicalTypeId::TIMESTAMP_NS: + case LogicalTypeId::TIMESTAMP_TZ: + case LogicalTypeId::VARCHAR: + return true; + default: + return false; + } +} + +//! VARCHAR bounds above this length are dropped (a truncated max is not an +//! upper bound; rather than truncate-and-increment like parquet, drop bounds). +constexpr idx_t MAX_STRING_STATS_LENGTH = 1024; + +template +bool ValueIsNan(T) { + return false; +} +template <> +bool ValueIsNan(float v) { + return std::isnan(v); +} +template <> +bool ValueIsNan(double v) { + return std::isnan(v); +} + +//! Per-column statistics accumulated from sunk chunks. +struct StatsAccumulator { + explicit StatsAccumulator(LogicalType type_p) + : type(std::move(type_p)), minmax_supported(TypeSupportsMinMaxStats(type)), + is_float(type.id() == LogicalTypeId::FLOAT || type.id() == LogicalTypeId::DOUBLE) { + } + + LogicalType type; + bool minmax_supported; + bool is_float; + //! Cleared when a bound cannot be represented (oversized string) + bool minmax_valid = true; + idx_t null_count = 0; + idx_t num_values = 0; + bool contains_nan = false; + //! NULL Values until the first non-null (non-NaN) value is seen + Value min_val; + Value max_val; + + void Update(Vector &vec, idx_t count); + void Merge(const StatsAccumulator &other); + +private: + void UpdateBound(const Value &candidate, bool is_min); + template + void TemplatedArgMinMax(const UnifiedVectorFormat &vdata, idx_t count, optional_idx &argmin, optional_idx &argmax); +}; + +template +void StatsAccumulator::TemplatedArgMinMax(const UnifiedVectorFormat &vdata, idx_t count, optional_idx &argmin, + optional_idx &argmax) { + auto data = UnifiedVectorFormat::GetData(vdata); + optional_idx min_row, max_row; + T min_v {}; + T max_v {}; + for (idx_t i = 0; i < count; i++) { + auto idx = vdata.sel->get_index(i); + if (!vdata.validity.RowIsValid(idx)) { + continue; + } + T v = data[idx]; + if (ValueIsNan(v)) { + contains_nan = true; + continue; + } + if (!min_row.IsValid() || LessThan::Operation(v, min_v)) { + min_v = v; + min_row = i; + } + if (!max_row.IsValid() || GreaterThan::Operation(v, max_v)) { + max_v = v; + max_row = i; + } + } + argmin = min_row; + argmax = max_row; +} + +void StatsAccumulator::UpdateBound(const Value &candidate, bool is_min) { + if (candidate.type().id() == LogicalTypeId::VARCHAR && + StringValue::Get(candidate).size() > MAX_STRING_STATS_LENGTH) { + minmax_valid = false; + min_val = Value(); + max_val = Value(); + return; + } + auto &bound = is_min ? min_val : max_val; + if (bound.IsNull()) { + bound = candidate; + return; + } + if (is_min ? candidate < bound : candidate > bound) { + bound = candidate; + } +} + +void StatsAccumulator::Update(Vector &vec, idx_t count) { + num_values += count; + + UnifiedVectorFormat vdata; + vec.ToUnifiedFormat(count, vdata); + if (!vdata.validity.AllValid()) { + for (idx_t i = 0; i < count; i++) { + if (!vdata.validity.RowIsValid(vdata.sel->get_index(i))) { + null_count++; + } + } + } + if (!minmax_supported || !minmax_valid) { + if (!is_float) { + return; + } + } + + optional_idx argmin, argmax; + switch (vec.GetType().InternalType()) { + case PhysicalType::BOOL: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::INT8: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::INT16: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::INT32: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::INT64: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::INT128: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::UINT8: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::UINT16: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::UINT32: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::UINT64: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::UINT128: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::FLOAT: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::DOUBLE: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + case PhysicalType::VARCHAR: + TemplatedArgMinMax(vdata, count, argmin, argmax); + break; + default: + // unreachable given TypeSupportsMinMaxStats, except float NaN scans + return; + } + if (!minmax_supported || !minmax_valid) { + // float column that only needed the NaN scan + return; + } + if (argmin.IsValid()) { + UpdateBound(vec.GetValue(argmin.GetIndex()), true); + } + if (argmax.IsValid() && minmax_valid) { + UpdateBound(vec.GetValue(argmax.GetIndex()), false); + } +} + +void StatsAccumulator::Merge(const StatsAccumulator &other) { + null_count += other.null_count; + num_values += other.num_values; + contains_nan = contains_nan || other.contains_nan; + if (!other.minmax_valid) { + minmax_valid = false; + min_val = Value(); + max_val = Value(); + } + if (!minmax_valid) { + return; + } + if (!other.min_val.IsNull()) { + UpdateBound(other.min_val, true); + } + if (!other.max_val.IsNull() && minmax_valid) { + UpdateBound(other.max_val, false); + } +} + +struct StatsCopyBindData : public FunctionData { + CopyFunction inner_function; + unique_ptr inner; + vector names; + vector types; + + unique_ptr Copy() const override { + auto result = make_uniq(); + result->inner_function = inner_function; + result->inner = inner->Copy(); + result->names = names; + result->types = types; + return std::move(result); + } + bool Equals(const FunctionData &other_p) const override { + auto &other = other_p.Cast(); + return names == other.names && types == other.types && inner->Equals(*other.inner); + } + + StatsCopyBindData() : inner_function("ducklake_stats_copy_wrapper") { + } +}; + +struct StatsCopyGlobalState : public GlobalFunctionData { + unique_ptr inner; + string file_path; + mutex merge_lock; + vector stats; + idx_t row_count = 0; + optional_ptr target; +}; + +struct StatsCopyLocalState : public LocalFunctionData { + unique_ptr inner; + vector stats; + idx_t row_count = 0; +}; + +vector MakeAccumulators(const vector &types) { + vector result; + result.reserve(types.size()); + for (auto &type : types) { + result.emplace_back(type); + } + return result; +} + +unique_ptr StatsCopyInitializeGlobal(ClientContext &context, FunctionData &bind_data_p, + const string &file_path) { + auto &bind_data = bind_data_p.Cast(); + auto result = make_uniq(); + result->inner = bind_data.inner_function.copy_to_initialize_global(context, *bind_data.inner, file_path); + result->file_path = file_path; + result->stats = MakeAccumulators(bind_data.types); + return std::move(result); +} + +unique_ptr StatsCopyInitializeLocal(ExecutionContext &context, FunctionData &bind_data_p) { + auto &bind_data = bind_data_p.Cast(); + auto result = make_uniq(); + result->inner = bind_data.inner_function.copy_to_initialize_local(context, *bind_data.inner); + result->stats = MakeAccumulators(bind_data.types); + return std::move(result); +} + +void StatsCopyGetWrittenStatistics(ClientContext &context, FunctionData &bind_data, GlobalFunctionData &gstate, + CopyFunctionFileStatistics &statistics) { + auto &global_state = gstate.Cast(); + global_state.target = &statistics; +} + +void StatsCopySink(ExecutionContext &context, FunctionData &bind_data_p, GlobalFunctionData &gstate, + LocalFunctionData &lstate, DataChunk &input) { + auto &bind_data = bind_data_p.Cast(); + auto &local_state = lstate.Cast(); + auto &global_state = gstate.Cast(); + for (idx_t c = 0; c < input.ColumnCount(); c++) { + local_state.stats[c].Update(input.data[c], input.size()); + } + local_state.row_count += input.size(); + bind_data.inner_function.copy_to_sink(context, *bind_data.inner, *global_state.inner, *local_state.inner, input); +} + +void StatsCopyCombine(ExecutionContext &context, FunctionData &bind_data_p, GlobalFunctionData &gstate, + LocalFunctionData &lstate) { + auto &bind_data = bind_data_p.Cast(); + auto &local_state = lstate.Cast(); + auto &global_state = gstate.Cast(); + { + lock_guard guard(global_state.merge_lock); + for (idx_t c = 0; c < global_state.stats.size(); c++) { + global_state.stats[c].Merge(local_state.stats[c]); + } + global_state.row_count += local_state.row_count; + } + if (bind_data.inner_function.copy_to_combine) { + bind_data.inner_function.copy_to_combine(context, *bind_data.inner, *global_state.inner, *local_state.inner); + } +} + +void StatsCopyFinalize(ClientContext &context, FunctionData &bind_data_p, GlobalFunctionData &gstate) { + auto &bind_data = bind_data_p.Cast(); + auto &global_state = gstate.Cast(); + if (bind_data.inner_function.copy_to_finalize) { + bind_data.inner_function.copy_to_finalize(context, *bind_data.inner, *global_state.inner); + } + if (!global_state.target) { + return; + } + auto &statistics = *global_state.target; + statistics.row_count = global_state.row_count; + auto &fs = FileSystem::GetFileSystem(context); + auto handle = fs.OpenFile(global_state.file_path, FileFlags::FILE_FLAGS_READ); + statistics.file_size_bytes = NumericCast(fs.GetFileSize(*handle)); + // footer_size_bytes stays NULL: the format's footer size is not surfaced here + + for (idx_t c = 0; c < global_state.stats.size(); c++) { + auto &acc = global_state.stats[c]; + case_insensitive_map_t column_stats; + column_stats["null_count"] = Value::UBIGINT(acc.null_count); + column_stats["num_values"] = Value::UBIGINT(acc.num_values); + if (acc.minmax_valid && !acc.min_val.IsNull()) { + column_stats["min"] = acc.min_val; + } + if (acc.minmax_valid && !acc.max_val.IsNull()) { + column_stats["max"] = acc.max_val; + } + if (acc.is_float) { + column_stats["has_nan"] = Value::BOOLEAN(acc.contains_nan); + } + vector name_path {bind_data.names[c]}; + statistics.column_statistics[DuckLakeUtil::ToQuotedList(name_path)] = std::move(column_stats); + } +} + +} // namespace + +bool DuckLakeStatsCopy::NeedsWrapping(const CopyFunction &fn) { + return !fn.copy_to_get_written_statistics; +} + +CopyFunction DuckLakeStatsCopy::WrapFunction(const CopyFunction &inner) { + if (inner.prepare_batch || inner.flush_batch || inner.rotate_files || inner.rotate_next_file) { + throw NotImplementedException( + "DuckLakeStatsCopy only supports single-file sink-based COPY functions"); + } + D_ASSERT(NeedsWrapping(inner)); + CopyFunction result(inner.name + "_with_stats"); + // binding happens on the INNER function; the caller wraps the bound data + // through WrapBindData, so the wrapper needs no copy_to_bind of its own + result.copy_to_initialize_global = StatsCopyInitializeGlobal; + result.copy_to_initialize_local = StatsCopyInitializeLocal; + result.copy_to_get_written_statistics = StatsCopyGetWrittenStatistics; + result.copy_to_sink = StatsCopySink; + result.copy_to_combine = StatsCopyCombine; + result.copy_to_finalize = StatsCopyFinalize; + result.execution_mode = inner.execution_mode; + result.initialize_operator = inner.initialize_operator; + return result; +} + +unique_ptr DuckLakeStatsCopy::WrapBindData(const CopyFunction &inner, + unique_ptr inner_bind, vector names, + vector types) { + auto result = make_uniq(); + result->inner_function = inner; + result->inner = std::move(inner_bind); + result->names = std::move(names); + result->types = std::move(types); + return std::move(result); +} + +} // namespace duckdb diff --git a/test/sql/vortex/vortex_mixed_format_stats.test b/test/sql/vortex/vortex_mixed_format_stats.test new file mode 100644 index 00000000..e31c2f1c --- /dev/null +++ b/test/sql/vortex/vortex_mixed_format_stats.test @@ -0,0 +1,98 @@ +# name: test/sql/vortex/vortex_mixed_format_stats.test +# description: Vortex writes record real per-file column statistics. Regression: the vortex COPY +# reports no statistics, and appending such a file left parquet-only bounds in +# ducklake_table_column_stats; readers trust those bounds, so filters on +# out-of-bounds values folded to empty (WHERE id = 300 -> no rows, count(*) WHERE +# id >= 250 -> 0) and DuckDB's sort-key compression packed ORDER BY columns into the +# stale range's width, returning values truncated mod 256 for ORDER BY ... LIMIT. +# The insert path now computes the statistics from the written file itself. +# group: [vortex] + +require ducklake + +require parquet + +require vortex + +statement ok +ATTACH 'ducklake:{TEST_DIR}/dl_mixed_stats.db' AS ducklake (DATA_PATH '{TEST_DIR}/dl_mixed_stats_files', METADATA_CATALOG 'ducklake_meta') + +statement ok +CALL ducklake.set_option('data_inlining_row_limit', 0) + +statement ok +CREATE TABLE ducklake.m(id INTEGER, s VARCHAR) + +# parquet file with ids 0..249 +statement ok +INSERT INTO ducklake.m SELECT i::INTEGER, 'val_' || i FROM range(0, 250) t(i) + +query II +SELECT min_value, max_value FROM ducklake_meta.ducklake_table_column_stats WHERE column_id = 1 +---- +0 249 + +# vortex file with ids 250..499 and some NULL strings +statement ok +CALL ducklake.set_option('data_file_format', 'vortex') + +statement ok +INSERT INTO ducklake.m SELECT i::INTEGER, CASE WHEN i % 13 = 0 THEN NULL ELSE 'val_' || i END FROM range(250, 500) t(i) + +# the vortex file records real per-file column statistics, like parquet +query IIIII +SELECT s.column_id, s.min_value, s.max_value, s.null_count, s.value_count +FROM ducklake_meta.ducklake_file_column_stats s +JOIN ducklake_meta.ducklake_data_file f USING (data_file_id) +WHERE f.file_format = 'vortex' +ORDER BY s.column_id +---- +1 250 499 0 250 +2 val_250 val_499 19 231 + +# and the global bounds now cover the whole table +query III +SELECT min_value, max_value, contains_null FROM ducklake_meta.ducklake_table_column_stats WHERE column_id = 1 +---- +0 499 false + +query III +SELECT min_value, max_value, contains_null FROM ducklake_meta.ducklake_table_column_stats WHERE column_id = 2 +---- +val_0 val_99 true + +# both files scanned as one table: filters on values above the old bound must match +query II +SELECT id, s FROM ducklake.m WHERE id = 300 +---- +300 val_300 + +query I +SELECT count(*) FROM ducklake.m WHERE id >= 250 +---- +250 + +# ORDER BY + LIMIT must not truncate the sort column to stale bounds' width +# (this returned 243/242 = 499/498 mod 256 before the fix) +query II +SELECT id, s FROM ducklake.m ORDER BY id DESC LIMIT 2 +---- +499 val_499 +498 val_498 + +query II +SELECT min(id), max(id) FROM ducklake.m +---- +0 499 + +# NOT NULL is now enforced for vortex writes from the computed null counts +statement ok +CREATE TABLE ducklake.nn(id INTEGER NOT NULL, s VARCHAR) + +statement ok +INSERT INTO ducklake.nn VALUES (1, 'a'), (2, 'b') + +statement error +INSERT INTO ducklake.nn VALUES (NULL, 'c') +---- +NOT NULL constraint failed diff --git a/test/sql/vortex/vortex_write.test b/test/sql/vortex/vortex_write.test index 3cf8a40a..b6060c38 100644 --- a/test/sql/vortex/vortex_write.test +++ b/test/sql/vortex/vortex_write.test @@ -112,14 +112,18 @@ SELECT r FROM ducklake.ctas WHERE r = 500 # carry the metadata these need) rather than crashing or silently corrupting. # --------------------------------------------------------------------------- -# NOT NULL columns (enforced via written null-count stats, which vortex lacks) +# NOT NULL columns: enforced from the null-count statistics computed for the +# written file (see vortex_mixed_format_stats.test for the stats themselves) statement ok CREATE TABLE ducklake.nn(i INTEGER NOT NULL) -statement error +statement ok INSERT INTO ducklake.nn VALUES (1) + +statement error +INSERT INTO ducklake.nn VALUES (NULL) ---- -NOT NULL columns are not yet supported +NOT NULL constraint failed # partitioned writes (need one file per partition value + recorded partition_values) statement ok diff --git a/test/sql/vortex/vortex_written_stats_types.test b/test/sql/vortex/vortex_written_stats_types.test new file mode 100644 index 00000000..0c24952e --- /dev/null +++ b/test/sql/vortex/vortex_written_stats_types.test @@ -0,0 +1,357 @@ +# name: test/sql/vortex/vortex_written_stats_types.test +# description: Adversarial coverage for in-stream written statistics (DuckLakeStatsCopy): +# every stats-capable type with extreme values, NaN/inf semantics, all-NULL and +# empty writes, the oversized-string bound cap, counts-only types, cross-chunk +# bounds, parallel sinks, multi-insert global merges, and the read-side consumers +# of the recorded bounds. +# group: [vortex] + +require ducklake + +require parquet + +require vortex + +statement ok +ATTACH 'ducklake:{TEST_DIR}/dl_vstats.db' AS ducklake (DATA_PATH '{TEST_DIR}/dl_vstats_files', METADATA_CATALOG 'ducklake_meta') + +statement ok +CALL ducklake.set_option('data_inlining_row_limit', 0) + +statement ok +CALL ducklake.set_option('data_file_format', 'vortex') + +# --------------------------------------------------------------------------- +# Every stats-capable type, with type extremes, negatives, and NULLs +# --------------------------------------------------------------------------- +statement ok +CREATE TABLE ducklake.t_types( + b BOOLEAN, i1 TINYINT, i2 SMALLINT, i4 INTEGER, i8 BIGINT, + u1 UTINYINT, u4 UINTEGER, u8 UBIGINT, + f FLOAT, d DOUBLE, dec1 DECIMAL(18,3), + dt DATE, tm TIME, ts TIMESTAMP, s VARCHAR) + +statement ok +INSERT INTO ducklake.t_types VALUES + (false, -128, -32768, -2147483648, -9223372036854775808, + 0, 0, 0, + -3.5, -1e300, -999999999999999.999, + DATE '1900-01-05', TIME '00:00:00', TIMESTAMP '1900-01-05 01:02:03.456789', ''), + (true, 127, 32767, 2147483647, 9223372036854775807, + 255, 4294967295, 18446744073709551615, + 3.5, 1e300, 999999999999999.999, + DATE '2999-12-31', TIME '23:59:59.999999', TIMESTAMP '2999-12-31 23:59:59.999999', 'zz "quoted" 日本語'), + (NULL, NULL, NULL, NULL, NULL, NULL, NULL, NULL, + NULL, NULL, NULL, NULL, NULL, NULL, NULL), + (false, 0, 0, 0, 0, 1, 1, 1, 0.0, 0.0, 0.0, + DATE '2024-06-15', TIME '12:00:00', TIMESTAMP '2024-06-15 12:00:00', 'middle') + +query IIIII +SELECT column_id, min_value, max_value, null_count, value_count +FROM ducklake_meta.ducklake_file_column_stats ORDER BY column_id +---- +1 false true 1 3 +2 -128 127 1 3 +3 -32768 32767 1 3 +4 -2147483648 2147483647 1 3 +5 -9223372036854775808 9223372036854775807 1 3 +6 0 255 1 3 +7 0 4294967295 1 3 +8 0 18446744073709551615 1 3 +9 -3.5 3.5 1 3 +10 -1e+300 1e+300 1 3 +11 -999999999999999.999 999999999999999.999 1 3 +12 1900-01-05 2999-12-31 1 3 +13 00:00:00 23:59:59.999999 1 3 +14 1900-01-05 01:02:03.456789 2999-12-31 23:59:59.999999 1 3 +15 (empty) zz "quoted" 日本語 1 3 + +# types outside vortex's type system fail cleanly at bind (pinned here so a +# silent change in that behaviour is noticed): 128-bit ints, UUID, INTERVAL +statement ok +CREATE TABLE ducklake.t_i128(h HUGEINT) + +statement error +INSERT INTO ducklake.t_i128 VALUES (1) +---- +not in Vortex type system + +statement ok +CREATE TABLE ducklake.t_uuid(u UUID) + +statement error +INSERT INTO ducklake.t_uuid VALUES ('00000000-0000-0000-0000-000000000001') +---- +conversion is not supported + +statement ok +CREATE TABLE ducklake.t_interval(iv INTERVAL) + +statement error +INSERT INTO ducklake.t_interval VALUES (INTERVAL 3 DAYS) +---- +conversion is not supported + +# timestamptz is stats-capable +statement ok +CREATE TABLE ducklake.t_tstz(tz TIMESTAMPTZ) + +statement ok +INSERT INTO ducklake.t_tstz VALUES (TIMESTAMP '2024-01-01 00:00:00'), (TIMESTAMP '2025-06-01 12:30:00'), (NULL) + +query IIII +SELECT min_value, max_value, null_count, value_count +FROM ducklake_meta.ducklake_file_column_stats +WHERE data_file_id = (SELECT data_file_id FROM ducklake_meta.ducklake_data_file d JOIN ducklake_meta.ducklake_table t USING (table_id) WHERE t.table_name = 't_tstz') +---- +2024-01-01 00:00:00+00 2025-06-01 12:30:00+00 1 2 + +# the recorded bounds feed the read side: filters at and beyond the bounds +query I +SELECT count(*) FROM ducklake.t_types WHERE i4 = -2147483648 +---- +1 + +query I +SELECT count(*) FROM ducklake.t_types WHERE i8 > 0 +---- +1 + +# negative-heavy ORDER BY .. LIMIT (sort-key compression guard with negatives) +query I +SELECT i2 FROM ducklake.t_types ORDER BY i2 ASC LIMIT 1 +---- +-32768 + +# --------------------------------------------------------------------------- +# Floats: NaN excluded from bounds and flagged; infinities ARE bounds +# --------------------------------------------------------------------------- +statement ok +CREATE TABLE ducklake.t_float(f FLOAT, d DOUBLE, all_nan DOUBLE) + +statement ok +INSERT INTO ducklake.t_float VALUES + (1.5, 'inf', 'nan'), ('nan', '-inf', 'nan'), (-2.5, 42.0, 'nan'), ('nan', 'nan', 'nan') + +query IIIIII +SELECT column_id, min_value, max_value, null_count, value_count, contains_nan +FROM ducklake_meta.ducklake_file_column_stats +WHERE data_file_id = (SELECT data_file_id FROM ducklake_meta.ducklake_data_file d JOIN ducklake_meta.ducklake_table t USING (table_id) WHERE t.table_name = 't_float') +ORDER BY column_id +---- +1 -2.5 1.5 0 4 true +2 -inf inf 0 4 true +3 NULL NULL 0 4 true + +# NaN rows are still found by scans (bounds must not prune them away) +query I +SELECT count(*) FROM ducklake.t_float WHERE isnan(f) +---- +2 + +# --------------------------------------------------------------------------- +# Counts-only types: no bounds recorded, counts exact +# --------------------------------------------------------------------------- +statement ok +CREATE TABLE ducklake.t_nominmax(bl BLOB, li INTEGER[], st STRUCT(a INTEGER)) + +statement ok +INSERT INTO ducklake.t_nominmax VALUES + ('\xAA\xBB', [1, 2], {'a': 1}), + ('\x00', [], {'a': NULL}), + (NULL, NULL, NULL) + +query IIIII +SELECT column_id, min_value, max_value, null_count, value_count +FROM ducklake_meta.ducklake_file_column_stats +WHERE data_file_id = (SELECT data_file_id FROM ducklake_meta.ducklake_data_file d JOIN ducklake_meta.ducklake_table t USING (table_id) WHERE t.table_name = 't_nominmax') +ORDER BY column_id +---- +1 NULL NULL 1 2 +2 NULL NULL 1 2 +4 NULL NULL 1 2 + +# equality filters still work without bounds +query I +SELECT count(*) FROM ducklake.t_nominmax WHERE bl = '\x00'::BLOB +---- +1 + +# --------------------------------------------------------------------------- +# All-NULL column and empty write +# --------------------------------------------------------------------------- +statement ok +CREATE TABLE ducklake.t_allnull(i INTEGER, s VARCHAR) + +statement ok +INSERT INTO ducklake.t_allnull SELECT NULL, NULL FROM range(100) + +query IIIII +SELECT column_id, min_value, max_value, null_count, value_count +FROM ducklake_meta.ducklake_file_column_stats +WHERE data_file_id = (SELECT data_file_id FROM ducklake_meta.ducklake_data_file d JOIN ducklake_meta.ducklake_table t USING (table_id) WHERE t.table_name = 't_allnull') +ORDER BY column_id +---- +1 NULL NULL 100 0 +2 NULL NULL 100 0 + +query I +SELECT count(*) FROM ducklake.t_allnull WHERE i = 5 +---- +0 + +# an INSERT writing zero rows registers no data file and no stats +statement ok +CREATE TABLE ducklake.t_empty(i INTEGER) + +statement ok +INSERT INTO ducklake.t_empty SELECT 1 WHERE 1 = 0 + +query I +SELECT count(*) FROM ducklake_meta.ducklake_data_file d +JOIN ducklake_meta.ducklake_table t USING (table_id) WHERE t.table_name = 't_empty' +---- +0 + +# --------------------------------------------------------------------------- +# Oversized string bounds are dropped (a truncated max is not an upper bound) +# --------------------------------------------------------------------------- +statement ok +CREATE TABLE ducklake.t_bigstr(s VARCHAR) + +# one value inside the 1024-byte cap, one outside +statement ok +INSERT INTO ducklake.t_bigstr VALUES (repeat('a', 1024)), (repeat('z', 1025)) + +query IIII +SELECT min_value IS NULL, max_value IS NULL, null_count, value_count +FROM ducklake_meta.ducklake_file_column_stats +WHERE data_file_id = (SELECT data_file_id FROM ducklake_meta.ducklake_data_file d JOIN ducklake_meta.ducklake_table t USING (table_id) WHERE t.table_name = 't_bigstr') +---- +true true 0 2 + +# values at exactly the cap keep their bounds +statement ok +CREATE TABLE ducklake.t_capstr(s VARCHAR) + +statement ok +INSERT INTO ducklake.t_capstr VALUES (repeat('a', 1024)), ('b') + +query II +SELECT min_value = repeat('a', 1024), max_value +FROM ducklake_meta.ducklake_file_column_stats +WHERE data_file_id = (SELECT data_file_id FROM ducklake_meta.ducklake_data_file d JOIN ducklake_meta.ducklake_table t USING (table_id) WHERE t.table_name = 't_capstr') +---- +true b + +# --------------------------------------------------------------------------- +# Cross-chunk bounds and parallel sinks: 500k rows, extremes far apart, exact +# --------------------------------------------------------------------------- +statement ok +CREATE TABLE ducklake.t_big(i INTEGER, s VARCHAR) + +statement ok +INSERT INTO ducklake.t_big +SELECT CASE WHEN x = 480000 THEN -1000000 WHEN x = 20000 THEN 1000000 + WHEN x % 97 = 0 THEN NULL ELSE (x % 5000)::INTEGER END, + 'k' || (x % 1000) +FROM range(500000) t(x) + +query IIII +SELECT min_value, max_value, null_count, value_count +FROM ducklake_meta.ducklake_file_column_stats s +JOIN ducklake_meta.ducklake_data_file d USING (data_file_id) +JOIN ducklake_meta.ducklake_table t ON d.table_id = t.table_id +WHERE t.table_name = 't_big' AND s.column_id = 1 +---- +-1000000 1000000 5155 494845 + +query III +SELECT count(*), min(i), max(i) FROM ducklake.t_big +---- +500000 -1000000 1000000 + +query I +SELECT count(*) FROM ducklake.t_big WHERE i = 1000000 +---- +1 + +# --------------------------------------------------------------------------- +# Multiple inserts merge the global bounds across snapshots +# --------------------------------------------------------------------------- +statement ok +CREATE TABLE ducklake.t_merge(i INTEGER) + +statement ok +INSERT INTO ducklake.t_merge VALUES (10), (20) + +statement ok +INSERT INTO ducklake.t_merge VALUES (-5) + +statement ok +INSERT INTO ducklake.t_merge VALUES (NULL), (99) + +query III +SELECT min_value, max_value, contains_null +FROM ducklake_meta.ducklake_table_column_stats s +JOIN ducklake_meta.ducklake_table t USING (table_id) +WHERE t.table_name = 't_merge' +---- +-5 99 true + +# --------------------------------------------------------------------------- +# NOT NULL enforced from written stats; a failed insert commits nothing +# --------------------------------------------------------------------------- +statement ok +CREATE TABLE ducklake.t_nn(i INTEGER NOT NULL, s VARCHAR) + +statement ok +INSERT INTO ducklake.t_nn VALUES (1, 'a') + +statement error +INSERT INTO ducklake.t_nn SELECT CASE WHEN x = 90000 THEN NULL ELSE x::INTEGER END, 'v' FROM range(100000) t(x) +---- +NOT NULL constraint failed + +query II +SELECT count(*), max(i) FROM ducklake.t_nn +---- +1 1 + +# --------------------------------------------------------------------------- +# Reverse mixed order: vortex file first, then parquet, bounds still merge +# --------------------------------------------------------------------------- +statement ok +CREATE TABLE ducklake.t_rev(i INTEGER) + +statement ok +INSERT INTO ducklake.t_rev SELECT x::INTEGER FROM range(500, 600) t(x) + +statement ok +CALL ducklake.set_option('data_file_format', 'parquet') + +statement ok +INSERT INTO ducklake.t_rev SELECT x::INTEGER FROM range(0, 100) t(x) + +query II +SELECT file_format, count(*) FROM ducklake_meta.ducklake_data_file d +JOIN ducklake_meta.ducklake_table t USING (table_id) +WHERE t.table_name = 't_rev' GROUP BY file_format ORDER BY file_format +---- +parquet 1 +vortex 1 + +query III +SELECT count(*), min(i), max(i) FROM ducklake.t_rev +---- +200 0 599 + +query I +SELECT count(*) FROM ducklake.t_rev WHERE i = 550 +---- +1 + +query I +SELECT i FROM ducklake.t_rev ORDER BY i DESC LIMIT 1 +---- +599