Skip to content
Draft
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
20 changes: 14 additions & 6 deletions datafusion/physical-expr/src/expressions/column.rs
Original file line number Diff line number Diff line change
Expand Up @@ -153,26 +153,34 @@ impl PhysicalExpr for Column {
_ctx: &datafusion_physical_expr_common::physical_expr::proto_encode::PhysicalExprEncodeCtx<'_>,
) -> Result<Option<datafusion_proto_models::protobuf::PhysicalExprNode>> {
use datafusion_proto_models::protobuf;
let Self { name, index } = self;
Ok(Some(protobuf::PhysicalExprNode {
expr_id: None,
expr_type: Some(protobuf::physical_expr_node::ExprType::Column(self.into())),
expr_type: Some(protobuf::physical_expr_node::ExprType::Column(
protobuf::PhysicalColumn {
name: name.clone(),
index: *index as u32,
},
)),
}))
}
}

#[cfg(feature = "proto")]
impl From<&datafusion_proto_models::protobuf::PhysicalColumn> for Column {
fn from(c: &datafusion_proto_models::protobuf::PhysicalColumn) -> Self {
Column::new(&c.name, c.index as usize)
let datafusion_proto_models::protobuf::PhysicalColumn { name, index } = c;
Column::new(name, *index as usize)
}
}

#[cfg(feature = "proto")]
impl From<&Column> for datafusion_proto_models::protobuf::PhysicalColumn {
fn from(c: &Column) -> Self {
let Column { name, index } = c;
Self {
name: c.name.clone(),
index: c.index as u32,
name: name.clone(),
index: *index as u32,
}
}
}
Expand All @@ -196,12 +204,12 @@ impl Column {
) -> Result<Arc<dyn PhysicalExpr>> {
use datafusion_physical_expr_common::expect_expr_variant;
use datafusion_proto_models::protobuf;
let column = expect_expr_variant!(
let protobuf::PhysicalColumn { name, index } = expect_expr_variant!(
node,
protobuf::physical_expr_node::ExprType::Column,
"Column",
);
Ok(Arc::new(Column::from(column)))
Ok(Arc::new(Column::new(name, *index as usize)))
}
}

Expand Down
11 changes: 5 additions & 6 deletions datafusion/physical-expr/src/expressions/is_not_null.rs
Original file line number Diff line number Diff line change
Expand Up @@ -110,11 +110,12 @@ impl PhysicalExpr for IsNotNullExpr {
) -> Result<Option<datafusion_proto_models::protobuf::PhysicalExprNode>> {
use datafusion_proto_models::protobuf;

let Self { arg } = self;
Ok(Some(protobuf::PhysicalExprNode {
expr_id: None,
expr_type: Some(protobuf::physical_expr_node::ExprType::IsNotNullExpr(
Box::new(protobuf::PhysicalIsNotNull {
expr: Some(Box::new(ctx.encode_child(&self.arg)?)),
expr: Some(Box::new(ctx.encode_child(arg)?)),
}),
)),
}))
Expand All @@ -136,11 +137,9 @@ impl IsNotNullExpr {
protobuf::physical_expr_node::ExprType::IsNotNullExpr,
"IsNotNullExpr",
);
let expr = ctx.decode_required_expression(
node.expr.as_deref(),
"IsNotNullExpr",
"expr",
)?;
let protobuf::PhysicalIsNotNull { expr } = node.as_ref();
let expr =
ctx.decode_required_expression(expr.as_deref(), "IsNotNullExpr", "expr")?;

Ok(Arc::new(IsNotNullExpr::new(expr)))
}
Expand Down
6 changes: 4 additions & 2 deletions datafusion/physical-expr/src/expressions/is_null.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,11 +109,12 @@ impl PhysicalExpr for IsNullExpr {
) -> Result<Option<datafusion_proto_models::protobuf::PhysicalExprNode>> {
use datafusion_proto_models::protobuf;

let Self { arg } = self;
Ok(Some(protobuf::PhysicalExprNode {
expr_id: None,
expr_type: Some(protobuf::physical_expr_node::ExprType::IsNullExpr(
Box::new(protobuf::PhysicalIsNull {
expr: Some(Box::new(ctx.encode_child(&self.arg)?)),
expr: Some(Box::new(ctx.encode_child(arg)?)),
}),
)),
}))
Expand All @@ -135,8 +136,9 @@ impl IsNullExpr {
protobuf::physical_expr_node::ExprType::IsNullExpr,
"IsNullExpr",
);
let protobuf::PhysicalIsNull { expr } = node.as_ref();
let expr =
ctx.decode_required_expression(node.expr.as_deref(), "IsNullExpr", "expr")?;
ctx.decode_required_expression(expr.as_deref(), "IsNullExpr", "expr")?;

Ok(Arc::new(IsNullExpr::new(expr)))
}
Expand Down
76 changes: 64 additions & 12 deletions datafusion/physical-expr/src/expressions/literal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -143,11 +143,21 @@ impl PhysicalExpr for Literal {
) -> Result<Option<datafusion_proto_models::protobuf::PhysicalExprNode>> {
use datafusion_proto_models::protobuf;

let Self { value, field } = self;
// The field name, type, and nullability are reconstructed by new_with_metadata.
let expr_type = if field.metadata().is_empty() {
protobuf::physical_expr_node::ExprType::Literal(value.try_into()?)
} else {
protobuf::physical_expr_node::ExprType::LiteralWithMetadata(
protobuf::PhysicalLiteralNode {
value: Some(value.try_into()?),
metadata: field.metadata().clone(),
},
)
};
Ok(Some(protobuf::PhysicalExprNode {
expr_id: None,
expr_type: Some(protobuf::physical_expr_node::ExprType::Literal(
(&self.value).try_into()?,
)),
expr_type: Some(expr_type),
}))
}
}
Expand All @@ -159,16 +169,38 @@ impl Literal {
node: &datafusion_proto_models::protobuf::PhysicalExprNode,
_ctx: &datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecodeCtx<'_>,
) -> Result<Arc<dyn PhysicalExpr>> {
use datafusion_physical_expr_common::expect_expr_variant;
use datafusion_common::{internal_datafusion_err, internal_err};
use datafusion_proto_models::protobuf;

let scalar_proto = expect_expr_variant!(
node,
protobuf::physical_expr_node::ExprType::Literal,
"Literal",
);
let value = ScalarValue::try_from(scalar_proto)?;
Ok(Arc::new(Literal::new(value)))
use protobuf::physical_expr_node::ExprType;

let protobuf::PhysicalExprNode {
expr_type,
// Expression IDs are handled by the enclosing proto converter.
expr_id: _,
} = node;
let (value, metadata) = match expr_type {
Some(ExprType::Literal(scalar)) => {
let datafusion_proto_models::datafusion_common::ScalarValue {
// The scalar payload is decoded by ScalarValue::try_from.
value: _,
} = scalar;
(ScalarValue::try_from(scalar)?, None)
}
Some(ExprType::LiteralWithMetadata(protobuf::PhysicalLiteralNode {
value,
metadata,
})) => {
let value = value.as_ref().ok_or_else(|| {
internal_datafusion_err!("Literal is missing required field 'value'")
})?;
(
ScalarValue::try_from(value)?,
Some(FieldMetadata::from(metadata)),
)
}
_ => return internal_err!("PhysicalExprNode is not a Literal"),
};
Ok(Arc::new(Literal::new_with_metadata(value, metadata)))
}
}

Expand Down Expand Up @@ -315,6 +347,26 @@ mod proto_tests {
assert_eq!(lit.value(), &ScalarValue::Int32(Some(42)));
}

#[test]
fn try_from_proto_rejects_missing_value() {
let node = datafusion_proto_models::protobuf::PhysicalExprNode {
expr_id: None,
expr_type: Some(physical_expr_node::ExprType::LiteralWithMetadata(
datafusion_proto_models::protobuf::PhysicalLiteralNode {
value: None,
metadata: Default::default(),
},
)),
};
let schema = Schema::empty();
let decoder = UnreachableDecoder;
let ctx = PhysicalExprDecodeCtx::new(&schema, &decoder);
let err = Literal::try_from_proto(&node, &ctx).unwrap_err();
assert!(
matches!(err, DataFusionError::Internal(msg) if msg.contains("Literal is missing required field 'value'"))
);
}

