Skip to content

Commit 6689e8f

Browse files
committed
feat(duckdb): read pushed-down struct extracts
DuckDB can ask a table function for a struct field instead of the whole column (`ColumnIndex::IsPushdownExtract()`), which is what lets `SELECT s.a.b` avoid reading `s`. Answer the callback for struct columns -- but not when the scan has aggregates bound, since the aggregate pushdown binds column indexes itself and cannot express an extracted path -- and swap `statistics` for `statistics_extended`, which is what makes DuckDB consult that callback at all. The extended callback returns no statistics for an extracted path, as Vortex has no nested statistics. `InitializeGlobalState` now passes the flattened child paths, their offsets and the type DuckDB expects each one to be emitted as, and `Projection::new` turns a path into a chain of `get_item`s with a cast when the requested type differs from the field's own. The packed output names include the path, because DuckDB binds one output column per extracted path. `test_vortex_scan_struct_extract_projection` is ignored: DuckDB's `MultiFileColumnMapper` still rebuilds the struct for a `PUSHDOWN_EXTRACT` column, so `MultiFileReader::FinalizeChunk` executes a struct expression against a vector typed as the extracted field and raises an internal error. Supporting this for a MultiFileReader-based scan needs a DuckDB-side patch. Part of #10297 Signed-off-by: Paolo Valletta <paovalletta@hotmail.it>
1 parent 4895bd4 commit 6689e8f

7 files changed

Lines changed: 365 additions & 24 deletions

File tree

‎vortex-duckdb/cpp/include/table_function.h‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,32 @@ typedef struct {
5050

5151
duckdb_vx_table_filter_set filters;
5252
duckdb_client_context client_context;
53+
54+
/**
55+
* Struct extract paths pushed down by DuckDB for the entries in
56+
* `column_ids` where DuckDB only needs a struct field
57+
* (`ColumnIndex::IsPushdownExtract()`).
58+
*
59+
* `column_extract_indexes` holds every path concatenated, in the order of
60+
* `column_ids`, and `column_extract_offsets` has `column_ids_count + 1`
61+
* entries giving each column the slice `[offsets[i], offsets[i + 1])`.
62+
* A column that is read whole has an empty slice. The indexes are struct
63+
* child positions, outermost first: `SELECT s.a.b` on `s = {a: {b: …}}`
64+
* yields `[0, 0]`.
65+
*/
66+
const idx_t *column_extract_indexes;
67+
size_t column_extract_indexes_count;
68+
const size_t *column_extract_offsets;
69+
70+
/**
71+
* For each entry in `column_ids`, the type DuckDB expects the scan to emit
72+
* (`ColumnIndex::GetScanType()`, which is the cast target when the query
73+
* casts the extracted field), or NULL when the column is read whole.
74+
* The types are borrowed from the bind data and only valid for the duration
75+
* of the init_global call.
76+
*/
77+
const duckdb_logical_type *column_extract_types;
78+
size_t column_extract_types_count;
5379
} duckdb_vx_tfunc_init_input;
5480

5581
// Result data returned from the cardinality callback.

‎vortex-duckdb/cpp/multi_file_reader.cpp‎

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -152,6 +152,29 @@ VortexReaderInterface::InitializeGlobalState(ClientContext &context,
152152
column_ids[i] = storage_index;
153153
}
154154

155+
// DuckDB asks for a struct field instead of the whole column when it knows
156+
// the scan can extract it. Pass the path and the type it expects back, so
157+
// Rust can build the extract expression: one flattened path per column, with
158+
// offsets, and the leaf type to cast to when the query asked for one.
159+
vector<idx_t> extract_indexes;
160+
vector<size_t> extract_offsets;
161+
vector<duckdb_logical_type> extract_types(input.column_indexes.size(), nullptr);
162+
extract_offsets.reserve(input.column_indexes.size() + 1);
163+
extract_offsets.push_back(0);
164+
for (size_t i = 0; i < input.column_indexes.size(); ++i) {
165+
const ColumnIndex *level = &input.column_indexes[i];
166+
if (level->IsPushdownExtract()) {
167+
while (level->HasChildren()) {
168+
level = &level->GetChildIndex(0);
169+
extract_indexes.push_back(level->GetPrimaryIndex());
170+
}
171+
// Borrowed from the bind data, valid for this call.
172+
extract_types[i] = reinterpret_cast<duckdb_logical_type>(
173+
const_cast<LogicalType *>(&level->GetScanType()));
174+
}
175+
extract_offsets.push_back(extract_indexes.size());
176+
}
177+
155178
// MultiFileGlobalState projection_ids are filled only when this call
156179
// returns. Take these from a physical operator.
157180
const idx_t *projection_ids = nullptr;
@@ -173,6 +196,11 @@ VortexReaderInterface::InitializeGlobalState(ClientContext &context,
173196
.projection_ids_count = projection_ids_count,
174197
.filters = reinterpret_cast<duckdb_vx_table_filter_set>(input.filters.get()),
175198
.client_context = reinterpret_cast<duckdb_client_context>(&context),
199+
.column_extract_indexes = extract_indexes.empty() ? nullptr : extract_indexes.data(),
200+
.column_extract_indexes_count = extract_indexes.size(),
201+
.column_extract_offsets = extract_offsets.data(),
202+
.column_extract_types = extract_types.data(),
203+
.column_extract_types_count = extract_types.size(),
176204
};
177205

