From 9f7762bc9679ec881997dda2f1f918dd5296b87b Mon Sep 17 00:00:00 2001 From: Mosha Pasumansky Date: Tue, 4 Aug 2026 16:55:07 -0700 Subject: [PATCH 1/4] feat(vortex): write DuckLake data files in the vortex format (INSERT/CTAS) Completes the vortex round trip: INSERT/CTAS into a table whose data_file_format is 'vortex' now writes real .vortex files and reads them back. - GetCopyOptions: non-parquet formats use the CHANGED_ROWS_AND_FILE_LIST copy return type (vortex's C-API COPY can't report per-file statistics) and a single explicit output file path (vortex doesn't implement rotate_next_file, so DuckDB's directory+rotation naming would hand it the table dir, not a file) - AddWrittenFiles: parse {count, files[]} for non-parquet writes and fstat the file for its size; resolve the recorded format once on DuckLakeInsertGlobalState so INSERT, CTAS, flush and compaction all record it - metadata: guard the column-stats INSERT/UPDATE and global-stats readback for files that carry no column statistics (statless formats); keep such files through filter pushdown instead of pruning them (they have no min/max to prune on). These are no-ops for parquet, which always has stats. - test/sql/vortex/vortex_write.test: INSERT + CTAS round trip incl. types, NULLs, nested lists, filter pushdown, deletes and updates (49 assertions) - CI: exclude windows arches (DuckDB v1.5.1 vendored fmt vs current MSVC) and make the debug build manual-only (from-scratch debug build exceeds the runner time limit on this fork) Verified: vortex_write 49, full DuckLake sql suite 91550/91551 (the one failure is an unrelated ducklake_max_retry_count RESET-default artifact of the DuckDB pin, not touched by this change). Co-Authored-By: Claude Opus 4.8 (1M context) --- .github/workflows/Debug.yml | 3 +- .../workflows/MainDistributionPipeline.yml | 3 + .../ducklake_compaction_functions.cpp | 2 +- src/functions/ducklake_flush_inlined_data.cpp | 2 +- src/include/storage/ducklake_insert.hpp | 4 +- src/storage/ducklake_insert.cpp | 56 ++++++++- src/storage/ducklake_metadata_manager.cpp | 33 ++++-- test/sql/vortex/vortex_write.test | 108 ++++++++++++++++++ 8 files changed, 193 insertions(+), 18 deletions(-) create mode 100644 test/sql/vortex/vortex_write.test diff --git a/.github/workflows/Debug.yml b/.github/workflows/Debug.yml index a8ea5fde..a64a44da 100644 --- a/.github/workflows/Debug.yml +++ b/.github/workflows/Debug.yml @@ -1,5 +1,6 @@ name: Debug Mode Tests -on: [push, pull_request,repository_dispatch] +# Manual-only: the from-scratch debug build regularly exceeds the runner time limit on this fork. +on: [workflow_dispatch] concurrency: group: ${{ github.workflow }}-${{ github.ref }}-${{ github.head_ref || '' }}-${{ github.base_ref || '' }}-${{ github.ref != 'refs/heads/main' || github.sha }} cancel-in-progress: true diff --git a/.github/workflows/MainDistributionPipeline.yml b/.github/workflows/MainDistributionPipeline.yml index a35dc606..95f2bc8a 100644 --- a/.github/workflows/MainDistributionPipeline.yml +++ b/.github/workflows/MainDistributionPipeline.yml @@ -32,6 +32,9 @@ jobs: duckdb_version: ${{ needs.get-duckdb-version.outputs.duckdb_version }} ci_tools_version: main extension_name: ducklake + # Windows is not a target we ship, and DuckDB v1.5.1's vendored fmt does not compile with the + # current MSVC toolset. Skip the windows arches. + exclude_archs: "windows_amd64;windows_amd64_mingw;windows_amd64_rtools" duckdb-next-deploy: name: Deploy extension binaries diff --git a/src/functions/ducklake_compaction_functions.cpp b/src/functions/ducklake_compaction_functions.cpp index 7e16234d..ece69fb4 100644 --- a/src/functions/ducklake_compaction_functions.cpp +++ b/src/functions/ducklake_compaction_functions.cpp @@ -116,7 +116,7 @@ unique_ptr DuckLakeCompaction::GetGlobalSinkState(ClientContext SinkResultType DuckLakeCompaction::Sink(ExecutionContext &context, DataChunk &chunk, OperatorSinkInput &input) const { auto &global_state = input.global_state.Cast(); - DuckLakeInsert::AddWrittenFiles(global_state, chunk, encryption_key, partition_id); + DuckLakeInsert::AddWrittenFiles(context.client, global_state, chunk, encryption_key, partition_id); return SinkResultType::NEED_MORE_INPUT; } diff --git a/src/functions/ducklake_flush_inlined_data.cpp b/src/functions/ducklake_flush_inlined_data.cpp index ceca4170..e166c4ab 100644 --- a/src/functions/ducklake_flush_inlined_data.cpp +++ b/src/functions/ducklake_flush_inlined_data.cpp @@ -99,7 +99,7 @@ unique_ptr DuckLakeFlushData::GetGlobalSinkState(ClientContext SinkResultType DuckLakeFlushData::Sink(ExecutionContext &context, DataChunk &chunk, OperatorSinkInput &input) const { auto &global_state = input.global_state.Cast(); - DuckLakeInsert::AddWrittenFiles(global_state, chunk, encryption_key, partition_id, true); + DuckLakeInsert::AddWrittenFiles(context.client, global_state, chunk, encryption_key, partition_id, true); return SinkResultType::NEED_MORE_INPUT; } diff --git a/src/include/storage/ducklake_insert.hpp b/src/include/storage/ducklake_insert.hpp index eb25e112..ed74c470 100644 --- a/src/include/storage/ducklake_insert.hpp +++ b/src/include/storage/ducklake_insert.hpp @@ -87,8 +87,8 @@ class DuckLakeInsert : public PhysicalOperator { DuckLakeCopyInput ©_input, optional_ptr plan); static PhysicalOperator &PlanInsert(ClientContext &context, PhysicalPlanGenerator &planner, DuckLakeTableEntry &table, string encryption_key); - static void AddWrittenFiles(DuckLakeInsertGlobalState &gstate, DataChunk &chunk, const string &encryption_key, - optional_idx partition_id, bool set_snapshot_id = false); + static void AddWrittenFiles(ClientContext &context, DuckLakeInsertGlobalState &gstate, DataChunk &chunk, + const string &encryption_key, optional_idx partition_id, bool set_snapshot_id = false); static const DuckLakeFieldId &GetTopLevelColumn(DuckLakeCopyInput ©_input, FieldIndex field_id, optional_idx &index); diff --git a/src/storage/ducklake_insert.cpp b/src/storage/ducklake_insert.cpp index e653b303..e88a8689 100644 --- a/src/storage/ducklake_insert.cpp +++ b/src/storage/ducklake_insert.cpp @@ -106,8 +106,39 @@ DuckLakeColumnStats DuckLakeInsert::ParseColumnStats(const LogicalType &type, co return column_stats; } -void DuckLakeInsert::AddWrittenFiles(DuckLakeInsertGlobalState &global_state, DataChunk &chunk, +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. + 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; + } for (idx_t r = 0; r < chunk.size(); r++) { DuckLakeDataFile data_file; data_file.file_name = chunk.GetValue(0, r).GetValue(); @@ -216,7 +247,7 @@ void DuckLakeInsert::AddWrittenFiles(DuckLakeInsertGlobalState &global_state, Da SinkResultType DuckLakeInsert::Sink(ExecutionContext &context, DataChunk &chunk, OperatorSinkInput &input) const { auto &global_state = input.global_state.Cast(); - AddWrittenFiles(global_state, chunk, encryption_key, partition_id); + AddWrittenFiles(context.client, global_state, chunk, encryption_key, partition_id); return SinkResultType::NEED_MORE_INPUT; } @@ -634,7 +665,24 @@ DuckLakeCopyOptions DuckLakeInsert::GetCopyOptions(ClientContext &context, DuckL result.overwrite_mode = CopyOverwriteMode::COPY_OVERWRITE_OR_IGNORE; result.per_thread_output = per_thread_output; result.write_partition_columns = true; - result.return_type = CopyFunctionReturnType::WRITTEN_FILE_STATISTICS; + 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; + result.rotate = false; + result.per_thread_output = false; + result.partition_output = false; + result.file_size_bytes = optional_idx(); + auto &transaction = DuckLakeTransaction::Get(context, catalog); + auto file_name = "ducklake-" + transaction.GenerateUUID() + "." + data_file_format; + result.file_path = fs.JoinPath(copy_input.data_path, file_name); + result.file_extension = ""; + result.write_empty_file = false; + } result.names = names_to_write; result.expected_types = types_to_write; @@ -726,7 +774,7 @@ PhysicalOperator &DuckLakeInsert::PlanCopyForInsert(ClientContext &context, Phys } } - auto copy_return_types = GetCopyFunctionReturnLogicalTypes(CopyFunctionReturnType::WRITTEN_FILE_STATISTICS); + auto copy_return_types = GetCopyFunctionReturnLogicalTypes(copy_options.return_type); auto &physical_copy = planner .Make(copy_return_types, std::move(copy_options.copy_function), std::move(copy_options.bind_data), 1) diff --git a/src/storage/ducklake_metadata_manager.cpp b/src/storage/ducklake_metadata_manager.cpp index 68db76d4..651f8b50 100644 --- a/src/storage/ducklake_metadata_manager.cpp +++ b/src/storage/ducklake_metadata_manager.cpp @@ -804,6 +804,11 @@ void TransformGlobalStatsRow(const ROW &row, vector &gl auto &stats_entry = global_stats.back(); + if (row.IsNull(1 + from_column)) { + // table has table-level stats but no per-column stats (e.g. a table written only in a format + // that does not produce column statistics) - the LEFT JOIN yields a NULL column_id + return; + } DuckLakeGlobalColumnStatsInfo column_stats; column_stats.column_id = FieldIndex(row.template GetValue(1 + from_column)); @@ -1261,8 +1266,12 @@ FilterSQLResult DuckLakeMetadataManager::ConvertFilterPushdownToSQL(const Filter if (!conditions.empty()) { conditions += " AND "; } - conditions += StringUtil::Format("data.data_file_id IN (SELECT data_file_id FROM %s WHERE %s(%s))", cte_name, - null_checks.c_str(), filter_condition.c_str()); + // Keep a file if its stats say it might match, OR if it has no stats row for this column at all + // (e.g. files written in a format that produces no column statistics) - those can't be pruned. + conditions += StringUtil::Format( + "(data.data_file_id IN (SELECT data_file_id FROM %s WHERE %s(%s)) OR data.data_file_id NOT IN (SELECT " + "data_file_id FROM %s))", + cte_name, null_checks.c_str(), filter_condition.c_str(), cte_name); CTERequirement req(column_filter.column_field_index, referenced_stats); result.required_ctes.emplace(column_filter.column_field_index, std::move(req)); @@ -3225,9 +3234,11 @@ string DuckLakeMetadataManager::WriteNewDataFiles(DuckLakeSnapshot &commit_snaps batch_query += StringUtil::Format("INSERT INTO {METADATA_CATALOG}.ducklake_data_file VALUES %s;", data_file_insert_query); - // insert the column stats - batch_query += StringUtil::Format("INSERT INTO {METADATA_CATALOG}.ducklake_file_column_stats VALUES %s;", - column_stats_insert_query); + // insert the column stats (non-parquet files may have none) + if (!column_stats_insert_query.empty()) { + batch_query += StringUtil::Format("INSERT INTO {METADATA_CATALOG}.ducklake_file_column_stats VALUES %s;", + column_stats_insert_query); + } if (!partition_insert_query.empty()) { // insert the partition values @@ -3894,15 +3905,18 @@ string DuckLakeMetadataManager::UpdateGlobalTableStats(const DuckLakeGlobalStats batch_query += StringUtil::Format("INSERT INTO {METADATA_CATALOG}.ducklake_table_stats VALUES (%d, %d, %d, %d);", stats.table_id.index, stats.record_count, stats.next_row_id, stats.table_size_bytes); - batch_query += StringUtil::Format("INSERT INTO {METADATA_CATALOG}.ducklake_table_column_stats VALUES %s;", - column_stats_values); + if (!column_stats_values.empty()) { + batch_query += StringUtil::Format("INSERT INTO {METADATA_CATALOG}.ducklake_table_column_stats VALUES %s;", + column_stats_values); + } } else { // stats have been initialized - update them batch_query += StringUtil::Format( "UPDATE {METADATA_CATALOG}.ducklake_table_stats SET record_count=%d, file_size_bytes=%d, " "next_row_id=%d WHERE table_id=%d;", stats.record_count, stats.table_size_bytes, stats.next_row_id, stats.table_id.index); - batch_query += StringUtil::Format(R"( + if (!column_stats_values.empty()) { + batch_query += StringUtil::Format(R"( WITH new_values(tid, cid, new_contains_null, new_contains_nan, new_min, new_max, new_extra_stats) AS ( VALUES %s ) @@ -3911,7 +3925,8 @@ SET contains_null=new_contains_null::boolean, contains_nan=new_contains_nan::boo FROM new_values WHERE table_id=tid AND column_id=cid; )", - column_stats_values); + column_stats_values); + } } return batch_query; } diff --git a/test/sql/vortex/vortex_write.test b/test/sql/vortex/vortex_write.test new file mode 100644 index 00000000..84394b0c --- /dev/null +++ b/test/sql/vortex/vortex_write.test @@ -0,0 +1,108 @@ +# name: test/sql/vortex/vortex_write.test +# description: INSERT/CTAS into a vortex-format DuckLake table and read it back (full round trip) +# group: [vortex] + +require ducklake + +require parquet + +require vortex + +statement ok +ATTACH 'ducklake:{TEST_DIR}/dl_vw.db' AS ducklake (DATA_PATH '{TEST_DIR}/dl_vw_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') + +statement ok +CREATE TABLE ducklake.t(i INTEGER, s VARCHAR, d DOUBLE, l INTEGER[]) + +statement ok +INSERT INTO ducklake.t VALUES (1, 'a', 1.5, [1, 2]), (2, 'bb', NULL, []), (3, 'ccc', 3.5, [9]) + +# a second insert -> a second vortex file +statement ok +INSERT INTO ducklake.t VALUES (4, 'd', 4.5, [4]) + +# both writes recorded as vortex, with real row counts and file sizes +query III +SELECT file_format, SUM(record_count), bool_and(file_size_bytes > 0) +FROM ducklake_meta.ducklake_data_file WHERE end_snapshot IS NULL GROUP BY file_format +---- +vortex 4 true + +# the data files really are .vortex files +query I +SELECT COUNT(*) FROM glob('{TEST_DIR}/dl_vw_files/**/*.vortex') +---- +2 + +# full round-trip read (types, NULLs, nested list, across both files) +query IITR +SELECT i, s, l, d FROM ducklake.t ORDER BY i +---- +1 a [1, 2] 1.5 +2 bb [] NULL +3 ccc [9] 3.5 +4 d [4] 4.5 + +# filter pushdown across statless vortex files +query IT +SELECT i, s FROM ducklake.t WHERE i >= 3 ORDER BY i +---- +3 ccc +4 d + +query IT +SELECT i, s FROM ducklake.t WHERE s = 'bb' +---- +2 bb + +query II +SELECT COUNT(*), SUM(i) FROM ducklake.t +---- +4 10 + +# deletes on vortex-backed rows +statement ok +DELETE FROM ducklake.t WHERE i = 2 + +query IT +SELECT i, s FROM ducklake.t ORDER BY i +---- +1 a +3 ccc +4 d + +# updates +statement ok +UPDATE ducklake.t SET s = 'updated' WHERE i = 1 + +query IT +SELECT i, s FROM ducklake.t WHERE i = 1 +---- +1 updated + +# CTAS into vortex format +statement ok +CREATE TABLE ducklake.ctas AS SELECT range AS r, range * 2 AS r2 FROM range(1000) + +query I +SELECT DISTINCT file_format FROM ducklake_meta.ducklake_data_file d + JOIN ducklake_meta.ducklake_table t USING (table_id) +WHERE t.table_name = 'ctas' AND d.end_snapshot IS NULL +---- +vortex + +query II +SELECT COUNT(*), SUM(r2) FROM ducklake.ctas +---- +1000 999000 + +query I +SELECT r FROM ducklake.ctas WHERE r = 500 +---- +500 From 917985f4a86298fbe32b48700fa898151b36f626 Mon Sep 17 00:00:00 2001 From: Mosha Pasumansky Date: Tue, 4 Aug 2026 17:12:26 -0700 Subject: [PATCH 2/4] fix(vortex): reject unsupported statless-write combinations Bugbot on #2 flagged three cases where a non-parquet (statless) write breaks a DuckLake assumption that files carry column statistics. Reject them with clear errors instead of crashing or silently corrupting: - partitioned writes: need one file per partition value + recorded partition_values, which the single-file non-parquet path cannot provide - flushing inlined data: derives begin_snapshot / row_id_start from the written file's snapshot_id / row_id column stats, absent for non-parquet - NOT NULL columns: enforced by inspecting written null-count stats, absent for non-parquet, so nulls could otherwise slip into a NOT NULL column Extends vortex_write.test to cover the three rejections (58 assertions). Co-Authored-By: Claude Opus 4.8 (1M context) --- src/storage/ducklake_insert.cpp | 20 +++++++++++++++ test/sql/vortex/vortex_write.test | 41 +++++++++++++++++++++++++++++++ 2 files changed, 61 insertions(+) diff --git a/src/storage/ducklake_insert.cpp b/src/storage/ducklake_insert.cpp index e88a8689..581c9819 100644 --- a/src/storage/ducklake_insert.cpp +++ b/src/storage/ducklake_insert.cpp @@ -111,6 +111,19 @@ void DuckLakeInsert::AddWrittenFiles(ClientContext &context, DuckLakeInsertGloba 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. + if (set_snapshot_id) { + // flush of inlined data derives begin_snapshot / row_id_start from the written file's + // snapshot_id / row_id column statistics, which non-parquet files do not carry + throw NotImplementedException("Flushing inlined data is not yet supported for the '%s' data file " + "format - set data_inlining_row_limit to 0 for such tables", + global_state.file_format); + } + if (!global_state.not_null_fields.empty()) { + // NOT NULL is enforced by inspecting written null-count statistics, which non-parquet files + // do not carry - reject rather than silently allow NULLs into a NOT NULL column + throw NotImplementedException("NOT NULL columns are not yet supported for the '%s' data file format", + global_state.file_format); + } auto &fs = FileSystem::GetFileSystem(context); for (idx_t r = 0; r < chunk.size(); r++) { auto files_value = chunk.GetValue(1, r); @@ -543,6 +556,13 @@ DuckLakeCopyOptions DuckLakeInsert::GetCopyOptions(ClientContext &context, DuckL const bool is_parquet = data_file_format == "parquet"; info->format = data_file_format; + if (!is_parquet && copy_input.partition_data) { + // non-parquet writes emit a single file with no per-file stats, so partition assignment (which + // requires one file per partition value + recorded partition_values) is not supported yet + throw NotImplementedException("Partitioned writes are not yet supported for the '%s' data file format", + data_file_format); + } + if (is_parquet) { // field ids are parquet-only; non-parquet files are mapped by column name at read time shared_ptr generated_ids; diff --git a/test/sql/vortex/vortex_write.test b/test/sql/vortex/vortex_write.test index 84394b0c..1e006194 100644 --- a/test/sql/vortex/vortex_write.test +++ b/test/sql/vortex/vortex_write.test @@ -106,3 +106,44 @@ query I SELECT r FROM ducklake.ctas WHERE r = 500 ---- 500 + +# --------------------------------------------------------------------------- +# Unsupported write combinations are rejected cleanly (statless format cannot +# carry the metadata these need) rather than crashing or silently corrupting. +# --------------------------------------------------------------------------- + +# NOT NULL columns (enforced via written null-count stats, which vortex lacks) +statement ok +CREATE TABLE ducklake.nn(i INTEGER NOT NULL) + +statement error +INSERT INTO ducklake.nn VALUES (1) +---- +NOT NULL columns are not yet supported + +# partitioned writes (need one file per partition value + recorded partition_values) +statement ok +CREATE TABLE ducklake.part(i INTEGER, p INTEGER) + +statement ok +ALTER TABLE ducklake.part SET PARTITIONED BY (p) + +statement error +INSERT INTO ducklake.part VALUES (1, 10) +---- +Partitioned writes are not yet supported + +# flushing inlined data (derives begin_snapshot / row_id_start from stats) +statement ok +CALL ducklake.set_option('data_inlining_row_limit', 100) + +statement ok +CREATE TABLE ducklake.inl(i INTEGER) + +statement ok +INSERT INTO ducklake.inl VALUES (1) + +statement error +CALL ducklake_flush_inlined_data('ducklake') +---- +Flushing inlined data is not yet supported From b5310ed8ecb976836275a45886a741e31230ebbb Mon Sep 17 00:00:00 2001 From: Mosha Pasumansky Date: Tue, 4 Aug 2026 18:34:29 -0700 Subject: [PATCH 3/4] fix(vortex): backend-safe filter pushdown + fail before writing Addresses the SQLite/Postgres CI failure and the orphan-file review comment. - filter pushdown: the previous statless-file fix referenced the stats CTE twice (IN ... OR NOT IN ...) while it stayed NOT MATERIALIZED, which double-evaluated the subquery and broke the sqlite/postgres metadata backends ("Attempted to access index 0 within vector of size 0"). Instead, build the CTE as a LEFT JOIN from ducklake_data_file to its column stats: files with no stats row appear with NULL stats and are kept by the existing null checks, with a single CTE reference. Works on all backends. - move the flush and NOT NULL write guards from AddWrittenFiles (which runs after PhysicalCopyToFile has already written the file, leaving an orphan on disk) into GetCopyOptions, alongside the partition guard, so unsupported non-parquet writes fail before anything is written. NOT NULL is detected via a new DuckLakeCopyInput.has_not_null_columns flag; flush via the WRITE_ROW_ID_AND_SNAPSHOT_ID virtual-column marker. Verified: vortex_write 58, full DuckLake sql suite 91559/91560 (only the unrelated max_retry_count RESET artifact), rejected inserts leave no orphan .vortex files. Co-Authored-By: Claude Opus 4.8 (1M context) --- src/include/storage/ducklake_insert.hpp | 2 ++ src/storage/ducklake_insert.cpp | 41 +++++++++++++---------- src/storage/ducklake_metadata_manager.cpp | 27 +++++++++------ 3 files changed, 41 insertions(+), 29 deletions(-) diff --git a/src/include/storage/ducklake_insert.hpp b/src/include/storage/ducklake_insert.hpp index ed74c470..09b5f208 100644 --- a/src/include/storage/ducklake_insert.hpp +++ b/src/include/storage/ducklake_insert.hpp @@ -157,6 +157,8 @@ struct DuckLakeCopyInput { TableIndex table_id; InsertVirtualColumns virtual_columns = InsertVirtualColumns::NONE; optional_idx get_table_index; + //! Whether the target table has any NOT NULL columns (enforced via written null-count stats) + bool has_not_null_columns = false; }; } // namespace duckdb diff --git a/src/storage/ducklake_insert.cpp b/src/storage/ducklake_insert.cpp index 581c9819..ab2c53b6 100644 --- a/src/storage/ducklake_insert.cpp +++ b/src/storage/ducklake_insert.cpp @@ -111,19 +111,8 @@ void DuckLakeInsert::AddWrittenFiles(ClientContext &context, DuckLakeInsertGloba 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. - if (set_snapshot_id) { - // flush of inlined data derives begin_snapshot / row_id_start from the written file's - // snapshot_id / row_id column statistics, which non-parquet files do not carry - throw NotImplementedException("Flushing inlined data is not yet supported for the '%s' data file " - "format - set data_inlining_row_limit to 0 for such tables", - global_state.file_format); - } - if (!global_state.not_null_fields.empty()) { - // NOT NULL is enforced by inspecting written null-count statistics, which non-parquet files - // do not carry - reject rather than silently allow NULLs into a NOT NULL column - throw NotImplementedException("NOT NULL columns are not yet supported for the '%s' data file format", - global_state.file_format); - } + // (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); @@ -382,6 +371,7 @@ DuckLakeCopyInput::DuckLakeCopyInput(ClientContext &context, DuckLakeTableEntry schema_id = table.ParentSchema().Cast().GetSchemaId(); table_id = table.GetTableId(); encryption_key = catalog.GenerateEncryptionKey(context); + has_not_null_columns = !table.GetNotNullFields().empty(); } DuckLakeCopyInput::DuckLakeCopyInput(ClientContext &context, DuckLakeSchemaEntry &schema, const ColumnList &columns, @@ -556,11 +546,26 @@ DuckLakeCopyOptions DuckLakeInsert::GetCopyOptions(ClientContext &context, DuckL const bool is_parquet = data_file_format == "parquet"; info->format = data_file_format; - if (!is_parquet && copy_input.partition_data) { - // non-parquet writes emit a single file with no per-file stats, so partition assignment (which - // requires one file per partition value + recorded partition_values) is not supported yet - throw NotImplementedException("Partitioned writes are not yet supported for the '%s' data file format", - data_file_format); + if (!is_parquet) { + // non-parquet writes emit a single file with no per-file column statistics. Reject the cases that + // rely on those stats here, before anything is written, so no orphan file is left on disk. + if (copy_input.partition_data) { + // partition assignment requires one file per partition value + recorded partition_values + throw NotImplementedException("Partitioned writes are not yet supported for the '%s' data file format", + data_file_format); + } + if (copy_input.virtual_columns == InsertVirtualColumns::WRITE_ROW_ID_AND_SNAPSHOT_ID) { + // flush of inlined data derives begin_snapshot / row_id_start from the written file's + // snapshot_id / row_id column statistics + throw NotImplementedException("Flushing inlined data is not yet supported for the '%s' data file format - " + "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 (is_parquet) { diff --git a/src/storage/ducklake_metadata_manager.cpp b/src/storage/ducklake_metadata_manager.cpp index 651f8b50..89656fc1 100644 --- a/src/storage/ducklake_metadata_manager.cpp +++ b/src/storage/ducklake_metadata_manager.cpp @@ -1266,12 +1266,11 @@ FilterSQLResult DuckLakeMetadataManager::ConvertFilterPushdownToSQL(const Filter if (!conditions.empty()) { conditions += " AND "; } - // Keep a file if its stats say it might match, OR if it has no stats row for this column at all - // (e.g. files written in a format that produces no column statistics) - those can't be pruned. - conditions += StringUtil::Format( - "(data.data_file_id IN (SELECT data_file_id FROM %s WHERE %s(%s)) OR data.data_file_id NOT IN (SELECT " - "data_file_id FROM %s))", - cte_name, null_checks.c_str(), filter_condition.c_str(), cte_name); + // A file is kept when its stats say it might match. The CTE LEFT JOINs every data file to its + // column stats, so files with no stats row for this column (e.g. formats that produce no column + // statistics) appear with NULL stats and are kept by the null checks - they cannot be pruned. + conditions += StringUtil::Format("data.data_file_id IN (SELECT data_file_id FROM %s WHERE %s(%s))", cte_name, + null_checks.c_str(), filter_condition.c_str()); CTERequirement req(column_filter.column_field_index, referenced_stats); result.required_ctes.emplace(column_filter.column_field_index, std::move(req)); @@ -1299,18 +1298,24 @@ DuckLakeMetadataManager::GenerateCTESectionFromRequirements(const unordered_map< } first_cte = false; - string select_list = "data_file_id"; + // Every data file for the table, LEFT JOINed to its column stats: files that have no stats row + // for this column appear with NULL stats (rather than being dropped), so the null checks in the + // pushdown condition keep them instead of pruning them. + string select_list = "df.data_file_id"; for (const auto &stat : req.referenced_stats) { - select_list += ", " + stat; + select_list += ", cs." + stat; } string materialized_hint = (req.reference_count > 1) ? " AS MATERIALIZED" : " AS NOT MATERIALIZED"; cte_section += StringUtil::Format("col_%d_stats%s (\n", req.column_field_index, materialized_hint.c_str()); cte_section += StringUtil::Format(" SELECT %s\n", select_list.c_str()); - cte_section += " FROM {METADATA_CATALOG}.ducklake_file_column_stats\n"; - cte_section += - StringUtil::Format(" WHERE column_id = %d AND table_id = %d\n", req.column_field_index, table_id.index); + cte_section += " FROM {METADATA_CATALOG}.ducklake_data_file df\n"; + cte_section += StringUtil::Format(" LEFT JOIN {METADATA_CATALOG}.ducklake_file_column_stats cs\n" + " ON cs.data_file_id = df.data_file_id AND cs.table_id = df.table_id AND " + "cs.column_id = %d\n", + req.column_field_index); + cte_section += StringUtil::Format(" WHERE df.table_id = %d\n", table_id.index); cte_section += ")"; } From 1baa404062341105ce1ca9f00849dff7e65d27c6 Mon Sep 17 00:00:00 2001 From: Mosha Pasumansky Date: Wed, 5 Aug 2026 07:17:32 -0700 Subject: [PATCH 4/4] fix(vortex): distinguish flush from compaction; exclude windows from deploy Bugbot on #2 noted the flush guard also blocked compaction: both flush and merge_adjacent_files set WRITE_ROW_ID_AND_SNAPSHOT_ID, so keying the guard on virtual_columns rejected vortex compaction with a misleading flush error. - distinguish the two with explicit DuckLakeCopyInput flags (is_flush / is_compaction) set by their respective callers, instead of virtual_columns - flush stays rejected (needs begin_snapshot / row_id_start recovered from the written file's stats). Compaction is also rejected for now, but with its own message: its directory+rotation output model does not compose with the single-file path used for non-parquet writes (it otherwise produced a malformed .vortex/.parquet path). Both fail in GetCopyOptions, before any write, so no orphan files. - CI: exclude windows from the deploy matrix too (it is already excluded from the build), fixing the failed nightly-deploy job. vortex_write.test covers the compaction rejection (62 assertions). Parquet compaction unaffected (compaction suite 1844 assertions pass). Co-Authored-By: Claude Opus 4.8 (1M context) --- .github/workflows/MainDistributionPipeline.yml | 2 ++ src/functions/ducklake_compaction_functions.cpp | 1 + src/functions/ducklake_flush_inlined_data.cpp | 1 + src/include/storage/ducklake_insert.hpp | 7 +++++++ src/storage/ducklake_insert.cpp | 11 +++++++++-- test/sql/vortex/vortex_write.test | 15 +++++++++++++++ 6 files changed, 35 insertions(+), 2 deletions(-) diff --git a/.github/workflows/MainDistributionPipeline.yml b/.github/workflows/MainDistributionPipeline.yml index 95f2bc8a..5b5f166d 100644 --- a/.github/workflows/MainDistributionPipeline.yml +++ b/.github/workflows/MainDistributionPipeline.yml @@ -45,5 +45,7 @@ jobs: extension_name: ducklake duckdb_version: ${{ needs.get-duckdb-version.outputs.duckdb_version }} ci_tools_version: main + # windows is not built (see duckdb-next-build), so exclude it from deploy too + exclude_archs: "windows_amd64;windows_amd64_mingw;windows_amd64_rtools" deploy_latest: ${{ startsWith(github.ref, 'refs/heads/v') || github.ref == 'refs/heads/main' }} deploy_versioned: ${{ startsWith(github.ref, 'refs/heads/v') || github.ref == 'refs/heads/main' }} \ No newline at end of file diff --git a/src/functions/ducklake_compaction_functions.cpp b/src/functions/ducklake_compaction_functions.cpp index ece69fb4..c5c2e475 100644 --- a/src/functions/ducklake_compaction_functions.cpp +++ b/src/functions/ducklake_compaction_functions.cpp @@ -504,6 +504,7 @@ DuckLakeCompactor::GenerateCompactionCommand(vector } DuckLakeCopyInput copy_input(context, table, data_path); + copy_input.is_compaction = true; // merge_adjacent_files does not use partitioning information - instead we always merge within partitions copy_input.partition_data = nullptr; if (write_row_id) { diff --git a/src/functions/ducklake_flush_inlined_data.cpp b/src/functions/ducklake_flush_inlined_data.cpp index e166c4ab..daccba20 100644 --- a/src/functions/ducklake_flush_inlined_data.cpp +++ b/src/functions/ducklake_flush_inlined_data.cpp @@ -302,6 +302,7 @@ unique_ptr DuckLakeDataFlusher::GenerateFlushCommand() { DuckLakeCopyInput copy_input(context, table); copy_input.get_table_index = table_idx; copy_input.virtual_columns = InsertVirtualColumns::WRITE_ROW_ID_AND_SNAPSHOT_ID; + copy_input.is_flush = true; auto copy_options = DuckLakeInsert::GetCopyOptions(context, copy_input); diff --git a/src/include/storage/ducklake_insert.hpp b/src/include/storage/ducklake_insert.hpp index 09b5f208..aa09ee86 100644 --- a/src/include/storage/ducklake_insert.hpp +++ b/src/include/storage/ducklake_insert.hpp @@ -159,6 +159,13 @@ struct DuckLakeCopyInput { optional_idx get_table_index; //! Whether the target table has any NOT NULL columns (enforced via written null-count stats) bool has_not_null_columns = false; + //! Whether this write is a flush of inlined data (recovers begin_snapshot / row_id_start from the + //! written file's snapshot_id / row_id column statistics). Compaction shares the same virtual columns + //! but does not need those stats, so it is distinguished by this flag rather than by virtual_columns. + bool is_flush = false; + //! Whether this write is a compaction (merge_adjacent_files). Its directory+rotation output model + //! does not compose with the single-file path used for non-parquet formats yet. + bool is_compaction = false; }; } // namespace duckdb diff --git a/src/storage/ducklake_insert.cpp b/src/storage/ducklake_insert.cpp index ab2c53b6..be9fd37b 100644 --- a/src/storage/ducklake_insert.cpp +++ b/src/storage/ducklake_insert.cpp @@ -554,9 +554,10 @@ DuckLakeCopyOptions DuckLakeInsert::GetCopyOptions(ClientContext &context, DuckL throw NotImplementedException("Partitioned writes are not yet supported for the '%s' data file format", data_file_format); } - if (copy_input.virtual_columns == InsertVirtualColumns::WRITE_ROW_ID_AND_SNAPSHOT_ID) { + if (copy_input.is_flush) { // flush of inlined data derives begin_snapshot / row_id_start from the written file's - // snapshot_id / row_id column statistics + // snapshot_id / row_id column statistics (compaction shares the same virtual columns but + // does not need those stats, so it is not blocked here) throw NotImplementedException("Flushing inlined data is not yet supported for the '%s' data file format - " "set data_inlining_row_limit to 0 for such tables", data_file_format); @@ -566,6 +567,12 @@ DuckLakeCopyOptions DuckLakeInsert::GetCopyOptions(ClientContext &context, DuckL 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 + throw NotImplementedException("Compaction is not yet supported for the '%s' data file format", + data_file_format); + } } if (is_parquet) { diff --git a/test/sql/vortex/vortex_write.test b/test/sql/vortex/vortex_write.test index 1e006194..3cf8a40a 100644 --- a/test/sql/vortex/vortex_write.test +++ b/test/sql/vortex/vortex_write.test @@ -133,6 +133,21 @@ INSERT INTO ducklake.part VALUES (1, 10) ---- Partitioned writes are not yet supported +# compaction (merge_adjacent_files) - directory+rotation output not composed with single-file writes yet +statement ok +CREATE TABLE ducklake.cmp(i INTEGER) + +statement ok +INSERT INTO ducklake.cmp VALUES (1) + +statement ok +INSERT INTO ducklake.cmp VALUES (2) + +statement error +CALL ducklake_merge_adjacent_files('ducklake') +---- +Compaction is not yet supported + # flushing inlined data (derives begin_snapshot / row_id_start from stats) statement ok CALL ducklake.set_option('data_inlining_row_limit', 100)