Skip to content

Commit de6b0db

Browse files
committed
fix
1 parent 7103cb5 commit de6b0db

12 files changed

Lines changed: 202 additions & 269 deletions

File tree

Cargo.lock

Lines changed: 1 addition & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -647,8 +647,7 @@ databend-meta-client = { git = "https://github.com/databendlabs/databend-meta.gi
647647
databend-meta-test-harness = { git = "https://github.com/databendlabs/databend-meta.git", tag = "v260629.2.0" }
648648
deltalake = { git = "https://github.com/delta-io/delta-rs", rev = "9954bff" }
649649
#jsonb = { git = "https://github.com/databendlabs/jsonb.git", rev = "a16344417d1fec6df4a62d8e77679211e909f86e" }
650-
#jsonb = { git = "https://github.com/b41sh/jsonb.git", rev = "816c7ab394b5732ea621e9315723a34a96d9047c" }
651-
jsonb = { git = "https://github.com/b41sh/jsonb.git", rev = "ec1130a0f6a5e63c6a741e06bc4733e1a27c47d3" }
650+
jsonb = { git = "https://github.com/b41sh/jsonb.git", rev = "9567de9a167744aad45ccc21a68c55b6864fa391" }
652651
lance-arrow = { git = "https://github.com/datafuse-extras/lance", rev = "85f9401d3246d52a287bca21457dbae0466858ee" }
653652
lance-core = { git = "https://github.com/datafuse-extras/lance", rev = "85f9401d3246d52a287bca21457dbae0466858ee" }
654653
lance-encoding = { git = "https://github.com/datafuse-extras/lance", rev = "85f9401d3246d52a287bca21457dbae0466858ee" }

src/query/service/tests/it/storages/fuse/operations/virtual_columns.rs

Lines changed: 0 additions & 110 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,6 @@
1414

1515
use std::sync::Arc;
1616

17-
use databend_common_sql::plans::RefreshSelection;
1817
use databend_common_storage::read_parquet_schema_async_rs;
1918
use databend_common_storages_fuse::FUSE_TBL_VIRTUAL_BLOCK_PREFIX;
2019
use databend_common_storages_fuse::FUSE_TBL_VIRTUAL_BLOCK_PREFIX_V1;
@@ -196,115 +195,6 @@ async fn test_fuse_do_refresh_virtual_column() -> anyhow::Result<()> {
196195
Ok(())
197196
}
198197

199-
#[tokio::test(flavor = "multi_thread")]
200-
async fn test_refresh_virtual_column_block_selection_refreshes_owning_segment() -> anyhow::Result<()>
201-
{
202-
let fixture = TestFixture::setup().await?;
203-
fixture.create_default_database().await?;
204-
fixture.create_variant_table().await?;
205-
fixture
206-
.execute_command(&format!(
207-
"alter table {}.{} set options(block_per_segment = 100)",
208-
fixture.default_db_name(),
209-
fixture.default_table_name()
210-
))
211-
.await?;
212-
fixture
213-
.execute_command(&format!(
214-
"alter table {}.{} set options(enable_virtual_column = false)",
215-
fixture.default_db_name(),
216-
fixture.default_table_name()
217-
))
218-
.await?;
219-
append_variant_sample_data(2, &fixture).await?;
220-
fixture
221-
.execute_command(&format!(
222-
"alter table {}.{} set options(enable_virtual_column = true)",
223-
fixture.default_db_name(),
224-
fixture.default_table_name()
225-
))
226-
.await?;
227-
228-
let ctx = fixture.new_query_ctx().await?;
229-
let table = fixture.latest_default_table().await?;
230-
let fuse = FuseTable::try_from_table(table.as_ref())?;
231-
let snapshot = fuse.read_table_snapshot().await?.unwrap();
232-
let reader = MetaReaders::segment_info_reader(fuse.get_operator(), table.schema());
233-
let (segment_location, version) = &snapshot.segments[0];
234-
let segment = reader
235-
.read(&LoadParams {
236-
location: segment_location.clone(),
237-
len_hint: None,
238-
ver: *version,
239-
put_cache: false,
240-
})
241-
.await?;
242-
let blocks = segment.block_metas()?;
243-
assert!(blocks.len() > 1);
244-
let selected_block = blocks[0].location.0.clone();
245-
246-
let results = prepare_refresh_virtual_column(
247-
ctx,
248-
fuse,
249-
None,
250-
true,
251-
Some(RefreshSelection::BlockLocation(selected_block)),
252-
)
253-
.await?;
254-
assert_eq!(results.len(), blocks.len());
255-
Ok(())
256-
}
257-
258-
#[tokio::test(flavor = "multi_thread")]
259-
async fn test_refresh_virtual_column_limit_keeps_segment_whole() -> anyhow::Result<()> {
260-
let fixture = TestFixture::setup().await?;
261-
fixture.create_default_database().await?;
262-
fixture.create_variant_table().await?;
263-
fixture
264-
.execute_command(&format!(
265-
"alter table {}.{} set options(block_per_segment = 100)",
266-
fixture.default_db_name(),
267-
fixture.default_table_name()
268-
))
269-
.await?;
270-
fixture
271-
.execute_command(&format!(
272-
"alter table {}.{} set options(enable_virtual_column = false)",
273-
fixture.default_db_name(),
274-
fixture.default_table_name()
275-
))
276-
.await?;
277-
append_variant_sample_data(2, &fixture).await?;
278-
fixture
279-
.execute_command(&format!(
280-
"alter table {}.{} set options(enable_virtual_column = true)",
281-
fixture.default_db_name(),
282-
fixture.default_table_name()
283-
))
284-
.await?;
285-
286-
let ctx = fixture.new_query_ctx().await?;
287-
let table = fixture.latest_default_table().await?;
288-
let fuse = FuseTable::try_from_table(table.as_ref())?;
289-
let snapshot = fuse.read_table_snapshot().await?.unwrap();
290-
let reader = MetaReaders::segment_info_reader(fuse.get_operator(), table.schema());
291-
let (segment_location, version) = &snapshot.segments[0];
292-
let segment = reader
293-
.read(&LoadParams {
294-
location: segment_location.clone(),
295-
len_hint: None,
296-
ver: *version,
297-
put_cache: false,
298-
})
299-
.await?;
300-
let segment_block_count = segment.block_metas()?.len();
301-
assert!(segment_block_count > 1);
302-
303-
let results = prepare_refresh_virtual_column(ctx, fuse, Some(1), true, None).await?;
304-
assert_eq!(results.len(), segment_block_count);
305-
Ok(())
306-
}
307-
308198
// Inject a legacy _vb/ directory with a file. Vacuum prepare should not remove
309199
// files before the commit/cleanup phase runs.
310200
#[tokio::test(flavor = "multi_thread")]

src/query/storages/common/pruner/Cargo.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,8 @@ publish = { workspace = true }
77
edition = { workspace = true }
88

99
[dependencies]
10-
databend-common-exception = { workspace = true }
1110
databend-common-catalog = { workspace = true }
11+
databend-common-exception = { workspace = true }
1212
databend-common-expression = { workspace = true }
1313
databend-common-functions = { workspace = true }
1414
databend-storages-common-index = { workspace = true }

src/query/storages/fuse/src/pruning/virtual_column_pruner.rs

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -296,6 +296,16 @@ impl VirtualColumnPruner {
296296
let Some(path) =
297297
schema.find_path_ref(field.source_column_id, &virtual_column_field.encoded_path)
298298
else {
299+
// The parent itself may have no leaf while descendant paths are
300+
// materialized (for example only `geo.lat` exists in this block).
301+
// BlockMeta cannot reconstruct that object, so inspect the
302+
// sidecar trie instead of incorrectly declaring the parent missing.
303+
if schema.has_descendant_paths(
304+
field.source_column_id,
305+
&virtual_column_field.encoded_path,
306+
) {
307+
return None;
308+
}
299309
if virtual_block_meta.virtual_columns_complete {
300310
virtual_column_read_plan
301311
.insert(field.query_column_id, vec![VirtualColumnReadPlan::Missing]);

src/query/storages/fuse/src/table_functions/fuse_virtual_column_block_meta.rs

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -187,7 +187,9 @@ fn build_virtual_column_metas(
187187
let mut in_memory_sizes = Vec::with_capacity(virtual_column_metas.len());
188188
let mut column_stat_bitmap = MutableBitmap::with_capacity(virtual_column_metas.len());
189189

190-
for (column_id, virtual_column_meta) in virtual_column_metas {
190+
let mut sorted_metas = virtual_column_metas.iter().collect::<Vec<_>>();
191+
sorted_metas.sort_unstable_by_key(|(column_id, _)| **column_id);
192+
for (column_id, virtual_column_meta) in sorted_metas {
191193
let path_name = path_map
192194
.get(column_id)
193195
.map(|(source_column_id, path)| {
@@ -296,12 +298,16 @@ fn build_virtual_path_statistic(
296298
let mut path_names = Vec::new();
297299
let mut path_counts = Vec::new();
298300

299-
for (source_column_id, path_stats) in virtual_path_statistics {
301+
let mut sorted_statistics = virtual_path_statistics.iter().collect::<Vec<_>>();
302+
sorted_statistics.sort_unstable_by_key(|(source_column_id, _)| **source_column_id);
303+
for (source_column_id, path_stats) in sorted_statistics {
304+
let mut sorted_path_counts = path_stats.path_counts.iter().collect::<Vec<_>>();
305+
sorted_path_counts.sort_unstable_by_key(|(column_id, _)| *column_id);
300306
let source_name = source_column_names
301307
.get(source_column_id)
302308
.cloned()
303309
.unwrap_or_else(|| source_column_id.to_string());
304-
for (column_id, path_count) in &path_stats.path_counts {
310+
for (column_id, path_count) in sorted_path_counts {
305311
let path_name = path_map
306312
.get(column_id)
307313
.map(|(_, path)| format!("{}.{}", &source_name, &path))

src/query/storages/fuse/src/table_functions/fuse_virtual_column_build.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -220,6 +220,8 @@ async fn build_segment_virtual_columns(
220220
}
221221
let mut direct_paths = Vec::new();
222222
let max_direct_columns = policy.max_direct_columns;
223+
let mut path_counts = path_counts.into_iter().collect::<Vec<_>>();
224+
path_counts.sort_unstable_by_key(|(source_column_id, _)| *source_column_id);
223225
for (source_column_id, counts) in path_counts {
224226
let mut counts = counts.into_iter().collect::<Vec<_>>();
225227
counts.sort_by(|left, right| right.1.cmp(&left.1).then_with(|| left.0.cmp(&right.0)));

src/query/storages/fuse/src/table_functions/fuse_virtual_column_parquet_meta.rs

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -270,15 +270,22 @@ pub(crate) fn collect_virtual_column_entries(
270270
let mut entries = Vec::new();
271271
let mut shared_buckets = HashMap::new();
272272

273-
for (source_column_id, node) in &virtual_meta.virtual_column_nodes {
273+
let mut source_column_ids = virtual_meta
274+
.virtual_column_nodes
275+
.keys()
276+
.copied()
277+
.collect::<Vec<_>>();
278+
source_column_ids.sort_unstable();
279+
for source_column_id in source_column_ids {
280+
let node = &virtual_meta.virtual_column_nodes[&source_column_id];
274281
let source_column_name = source_column_names
275-
.get(source_column_id)
282+
.get(&source_column_id)
276283
.cloned()
277284
.unwrap_or_else(|| source_column_id.to_string());
278285
let mut key_paths = OwnedKeyPaths { paths: Vec::new() };
279286
collect_virtual_column_leaves(
280287
virtual_meta,
281-
*source_column_id,
288+
source_column_id,
282289
&source_column_name,
283290
node,
284291
&mut key_paths,
@@ -287,6 +294,9 @@ pub(crate) fn collect_virtual_column_entries(
287294
);
288295
}
289296

297+
let mut shared_buckets = shared_buckets.into_iter().collect::<Vec<_>>();
298+
shared_buckets
299+
.sort_by_key(|((source_column_id, data_type), _)| (*source_column_id, *data_type));
290300
for ((source_column_id, data_type), mut bucket) in shared_buckets {
291301
bucket.paths.sort();
292302
let metas = virtual_meta

src/query/storages/system/src/virtual_columns_table.rs

Lines changed: 29 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ use databend_common_expression::TableDataType;
2525
use databend_common_expression::TableField;
2626
use databend_common_expression::TableSchemaRefExt;
2727
use databend_common_expression::types::StringType;
28+
use databend_common_functions::BUILTIN_FUNCTIONS;
2829
use databend_common_meta_app::schema::TableIdent;
2930
use databend_common_meta_app::schema::TableInfo;
3031
use databend_common_meta_app::schema::TableMeta;
@@ -38,6 +39,7 @@ use jsonb::keypath::OwnedKeyPaths;
3839

3940
use crate::table::AsyncOneBlockSystemTable;
4041
use crate::table::AsyncSystemTable;
42+
use crate::util::extract_leveled_strings;
4143

4244
pub struct VirtualColumnsTable {
4345
table_info: TableInfo,
@@ -54,8 +56,22 @@ impl AsyncSystemTable for VirtualColumnsTable {
5456
async fn get_full_data(
5557
&self,
5658
ctx: Arc<dyn TableContext>,
57-
_push_downs: Option<PushDownInfo>,
59+
push_downs: Option<PushDownInfo>,
5860
) -> Result<DataBlock> {
61+
let mut filtered_db_names = Vec::new();
62+
let mut filtered_table_names = Vec::new();
63+
if let Some(filters) = push_downs.and_then(|push_down| push_down.filters) {
64+
let expr = filters.filter.as_expr(&BUILTIN_FUNCTIONS);
65+
(filtered_db_names, filtered_table_names) = extract_leveled_strings(
66+
&expr,
67+
&["database", "table"],
68+
&ctx.get_function_context()?,
69+
)?;
70+
}
71+
let filtered_db_names = (!filtered_db_names.is_empty()).then_some(filtered_db_names);
72+
let filtered_table_names =
73+
(!filtered_table_names.is_empty()).then_some(filtered_table_names);
74+
5975
let tenant = ctx.get_tenant();
6076
let session_state = ctx.session_state()?;
6177

@@ -70,8 +86,20 @@ impl AsyncSystemTable for VirtualColumnsTable {
7086

7187
let dbs = catalog.list_databases(&tenant).await?;
7288
for db in dbs {
89+
if filtered_db_names
90+
.as_ref()
91+
.is_some_and(|names| !names.iter().any(|name| name == db.name()))
92+
{
93+
continue;
94+
}
7395
let tables = catalog.list_tables(&tenant, db.name()).await?;
7496
for table in tables {
97+
if filtered_table_names
98+
.as_ref()
99+
.is_some_and(|names| !names.iter().any(|name| name == table.name()))
100+
{
101+
continue;
102+
}
75103
if !table.storage_format_as_parquet() {
76104
continue;
77105
}

tests/sqllogictests/suites/ee/01_ee_system/01_0002_virtual_column.test

Lines changed: 11 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -56,13 +56,13 @@ statement ok
5656
REFRESH VIRTUAL COLUMN FOR vacuum_test;
5757

5858
# Schema has 4 fields (a, b, c, d), 2 _vb_v2 files (one per block)
59-
query TTTITT
59+
query TTTTT
6060
show virtual columns from vacuum_test;
6161
----
62-
test_vacuum_virtual_column vacuum_test v 3000000000 ['a'] UInt64
63-
test_vacuum_virtual_column vacuum_test v 3000000001 ['b'] UInt64
64-
test_vacuum_virtual_column vacuum_test v 3000000002 ['c'] UInt64
65-
test_vacuum_virtual_column vacuum_test v 3000000003 ['d'] UInt64
62+
test_vacuum_virtual_column vacuum_test v a UInt64
63+
test_vacuum_virtual_column vacuum_test v b UInt64
64+
test_vacuum_virtual_column vacuum_test v c UInt64
65+
test_vacuum_virtual_column vacuum_test v d UInt64
6666

6767
query I
6868
select count() from list_stage(location=> '@vb_stage') where name like '%_vb_v2%';
@@ -78,13 +78,11 @@ DELETE FROM vacuum_test WHERE id >= 4;
7878
statement ok
7979
REFRESH VIRTUAL COLUMN FOR vacuum_test;
8080

81-
query TTTITT
81+
query TTTTT
8282
show virtual columns from vacuum_test;
8383
----
84-
test_vacuum_virtual_column vacuum_test v 3000000000 ['a'] UInt64
85-
test_vacuum_virtual_column vacuum_test v 3000000001 ['b'] UInt64
86-
test_vacuum_virtual_column vacuum_test v 3000000002 ['c'] UInt64
87-
test_vacuum_virtual_column vacuum_test v 3000000003 ['d'] UInt64
84+
test_vacuum_virtual_column vacuum_test v a UInt64
85+
test_vacuum_virtual_column vacuum_test v b UInt64
8886

8987
# Purge old snapshots so orphan _vb_v2 files are no longer referenced.
9088
statement ok
@@ -97,11 +95,11 @@ VACUUM VIRTUAL COLUMN FROM vacuum_test;
9795
1
9896

9997
# After vacuum: only a and b remain in schema
100-
query TTTITT
98+
query TTTTT
10199
show virtual columns from vacuum_test;
102100
----
103-
test_vacuum_virtual_column vacuum_test v 3000000000 ['a'] UInt64
104-
test_vacuum_virtual_column vacuum_test v 3000000001 ['b'] UInt64
101+
test_vacuum_virtual_column vacuum_test v a UInt64
102+
test_vacuum_virtual_column vacuum_test v b UInt64
105103

106104
# Orphan _vb_v2 file removed, only 1 remains
107105
query I

0 commit comments

Comments
 (0)