178206
duckdb_vx_error error_out = nullptr;

‎vortex-duckdb/cpp/table_function.cpp‎

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -209,7 +209,31 @@ duckdb_state register_table_function(DatabaseInstance &db, LogicalType parameter
209209
return {COLUMN_IDENTIFIER_FILE_INDEX, COLUMN_IDENTIFIER_FILE_ROW_NUMBER};
210210
};
211211

212-
fn.statistics = MultiFileFunction<VortexReaderInterface>::MultiFileScanStats;
212+
// DuckDB only pushes a struct extract into a scan when the table function has
213+
// no `statistics` callback (RemoveUnusedColumns::CheckPushdownExtract), so use
214+
// the extended variant: it sees the whole column index, and can report no
215+
// statistics for an extracted path, since Vortex has no nested statistics.
216+
fn.statistics_extended = [](ClientContext &context, TableFunctionGetStatisticsInput &input) {
217+
if (input.column_index.IsPushdownExtract()) {
218+
return unique_ptr<BaseStatistics>();
219+
}
220+
return MultiFileFunction<VortexReaderInterface>::MultiFileScanStats(
221+
context, input.bind_data.get(), input.column_index.GetPrimaryIndex());
222+
};
223+
224+
// Only a struct column can be read as an extracted path, and never when
225+
// DuckDB bound aggregates to this scan: the aggregate pushdown binds column
226+
// indexes itself and cannot express an extracted path.
227+
fn.supports_pushdown_extract = [](const FunctionData &bind_data, const LogicalIndex &col_idx) {
228+
if (duckdb_reader_is_aggregate(get_ffi_bind(&bind_data))) {
229+
return false;
230+
}
231+
const auto &bind = bind_data.Cast<MultiFileBindData>();
232+
if (col_idx.index >= bind.types.size()) {
233+
return false;
234+
}
235+
return bind.types[col_idx.index].id() == LogicalTypeId::STRUCT;
236+
};
213237
fn.get_partition_stats = get_partition_stats;
214238
fn.get_multi_file_reader = get_multi_file_reader;
215239

‎vortex-duckdb/src/duckdb/table_init_input.rs‎

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,9 +6,56 @@ use std::fmt::Formatter;
66
use std::fmt::Result;
77

88
use crate::cpp;
9+
use crate::duckdb::LogicalType;
10+
use crate::duckdb::LogicalTypeRef;
911
use crate::duckdb::TableFilterSet;
1012
use crate::duckdb::TableFilterSetRef;
1113