#[test]
fn try_from_proto_rejects_non_literal_node() {
let node = column_node("a");
Expand Down
6 changes: 4 additions & 2 deletions datafusion/physical-expr/src/expressions/negative.rs
Original file line number Diff line number Diff line change
Expand Up @@ -184,11 +184,12 @@ impl PhysicalExpr for NegativeExpr {
) -> Result<Option<datafusion_proto_models::protobuf::PhysicalExprNode>> {
use datafusion_proto_models::protobuf;

let Self { arg } = self;
Ok(Some(protobuf::PhysicalExprNode {
expr_id: None,
expr_type: Some(protobuf::physical_expr_node::ExprType::Negative(Box::new(
protobuf::PhysicalNegativeNode {
expr: Some(Box::new(ctx.encode_child(&self.arg)?)),
expr: Some(Box::new(ctx.encode_child(arg)?)),
},
))),
}))
Expand All @@ -210,8 +211,9 @@ impl NegativeExpr {
protobuf::physical_expr_node::ExprType::Negative,
"Negative",
);
let protobuf::PhysicalNegativeNode { expr } = n.as_ref();
let expr =
ctx.decode_required_expression(n.expr.as_deref(), "NegativeExpr", "expr")?;
ctx.decode_required_expression(expr.as_deref(), "NegativeExpr", "expr")?;

Ok(Arc::new(NegativeExpr::new(expr)))
}
Expand Down
7 changes: 4 additions & 3 deletions datafusion/physical-expr/src/expressions/not.rs
Original file line number Diff line number Diff line change
Expand Up @@ -189,11 +189,12 @@ impl PhysicalExpr for NotExpr {
) -> Result<Option<datafusion_proto_models::protobuf::PhysicalExprNode>> {
use datafusion_proto_models::protobuf;

let Self { arg } = self;
Ok(Some(protobuf::PhysicalExprNode {
expr_id: None,
expr_type: Some(protobuf::physical_expr_node::ExprType::NotExpr(Box::new(
protobuf::PhysicalNot {
expr: Some(Box::new(ctx.encode_child(&self.arg)?)),
expr: Some(Box::new(ctx.encode_child(arg)?)),
},
))),
}))
Expand All @@ -215,8 +216,8 @@ impl NotExpr {
protobuf::physical_expr_node::ExprType::NotExpr,
"NotExpr",
);
let expr =
ctx.decode_required_expression(not_expr.expr.as_deref(), "NotExpr", "expr")?;
let protobuf::PhysicalNot { expr } = not_expr.as_ref();
let expr = ctx.decode_required_expression(expr.as_deref(), "NotExpr", "expr")?;

Ok(Arc::new(NotExpr::new(expr)))
}
Expand Down
9 changes: 4 additions & 5 deletions datafusion/physical-expr/src/expressions/unknown_column.rs
Original file line number Diff line number Diff line change
Expand Up @@ -93,12 +93,11 @@ impl PhysicalExpr for UnKnownColumn {
) -> Result<Option<datafusion_proto_models::protobuf::PhysicalExprNode>> {
use datafusion_proto_models::protobuf;

let Self { name } = self;
Ok(Some(protobuf::PhysicalExprNode {
expr_id: None,
expr_type: Some(protobuf::physical_expr_node::ExprType::UnknownColumn(
protobuf::UnknownColumn {
name: self.name.clone(),
},
protobuf::UnknownColumn { name: name.clone() },
)),
}))
}
Expand All @@ -114,12 +113,12 @@ impl UnKnownColumn {
use datafusion_physical_expr_common::expect_expr_variant;
use datafusion_proto_models::protobuf;

let unknown_col = expect_expr_variant!(
let protobuf::UnknownColumn { name } = expect_expr_variant!(
node,
protobuf::physical_expr_node::ExprType::UnknownColumn,
"UnKnownColumn",
);
Ok(Arc::new(UnKnownColumn::new(&unknown_col.name)))
Ok(Arc::new(UnKnownColumn::new(name)))
}
}

Expand Down
6 changes: 6 additions & 0 deletions datafusion/proto-models/proto/datafusion.proto
Original file line number Diff line number Diff line change
Expand Up @@ -1074,9 +1074,15 @@ message PhysicalExprNode {
PhysicalLambdaVariableExprNode lambda_variable = 26;
PhysicalRangeExprNode range_expr = 27;
PhysicalSqlSimilarToPatternNode sql_similar_to_pattern = 28;
PhysicalLiteralNode literal_with_metadata = 29;
}
}

message PhysicalLiteralNode {
datafusion_common.ScalarValue value = 1;
map<string, string> metadata = 2;
}

message PhysicalDynamicFilterNode {
repeated PhysicalExprNode children = 1;
repeated PhysicalExprNode remapped_children = 2;
Expand Down
Loading