Skip to content
Merged
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
3 changes: 2 additions & 1 deletion .github/workflows/Debug.yml
Original file line number Diff line number Diff line change
@@ -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
Expand Down
5 changes: 5 additions & 0 deletions .github/workflows/MainDistributionPipeline.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -42,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' }}
3 changes: 2 additions & 1 deletion src/functions/ducklake_compaction_functions.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,7 @@ unique_ptr<GlobalSinkState> DuckLakeCompaction::GetGlobalSinkState(ClientContext

SinkResultType DuckLakeCompaction::Sink(ExecutionContext &context, DataChunk &chunk, OperatorSinkInput &input) const {
auto &global_state = input.global_state.Cast<DuckLakeInsertGlobalState>();
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;
}

Expand Down Expand Up @@ -504,6 +504,7 @@ DuckLakeCompactor::GenerateCompactionCommand(vector<DuckLakeCompactionFileEntry>
}

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) {
Expand Down
3 changes: 2 additions & 1 deletion src/functions/ducklake_flush_inlined_data.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,7 @@ unique_ptr<GlobalSinkState> DuckLakeFlushData::GetGlobalSinkState(ClientContext

SinkResultType DuckLakeFlushData::Sink(ExecutionContext &context, DataChunk &chunk, OperatorSinkInput &input) const {
auto &global_state = input.global_state.Cast<DuckLakeInsertGlobalState>();
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;
}

Expand Down Expand Up @@ -302,6 +302,7 @@ unique_ptr<LogicalOperator> 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);

Expand Down
13 changes: 11 additions & 2 deletions src/include/storage/ducklake_insert.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -87,8 +87,8 @@ class DuckLakeInsert : public PhysicalOperator {
DuckLakeCopyInput &copy_input, optional_ptr<PhysicalOperator> 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 &copy_input, FieldIndex field_id,
optional_idx &index);
Expand Down Expand Up @@ -157,6 +157,15 @@ 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;
//! 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
88 changes: 84 additions & 4 deletions src/storage/ducklake_insert.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -106,8 +106,41 @@ 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.
// (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;
Comment thread
cursor[bot] marked this conversation as resolved.
Comment thread
cursor[bot] marked this conversation as resolved.
}
for (idx_t r = 0; r < chunk.size(); r++) {
DuckLakeDataFile data_file;
data_file.file_name = chunk.GetValue(0, r).GetValue<string>();
Expand Down Expand Up @@ -216,7 +249,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<DuckLakeInsertGlobalState>();
AddWrittenFiles(global_state, chunk, encryption_key, partition_id);
AddWrittenFiles(context.client, global_state, chunk, encryption_key, partition_id);
return SinkResultType::NEED_MORE_INPUT;
}

Expand Down Expand Up @@ -338,6 +371,7 @@ DuckLakeCopyInput::DuckLakeCopyInput(ClientContext &context, DuckLakeTableEntry
schema_id = table.ParentSchema().Cast<DuckLakeSchemaEntry>().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,
Expand Down Expand Up @@ -512,6 +546,35 @@ DuckLakeCopyOptions DuckLakeInsert::GetCopyOptions(ClientContext &context, DuckL
const bool is_parquet = data_file_format == "parquet";
info->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.is_flush) {
// flush of inlined data derives 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 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);
}
Comment thread
cursor[bot] marked this conversation as resolved.
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
throw NotImplementedException("Compaction is 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<DuckLakeFieldData> generated_ids;
Expand Down Expand Up @@ -634,7 +697,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;
}
Comment thread
cursor[bot] marked this conversation as resolved.
result.names = names_to_write;
result.expected_types = types_to_write;

Expand Down Expand Up @@ -726,7 +806,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<PhysicalCopyToFile>(copy_return_types, std::move(copy_options.copy_function),
std::move(copy_options.bind_data), 1)
Expand Down
44 changes: 32 additions & 12 deletions src/storage/ducklake_metadata_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -804,6 +804,11 @@ void TransformGlobalStatsRow(const ROW &row, vector<DuckLakeGlobalStatsInfo> &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<uint64_t>(1 + from_column));

Expand Down Expand Up @@ -1261,6 +1266,9 @@ FilterSQLResult DuckLakeMetadataManager::ConvertFilterPushdownToSQL(const Filter
if (!conditions.empty()) {
conditions += " AND ";
}
// 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());

Expand Down Expand Up @@ -1290,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 += ")";
}

Expand Down Expand Up @@ -3225,9 +3239,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
Expand Down Expand Up @@ -3894,15 +3910,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
)
Expand All @@ -3911,7 +3930,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;
}
Expand Down
Loading
Loading