14+
/// A struct field DuckDB asked the scan to read instead of the whole column.
15+
#[derive(Debug, Clone, Copy)]
16+
pub struct StructExtract<'a> {
17+
/// Struct child indexes, outermost first.
18+
pub indexes: &'a [u64],
19+
/// The type DuckDB expects the scan to emit for this field, which is the
20+
/// cast target when the query casts the extracted field itself.
21+
pub datatype: &'a LogicalTypeRef,
22+
}
23+
24+
/// The struct extracts DuckDB pushed into the scan, one entry per column id.
25+
#[derive(Debug, Clone, Copy, Default)]
26+
pub struct StructExtracts<'a> {
27+
pub(crate) indexes: &'a [u64],
28+
pub(crate) offsets: &'a [usize],
29+
pub(crate) datatypes: &'a [cpp::duckdb_logical_type],
30+
}
31+
32+
impl<'a> StructExtracts<'a> {
33+
/// The extract for the column at `column`, if DuckDB asked for one.
34+
pub fn get(&self, column: usize) -> Option<StructExtract<'a>> {
35+
let start = usize::try_from(*self.offsets.get(column)?).ok()?;
36+
let end = usize::try_from(*self.offsets.get(column + 1)?).ok()?;
37+
if start == end {
38+
return None;
39+
}
40+
41+
let indexes = self.indexes.get(start..end)?;
42+
let datatype = unsafe { LogicalType::borrow(*self.datatypes.get(column)?) };
43+
Some(StructExtract { indexes, datatype })
44+
}
45+
}
46+
47+
/// Borrows a C++ array as a slice, reading a null pointer as an empty array.
48+
///
49+
/// # Safety
50+
///
51+
/// `ptr` must point to `len` initialized values that stay alive for `'a`.
52+
unsafe fn borrow_array<'a, T>(ptr: *const T, len: usize) -> &'a [T] {
53+
if ptr.is_null() {
54+
return &[];
55+
}
56+
unsafe { std::slice::from_raw_parts(ptr, len) }
57+
}
58+
1259
pub struct TableInitInput<'a> {
1360
pub input: &'a cpp::duckdb_vx_tfunc_init_input,
1461
}
@@ -52,4 +99,35 @@ impl<'a> TableInitInput<'a> {
5299
Some(unsafe { TableFilterSet::borrow(ptr) })
53100
}
54101
}
102+
103+
/// The struct extracts DuckDB pushed into the scan, one per column id.
104+
///
105+
/// The paths and types are borrowed from the bind data and only valid for
106+
/// the duration of the `init_global` call.
107+
pub fn struct_extracts(&self) -> StructExtracts<'_> {
108+
// The offsets array carries one entry per column id plus the end offset.
109+
let offsets_len = if self.input.column_extract_offsets.is_null() {
110+
0
111+
} else {
112+
self.input.column_ids_count + 1
113+
};
114+
115+
StructExtracts {
116+
indexes: unsafe {
117+
borrow_array(
118+
self.input.column_extract_indexes,
119+
self.input.column_extract_indexes_count,
120+
)
121+
},
122+
offsets: unsafe {
123+
borrow_array(self.input.column_extract_offsets, offsets_len)
124+
},
125+
datatypes: unsafe {
126+
borrow_array(
127+
self.input.column_extract_types,
128+
self.input.column_extract_types_count,
129+
)
130+
},
131+
}
132+
}
55133
}

‎vortex-duckdb/src/e2e_test/vortex_scan_test.rs‎

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1034,6 +1034,54 @@ fn test_geometry() {
10341034
assert_eq!(area, 1000.0);
10351035
}
10361036

1037+
/// `SELECT s.x.a` is pushed down as a struct extract: DuckDB asks the scan to read
1038+
/// the field at that child path instead of the whole `s` struct, and the scan
1039+
/// emits the extracted field directly.
1040+
///
1041+
/// Ignored until DuckDB's `MultiFileColumnMapper` stops rebuilding the struct for a
1042+
/// `PUSHDOWN_EXTRACT` column: `MultiFileReader::FinalizeChunk` currently applies a
1043+
/// struct expression to a vector DuckDB typed as the extracted field.
1044+
#[test]
1045+
#[ignore = "needs a DuckDB patch in MultiFileColumnMapper"]
1046+
fn test_vortex_scan_struct_extract_projection() {
1047+
let file = RUNTIME.block_on(async {
1048+
let inner = StructArray::try_from_iter([
1049+
("a", PrimitiveArray::from_iter([1i32, 2, 3]).into_array()),
1050+
("b", PrimitiveArray::from_iter([4i32, 5, 6]).into_array()),
1051+
])
1052+
.unwrap();
1053+
let top = StructArray::try_from_iter([
1054+
("x", inner.into_array()),
1055+
("y", PrimitiveArray::from_iter([7i32, 8, 9]).into_array()),
1056+
])
1057+
.unwrap();
1058+
1059+
write_single_column_vortex_file("s", top).await
1060+
});
1061+
1062+
let conn = database_connection();
1063+
let file_path = file.path().to_string_lossy();
1064+
1065+
// One level down, and two levels down, both requested by child index.
1066+
let result = conn
1067+
.query(&format!("SELECT s.y, s.x.b, s.x.a FROM '{file_path}'"))
1068+
.unwrap();
1069+
1070+
let mut y = Vec::new();
1071+
let mut xb = Vec::new();
1072+
let mut xa = Vec::new();
1073+
for chunk in result {
1074+
let len = chunk.len().as_();
1075+
y.extend_from_slice(chunk.get_vector(0).as_slice_with_len::<i32>(len));
1076+
xb.extend_from_slice(chunk.get_vector(1).as_slice_with_len::<i32>(len));
1077+
xa.extend_from_slice(chunk.get_vector(2).as_slice_with_len::<i32>(len));
1078+
}
1079+
1080+
assert_eq!(y, vec![7, 8, 9], "s.y mismatch");
1081+
assert_eq!(xb, vec![4, 5, 6], "s.x.b mismatch");
1082+
assert_eq!(xa, vec![1, 2, 3], "s.x.a mismatch");
1083+
}
1084+
10371085
/// `SELECT array_length(list)` / `len(list)` / `length(list)` should push the list-length
10381086
/// computation into the Vortex scan (computed from offsets, without materializing the list
10391087
/// elements) and return the per-row element counts.

0 commit comments

Comments
 (0)