Skip to content

Commit 6a1ee17

Browse files
committed
fix
1 parent de6b0db commit 6a1ee17

14 files changed

Lines changed: 188 additions & 173 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
@@ -646,8 +646,7 @@ databend-meta = { git = "https://github.com/databendlabs/databend-meta.git", tag
646646
databend-meta-client = { git = "https://github.com/databendlabs/databend-meta.git", tag = "v260205.13.2" }
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" }
649-
#jsonb = { git = "https://github.com/databendlabs/jsonb.git", rev = "a16344417d1fec6df4a62d8e77679211e909f86e" }
650-
jsonb = { git = "https://github.com/b41sh/jsonb.git", rev = "9567de9a167744aad45ccc21a68c55b6864fa391" }
649+
jsonb = { git = "https://github.com/databendlabs/jsonb.git", rev = "fba895c" }
651650
lance-arrow = { git = "https://github.com/datafuse-extras/lance", rev = "85f9401d3246d52a287bca21457dbae0466858ee" }
652651
lance-core = { git = "https://github.com/datafuse-extras/lance", rev = "85f9401d3246d52a287bca21457dbae0466858ee" }
653652
lance-encoding = { git = "https://github.com/datafuse-extras/lance", rev = "85f9401d3246d52a287bca21457dbae0466858ee" }

src/query/catalog/src/plan/pushdown.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -62,7 +62,7 @@ pub struct VirtualColumnField {
6262
/// virtual column id; each segment may assign a different real column id
6363
/// for the same path, so readers must map this id through the segment schema.
6464
pub query_column_id: u32,
65-
/// Full query field name, including the source column, e.g. `v.user.name` or `v[0].id`.
65+
/// Full query field name using bracket path notation.
6666
pub name: String,
6767
/// Paths to generate virtual column from source column.
6868
pub key_paths: OwnedKeyPaths,
@@ -77,7 +77,7 @@ pub struct VirtualColumnField {
7777
/// source-column/canonical-path identity used to locate segment-local stats.
7878
#[derive(Clone, Debug)]
7979
pub struct VirtualPredicateRef {
80-
/// Full query/pipeline field name used by expression column references.
80+
/// Unique internal query/pipeline field name used by expression column references.
8181
pub name: String,
8282
/// Id of the authoritative source Variant column.
8383
pub source_column_id: u32,

src/query/service/src/physical_plans/format/common.rs

Lines changed: 14 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -281,22 +281,22 @@ pub fn format_output_columns(
281281
return String::from("dummy value");
282282
}
283283
let column_entry = metadata.column(column_index);
284+
let column_name = column_entry.name();
284285
match column_entry.table_index() {
285-
Some(table_index) if format_table => match metadata
286-
.table(table_index)
287-
.alias_name()
288-
{
289-
Some(alias_name) => {
290-
format!("{}.{} (#{})", alias_name, column_entry.name(), column_index)
286+
Some(table_index) if format_table => {
287+
match metadata.table(table_index).alias_name() {
288+
Some(alias_name) => {
289+
format!("{}.{} (#{})", alias_name, column_name, column_index)
290+
}
291+
None => format!(
292+
"{}.{} (#{})",
293+
metadata.table(table_index).name(),
294+
column_name,
295+
column_index,
296+
),
291297
}
292-
None => format!(
293-
"{}.{} (#{})",
294-
metadata.table(table_index).name(),
295-
column_entry.name(),
296-
column_index,
297-
),
298-
},
299-
_ => format!("{} (#{})", column_entry.name(), column_index),
298+
}
299+
_ => format!("{} (#{})", column_name, column_index),
300300
}
301301
}
302302
_ => format!("#{}", field.name()),

src/query/service/src/physical_plans/format/format_table_scan.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -85,7 +85,7 @@ impl<'a> PhysicalFormat for TableScanFormatter<'a> {
8585
let mut names = virtual_column
8686
.virtual_column_fields
8787
.iter()
88-
.map(|c| c.name.clone())
88+
.map(|column| column.name.clone())
8989
.collect::<Vec<_>>();
9090
names.sort();
9191
names.iter().join(", ")

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

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ use databend_storages_common_table_meta::meta::VirtualSegmentColumnPath;
4343
use databend_storages_common_table_meta::meta::VirtualSegmentPath;
4444
use databend_storages_common_table_meta::meta::VirtualSegmentSchema;
4545
use jsonb::OwnedJsonb;
46+
use jsonb::keypath::OwnedKeyPath;
4647
use jsonb::keypath::OwnedKeyPaths;
4748
use jsonb::keypath::parse_key_paths;
4849

@@ -456,20 +457,26 @@ fn build_virtual_column_field(
456457
query_column_id: u32,
457458
key_paths: OwnedKeyPaths,
458459
) -> VirtualColumnField {
459-
let name = format_virtual_column_name(source_name, &key_paths);
460460
VirtualColumnField {
461461
source_column_id,
462462
source_name: source_name.to_string(),
463463
query_column_id,
464-
name,
464+
name: format_virtual_column_name(source_name, &key_paths),
465465
key_paths,
466466
cast_func_name: None,
467467
data_type: Box::new(TableDataType::Variant),
468468
}
469469
}
470470

471471
fn format_virtual_column_name(source: &str, key_paths: &OwnedKeyPaths) -> String {
472-
format!("{source}.{}", key_paths.to_canonical_path())
472+
let mut name = source.to_string();
473+
for path in &key_paths.paths {
474+
match path {
475+
OwnedKeyPath::Index(index) => name.push_str(&format!("[{index}]")),
476+
OwnedKeyPath::Name(key) => name.push_str(&format!("['{key}']")),
477+
}
478+
}
479+
name
473480
}
474481

475482
fn collect_plan_kinds(plan: &VirtualColumnReadPlan, kinds: &mut HashSet<&'static str>) {

src/query/sql/src/planner/binder/bind_context.rs

Lines changed: 13 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@ use databend_common_expression::infer_schema_type;
4444
use databend_common_meta_app::principal::UserDefinedFunction;
4545
use enum_as_inner::EnumAsInner;
4646
use indexmap::IndexMap;
47+
use jsonb::keypath::OwnedKeyPath;
4748
use jsonb::keypath::OwnedKeyPaths;
4849
use parking_lot::RwLock;
4950

@@ -1002,12 +1003,7 @@ impl BindContext {
10021003
});
10031004

10041005
let source_column_id = virtual_column_name.source_column_id;
1005-
let path_name = &virtual_column_name.key_name;
1006-
let column_name = if path_name.starts_with('[') {
1007-
format!("{source_column_name}{path_name}")
1008-
} else {
1009-
format!("{source_column_name}.{path_name}")
1010-
};
1006+
let column_name = format_virtual_column_name(source_column_name, &key_paths);
10111007
// todo
10121008
let table_data_type = TableDataType::Nullable(Box::new(TableDataType::Variant));
10131009
let is_try = true;
@@ -1113,3 +1109,14 @@ pub fn apply_alias_for_columns(
11131109
}
11141110
Ok(())
11151111
}
1112+
1113+
fn format_virtual_column_name(source_column_name: &str, key_paths: &OwnedKeyPaths) -> String {
1114+
let mut name = source_column_name.to_string();
1115+
for path in &key_paths.paths {
1116+
match path {
1117+
OwnedKeyPath::Index(index) => name.push_str(&format!("[{index}]")),
1118+
OwnedKeyPath::Name(key) => name.push_str(&format!("['{key}']")),
1119+
}
1120+
}
1121+
name
1122+
}

src/query/sql/src/planner/metadata/metadata.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -856,7 +856,7 @@ pub struct VirtualColumn {
856856
/// Query-time temporary column id.
857857
pub query_column_id: u32,
858858
pub column_index: Symbol,
859-
/// Full query/pipeline name including the source column, e.g. `v.user.name` or `v[0]`.
859+
/// Full query/pipeline name using bracket path notation.
860860
pub column_name: String,
861861
pub key_paths: OwnedKeyPaths,
862862
pub data_type: TableDataType,

src/query/sql/tests/it/semantic/binder_materialized_cte_virtual_column.txt

Lines changed: 15 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -13,30 +13,30 @@ Sequence(Sequence)
1313
│ ├── ref_count: 0
1414
│ ├── channel_size: None
1515
│ └── EvalScalar
16-
│ ├── scalars: [default.t.v.message (#1) AS (#1), default.t.v.message.attribute.account_id (#2) AS (#2), default.t.v.message.attribute.user_id (#3) AS (#3)]
16+
│ ├── scalars: [default.t.v['message'] (#1) AS (#1), default.t.v['message']['attribute']['account_id'] (#2) AS (#2), default.t.v['message']['attribute']['user_id'] (#3) AS (#3)]
1717
│ └── Scan
1818
│ ├── table: default.t (#0)
1919
│ ├── filters: []
2020
│ ├── order by: []
2121
│ └── limit: NONE
2222
└── UnionAll
2323
├── output: [__databend_virtual_column__0 (#12)]
24-
├── left: [default.t.v.message.attribute.user_id (#7)]
25-
├── right: [default.t.v.message.attribute.account_id (#10)]
24+
├── left: [default.t.v['message']['attribute']['user_id'] (#7)]
25+
├── right: [default.t.v['message']['attribute']['account_id'] (#10)]
2626
├── cte_scan_names: []
2727
├── logical_recursive_cte_id: None
2828
├── EvalScalar
29-
│ ├── scalars: [default.t.v.message.attribute.user_id (#7) AS (#7)]
29+
│ ├── scalars: [default.t.v['message']['attribute']['user_id'] (#7) AS (#7)]
3030
│ └── MaterializedCTERef
3131
│ ├── cte_name: __materialized_cte_0_logs
32-
│ ├── output columns: [default.t.v.message (#5), default.t.v.message.attribute.account_id (#6), default.t.v.message.attribute.user_id (#7)]
33-
│ └── column mapping: [default.t.v.message (#5) -> default.t.v.message (#1), default.t.v.message.attribute.account_id (#6) -> default.t.v.message.attribute.account_id (#2), default.t.v.message.attribute.user_id (#7) -> default.t.v.message.attribute.user_id (#3)]
32+
│ ├── output columns: [default.t.v['message'] (#5), default.t.v['message']['attribute']['account_id'] (#6), default.t.v['message']['attribute']['user_id'] (#7)]
33+
│ └── column mapping: [default.t.v['message'] (#5) -> default.t.v['message'] (#1), default.t.v['message']['attribute']['account_id'] (#6) -> default.t.v['message']['attribute']['account_id'] (#2), default.t.v['message']['attribute']['user_id'] (#7) -> default.t.v['message']['attribute']['user_id'] (#3)]
3434
└── EvalScalar
35-
├── scalars: [default.t.v.message.attribute.account_id (#10) AS (#10)]
35+
├── scalars: [default.t.v['message']['attribute']['account_id'] (#10) AS (#10)]
3636
└── MaterializedCTERef
3737
├── cte_name: __materialized_cte_0_logs
38-
├── output columns: [default.t.v.message (#9), default.t.v.message.attribute.account_id (#10), default.t.v.message.attribute.user_id (#11)]
39-
└── column mapping: [default.t.v.message (#9) -> default.t.v.message (#1), default.t.v.message.attribute.account_id (#10) -> default.t.v.message.attribute.account_id (#2), default.t.v.message.attribute.user_id (#11) -> default.t.v.message.attribute.user_id (#3)]
38+
├── output columns: [default.t.v['message'] (#9), default.t.v['message']['attribute']['account_id'] (#10), default.t.v['message']['attribute']['user_id'] (#11)]
39+
└── column mapping: [default.t.v['message'] (#9) -> default.t.v['message'] (#1), default.t.v['message']['attribute']['account_id'] (#10) -> default.t.v['message']['attribute']['account_id'] (#2), default.t.v['message']['attribute']['user_id'] (#11) -> default.t.v['message']['attribute']['user_id'] (#3)]
4040

4141

4242
=== materialized_cte_virtual_column_rewrites_chained_ctes ===
@@ -50,13 +50,13 @@ sql:
5050

5151
status: ok
5252
EvalScalar
53-
├── scalars: [default.t.v.message.attribute.user_id (#3) AS (#3), default.t.v.message.attribute.name (#4) AS (#4)]
53+
├── scalars: [default.t.v['message']['attribute']['user_id'] (#3) AS (#3), default.t.v['message']['attribute']['name'] (#4) AS (#4)]
5454
└── EvalScalar
55-
├── scalars: [default.t.v.message.attribute.user_id (#3) AS (#3), default.t.v.message.attribute.name (#4) AS (#4)]
55+
├── scalars: [default.t.v['message']['attribute']['user_id'] (#3) AS (#3), default.t.v['message']['attribute']['name'] (#4) AS (#4)]
5656
└── EvalScalar
57-
├── scalars: [default.t.v.message.attribute (#2) AS (#2)]
57+
├── scalars: [default.t.v['message']['attribute'] (#2) AS (#2)]
5858
└── EvalScalar
59-
├── scalars: [default.t.v.message (#1) AS (#1)]
59+
├── scalars: [default.t.v['message'] (#1) AS (#1)]
6060
└── Scan
6161
├── table: default.t (#0)
6262
├── filters: []
@@ -80,9 +80,9 @@ status: ok
8080
EvalScalar
8181
├── scalars: [user_id (#8) AS (#8), trace_id (#9) AS (#9), level (#10) AS (#10), status_code (#11) AS (#11), error_message (#12) AS (#12)]
8282
└── EvalScalar
83-
├── scalars: [CAST(default.t.v.message.attribute.user_id (#3) AS Int32 NULL) AS (#8), CAST(default.t.v.message.attribute.trace_id (#4) AS Int32 NULL) AS (#9), CAST(default.t.v.message.attribute.level (#5) AS String NULL) AS (#10), CAST(default.t.v.response.status_code (#6) AS Int64 NULL) AS (#11), CAST(default.t.v.response.error_message (#7) AS String NULL) AS (#12)]
83+
├── scalars: [CAST(default.t.v['message']['attribute']['user_id'] (#3) AS Int32 NULL) AS (#8), CAST(default.t.v['message']['attribute']['trace_id'] (#4) AS Int32 NULL) AS (#9), CAST(default.t.v['message']['attribute']['level'] (#5) AS String NULL) AS (#10), CAST(default.t.v['response']['status_code'] (#6) AS Int64 NULL) AS (#11), CAST(default.t.v['response']['error_message'] (#7) AS String NULL) AS (#12)]
8484
└── EvalScalar
85-
├── scalars: [default.t.v.message (#1) AS (#1), default.t.v.response (#2) AS (#2)]
85+
├── scalars: [default.t.v['message'] (#1) AS (#1), default.t.v['response'] (#2) AS (#2)]
8686
└── Scan
8787
├── table: default.t (#0)
8888
├── filters: []

src/query/sql/tests/it/semantic/type_check/variant.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -86,8 +86,8 @@ async fn nested_get_virtual_column_rewrite_skips_intermediate_paths() -> Result<
8686

8787
let metadata = metadata.read();
8888
assert_eq!(metadata.columns().len(), 4);
89-
assert_virtual_column(metadata.column(Symbol::new(2)), "v.a[0]");
90-
assert_virtual_column(metadata.column(Symbol::new(3)), "v.b.c");
89+
assert_virtual_column(metadata.column(Symbol::new(2)), "v['a'][0]");
90+
assert_virtual_column(metadata.column(Symbol::new(3)), "v['b']['c']");
9191

9292
let case = SqlTestCase {
9393
name: "nested_get_virtual_column_rewrite_skips_intermediate_paths",

0 commit comments

Comments
 (0)