Skip to content

Commit 7d2603c

Browse files
committed
fix: preserve dotted column relation qualifiers
1 parent a0631ed commit 7d2603c

6 files changed

Lines changed: 110 additions & 2 deletions

File tree

‎datafusion/proto-common/proto/datafusion_common.proto‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ package datafusion_common;
2222

2323
message ColumnRelation {
2424
string relation = 1;
25+
repeated string parts = 2;
2526
}
2627

2728
message Column {

‎datafusion/proto-common/src/from_proto/mod.rs‎

Lines changed: 85 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -145,7 +145,14 @@ where
145145

146146
impl From<protobuf::ColumnRelation> for TableReference {
147147
fn from(rel: protobuf::ColumnRelation) -> Self {
148-
Self::parse_str_normalized(rel.relation.as_str(), true)
148+
match rel.parts.as_slice() {
149+
[table] => Self::bare(table.as_str()),
150+
[schema, table] => Self::partial(schema.as_str(), table.as_str()),
151+
[catalog, schema, table] => {
152+
Self::full(catalog.as_str(), schema.as_str(), table.as_str())
153+
}
154+
_ => Self::parse_str_normalized(rel.relation.as_str(), true),
155+
}
149156
}
150157
}
151158

@@ -1386,6 +1393,7 @@ pub(crate) fn csv_writer_options_from_proto(
13861393

13871394
#[cfg(test)]
13881395
mod tests {
1396+
use datafusion_common::TableReference;
13891397
use datafusion_common::config::{
13901398
MaxRowGroupBytes, ParquetCdcOptions, ParquetOptions, TableParquetOptions,
13911399
};
@@ -1405,6 +1413,82 @@ mod tests {
14051413
);
14061414
}
14071415

1416+
#[test]
1417+
fn column_relation_round_trip_preserves_dotted_bare_table() {
1418+
let column = datafusion_common::Column::new(
1419+
Some(TableReference::bare("has.dot")),
1420+
"column",
1421+
);
1422+
1423+
let proto: crate::protobuf_common::Column = (&column).into();
1424+
let relation = proto.relation.expect("relation should be present");
1425+
1426+
assert_eq!(relation.relation, "has.dot");
1427+
assert_eq!(relation.parts, vec!["has.dot".to_string()]);
1428+
1429+
let recovered = TableReference::from(relation);
1430+
assert_eq!(recovered, TableReference::bare("has.dot"));
1431+
}
1432+
1433+
#[test]
1434+
fn column_relation_round_trip_preserves_dotted_partial_reference() {
1435+
let column = datafusion_common::Column::new(
1436+
Some(TableReference::partial("my.schema", "table")),
1437+
"column",
1438+
);
1439+
1440+
let proto: crate::protobuf_common::Column = (&column).into();
1441+
let relation = proto.relation.expect("relation should be present");
1442+
1443+
assert_eq!(relation.relation, "my.schema.table");
1444+
assert_eq!(
1445+
relation.parts,
1446+
vec!["my.schema".to_string(), "table".to_string()]
1447+
);
1448+
1449+
let recovered = TableReference::from(relation);
1450+
assert_eq!(recovered, TableReference::partial("my.schema", "table"));
1451+
}
1452+
1453+
#[test]
1454+
fn column_relation_round_trip_preserves_dotted_full_reference() {
1455+
let column = datafusion_common::Column::new(
1456+
Some(TableReference::full("catalog", "my.schema", "table")),
1457+
"column",
1458+
);
1459+
1460+
let proto: crate::protobuf_common::Column = (&column).into();
1461+
let relation = proto.relation.expect("relation should be present");
1462+
1463+
assert_eq!(relation.relation, "catalog.my.schema.table");
1464+
assert_eq!(
1465+
relation.parts,
1466+
vec![
1467+
"catalog".to_string(),
1468+
"my.schema".to_string(),
1469+
"table".to_string()
1470+
]
1471+
);
1472+
1473+
let recovered = TableReference::from(relation);
1474+
assert_eq!(
1475+
recovered,
1476+
TableReference::full("catalog", "my.schema", "table")
1477+
);
1478+
}
1479+
1480+
#[test]
1481+
fn column_relation_decodes_legacy_relation() {
1482+
let proto = crate::protobuf_common::ColumnRelation {
1483+
relation: "schema.table".to_string(),
1484+
parts: vec![],
1485+
};
1486+
1487+
let recovered = TableReference::from(proto);
1488+
1489+
assert_eq!(recovered, TableReference::partial("schema", "table"));
1490+
}
1491+
14081492
#[test]
14091493
fn table_parquet_options_defaults_missing_global() {
14101494
let recovered = TableParquetOptions::try_from(

‎datafusion/proto-common/src/generated/pbjson.rs‎

Lines changed: 18 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1030,10 +1030,16 @@ impl serde::Serialize for ColumnRelation {
10301030
if !self.relation.is_empty() {
10311031
len += 1;
10321032
}
1033+
if !self.parts.is_empty() {
1034+
len += 1;
1035+
}
10331036
let mut struct_ser = serializer.serialize_struct("datafusion_common.ColumnRelation", len)?;
10341037
if !self.relation.is_empty() {
10351038
struct_ser.serialize_field("relation", &self.relation)?;
10361039
}
1040+
if !self.parts.is_empty() {
1041+
struct_ser.serialize_field("parts", &self.parts)?;
1042+
}
10371043
struct_ser.end()
10381044
}
10391045
}
@@ -1045,11 +1051,13 @@ impl<'de> serde::Deserialize<'de> for ColumnRelation {
10451051
{
10461052
const FIELDS: &[&str] = &[
10471053
"relation",
1054+
"parts",
10481055
];
10491056

10501057
#[allow(clippy::enum_variant_names)]
10511058
enum GeneratedField {
10521059
Relation,
1060+
Parts,
10531061
}
10541062
impl<'de> serde::Deserialize<'de> for GeneratedField {
10551063
fn deserialize<D>(deserializer: D) -> std::result::Result<GeneratedField, D::Error>
@@ -1072,6 +1080,7 @@ impl<'de> serde::Deserialize<'de> for ColumnRelation {
10721080
{
10731081
match value {
10741082
"relation" => Ok(GeneratedField::Relation),
1083+
"parts" => Ok(GeneratedField::Parts),
10751084
_ => Err(serde::de::Error::unknown_field(value, FIELDS)),
10761085
}
10771086
}
@@ -1092,6 +1101,7 @@ impl<'de> serde::Deserialize<'de> for ColumnRelation {
10921101
V: serde::de::MapAccess<'de>,
10931102
{
10941103
let mut relation__ = None;
1104+
let mut parts__ = None;
10951105
while let Some(k) = map_.next_key()? {
10961106
match k {
10971107
GeneratedField::Relation => {
@@ -1100,10 +1110,17 @@ impl<'de> serde::Deserialize<'de> for ColumnRelation {
11001110
}
11011111
relation__ = Some(map_.next_value()?);
11021112
}
1113+
GeneratedField::Parts => {
1114+
if parts__.is_some() {
1115+
return Err(serde::de::Error::duplicate_field("parts"));
1116+
}
1117+
parts__ = Some(map_.next_value()?);
1118+
}
11031119
}
11041120
}
11051121
Ok(ColumnRelation {
11061122
relation: relation__.unwrap_or_default(),
1123+
parts: parts__.unwrap_or_default(),
11071124
})
11081125
}
11091126
}
@@ -3997,7 +4014,7 @@ impl serde::Serialize for ExplainAnalyzeCategoriesNode {
39974014
struct_ser.serialize_field("all", &self.all)?;
39984015
}
39994016
if !self.only.is_empty() {
4000-
let v = self.only.iter().copied().map(|v| {
4017+
let v = self.only.iter().cloned().map(|v| {
40014018
MetricCategory::try_from(v)
40024019
.map_err(|_| serde::ser::Error::custom(format!("Invalid variant {}", v)))
40034020
}).collect::<std::result::Result<Vec<_>, _>>()?;

‎datafusion/proto-common/src/generated/prost.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@
33
pub struct ColumnRelation {
44
#[prost(string, tag = "1")]
55
pub relation: ::prost::alloc::string::String,
6+
#[prost(string, repeated, tag = "2")]
7+
pub parts: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
68
}
79
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
810
pub struct Column {

‎datafusion/proto-common/src/to_proto/mod.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -247,6 +247,7 @@ impl From<Column> for protobuf::Column {
247247
Self {
248248
relation: c.relation.map(|relation| protobuf::ColumnRelation {
249249
relation: relation.to_string(),
250+
parts: relation.to_vec(),
250251
}),
251252
name: c.name,
252253
}
@@ -292,6 +293,7 @@ impl TryFrom<&DFSchema> for protobuf::DfSchema {
292293
field: Some(field.as_ref().try_into()?),
293294
qualifier: qualifier.map(|r| protobuf::ColumnRelation {
294295
relation: r.to_string(),
296+
parts: r.to_vec(),
295297
}),
296298
})
297299
})

‎datafusion/proto-models/src/generated/datafusion_proto_common.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@
33
pub struct ColumnRelation {
44
#[prost(string, tag = "1")]
55
pub relation: ::prost::alloc::string::String,
6+
#[prost(string, repeated, tag = "2")]
7+
pub parts: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
68
}
79
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
810
pub struct Column {

0 commit comments

Comments
 (0)