Skip to content

isthmus: visit(VirtualTableScan) decides between LogicalValues and VirtualTable over the columns a projection drops #1345

Description

@nielspardon

With #1280, visit(VirtualTableScan) applies a read's projection as a Project above the node it builds, but it still converts every column of every row and decides encodableAsTuples over all of them:

final RelDataType rowType =
typeFactory.createTypeWithNullability(
typeConverter.toCalcite(typeFactory, schema.struct(), schema.names()), false);
List<List<RexNode>> convertedRows = new ArrayList<>();
// The same values with the cast a nullable literal converts as taken off, which is the form a
// LogicalValues tuple would hold them in -- where they are literals at all.
List<List<RexNode>> tupleValues = new ArrayList<>();
for (final Expression.NestedStruct rowExpr : virtualTableScan.getRows()) {
List<RexNode> convertedRow = new ArrayList<>();
List<RexNode> tupleRow = new ArrayList<>();
for (int column = 0; column < rowExpr.fields().size(); column++) {
Expression field = rowExpr.fields().get(column);
RelDataType declaredType = rowType.getFieldList().get(column).getType();
RexNode value =
valueAsDeclared(field.accept(expressionRexConverter, context), declaredType);
convertedRow.add(value);
// Only a converted literal has its nullability cast taken off. A cast the plan carries
// itself is doing work -- it has a failure behavior, and it is part of what round-trips --
// so a row holding one is computed rather than tabulated.
tupleRow.add(field instanceof Expression.Literal ? unwrapNullabilityCast(value) : value);
}
convertedRows.add(convertedRow);
tupleValues.add(tupleRow);
}
// A LogicalValues tuple holds nothing but literals, and whether a value is one is a property of
// what it converts to rather than of the Substrait expression it came from: a struct converts
// to a ROW call and a list to an array constructor, literal or not. Neither belongs in a tuple
// anyway -- Calcite orders tuples by casting each value to Comparable, and the value of a row
// literal is a list of RexLiterals, which are not.
boolean encodableAsTuples =
tupleValues.stream().flatMap(List::stream).allMatch(value -> value instanceof RexLiteral);
if (encodableAsTuples) {
ImmutableList.Builder<ImmutableList<RexLiteral>> tuplesBuilder = ImmutableList.builder();
for (final List<RexNode> tupleRow : tupleValues) {
ImmutableList.Builder<RexLiteral> tupleBuilder = ImmutableList.builder();
for (RexNode value : tupleRow) {
tupleBuilder.add((RexLiteral) value);
}
tuplesBuilder.add(tupleBuilder.build());
}
return applyRelCommon(
applyProjection(
LogicalValues.create(relBuilder.getCluster(), rowType, tuplesBuilder.build()),
virtualTableScan.getProjection()),
virtualTableScan);
} else {
// A row that does not fit a LogicalValues tuple keeps its expressions, in a relation of our
// own: Calcite has none that holds them, and expanding the table into a projection per row
// does not come back -- the projection is what converts back, and the table is gone. A
// consumer whose planner only knows Calcite's own relations can expand it with
// VirtualTableExpansionRule.
return applyRelCommon(
applyProjection(
VirtualTable.create(relBuilder.getCluster(), rowType, convertedRows),
virtualTableScan.getProjection()),
virtualTableScan);
}

One non-literal value in a column the mask drops therefore sends the table down the VirtualTable branch, even when every column the read produces is a literal:

VirtualTableScan(col1 i32, col2 fp64, col3 string), row [2, add(4.4, 4.5), 'a'], projection = [0, 2]
  converted       LogicalProject(col1=[$0], col3=[$2])
                    VirtualTable(rows=[[{ 2, +(4.4E0, 4.5E0), 'a' }]])
  columns (0, 2)  LogicalValues(tuples=[[{ 2, 'a' }]])      -- the same table declared without a mask

A consumer whose planner knows only Calcite's own relations then has to register VirtualTableExpansionRule for a table whose output is all literals. The masked-away expression is also still converted and carried: converting back gives a Project over a VirtualTableScan whose row still holds add(4.4, 4.5).

Deciding encodableAsTuples over the columns the mask keeps, and converting only those, would keep such a table a plain LogicalValues. That changes what a round trip returns: building over the masked row type means neither the masked-away columns nor the mask itself come back. #1218 covers the same method's per-row allocation, not this.

Measured at #1280's head, 146f54f.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions