Skip to content
Merged
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
56 changes: 53 additions & 3 deletions pkg/partitionprune/filter.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ package partitionprune
import (
"context"
"sort"
"strings"

"github.com/matrixorigin/matrixone/pkg/common/moerr"
"github.com/matrixorigin/matrixone/pkg/container/batch"
Expand Down Expand Up @@ -147,7 +148,7 @@ func hashFilterExpr(
if !ok {
return nil, false, nil
}
if left.Col.ColPos != colPosition {
if !matchesPartitionColumn(left.Col, colPosition, metadata.Partitions[0].Expr) {
return nil, false, nil
}
value, canPrune := normalizePartitionValue(exprImpl.F.Args[1])
Expand Down Expand Up @@ -262,7 +263,7 @@ func rangeFilterExpr(
if !ok {
return nil, false, nil
}
if left.Col.ColPos != colPosition {
if !matchesPartitionColumn(left.Col, colPosition, metadata.Partitions[0].Expr) {
return nil, false, nil
}
value, canPrune := normalizePartitionValue(exprImpl.F.Args[1])
Expand Down Expand Up @@ -460,6 +461,55 @@ func mustGetColPosition(expr *plan.Expr) int32 {
return -1
}

// Scan column positions are compacted independently of the stored partition
// expression. Use column identity when the scan supplies a name; retain the
// positional fallback for older nameless expressions.
func matchesPartitionColumn(col *plan.ColRef, position int32, partitionExpr *plan.Expr) bool {
if col.Name != "" {
partitionName, consistent := partitionExpressionColumnName(partitionExpr)
scanName, unambiguous := scanColumnName(col.Name)
return consistent && unambiguous && partitionName != "" &&
strings.EqualFold(scanName, partitionName)
}
return col.ColPos == position
}

// A predicate on one column cannot prune a partition expression that depends
// on a different column or on several columns.
func partitionExpressionColumnName(expr *plan.Expr) (string, bool) {
if expr == nil {
return "", true
}
if col := expr.GetCol(); col != nil {
return col.Name, col.Name != ""
}
if fn := expr.GetF(); fn != nil {
var name string
for _, arg := range fn.Args {
other, ok := partitionExpressionColumnName(arg)
if !ok || (name != "" && other != "" && !strings.EqualFold(name, other)) {
return "", false
}
if other != "" {
name = other
}
}
return name, true
}
return "", true
}

// The planner emits alias.column. More than one dot is ambiguous because the
// quoted alias or the physical column name may itself contain a dot. In that
// case partition pruning must leave the filter to the row reader.
func scanColumnName(name string) (string, bool) {
if idx := strings.IndexByte(name, '.'); idx >= 0 {
column := name[idx+1:]
return column, !strings.ContainsRune(column, '.')
}
return name, true
}

// listFilter handles partition pruning for list-based partitioning.
// It evaluates the filters against list partition expressions and returns matching partition positions.
func listFilter(
Expand Down Expand Up @@ -622,7 +672,7 @@ func listFilterExprNormalized(
if !ok {
return nil, false, nil
}
if left.Col.ColPos != colPosition {
if !matchesPartitionColumn(left.Col, colPosition, metadata.Partitions[0].Expr) {
return nil, false, nil
}
left.Col.ColPos = 0
Expand Down
41 changes: 41 additions & 0 deletions pkg/partitionprune/filter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,30 @@ func TestFilter(t *testing.T) {
want []int
wantErr bool
}{
{
name: "different column sharing scan position cannot prune partitions",
filters: []*plan.Expr{makeNamedEqualExpr(0, "hidden_key", 1)},
metadata: namedRangePartitionMetadata("a"),
want: []int{0, 1, 2},
},
{
name: "same column at compact scan position prunes partitions",
filters: []*plan.Expr{makeNamedEqualExpr(1, "range_probe.a", 1)},
metadata: namedRangePartitionMetadata("a"),
want: []int{1},
},
{
name: "dotted column name must not alias partition column",
filters: []*plan.Expr{makeNamedEqualExpr(0, "t.a.b", 1)},
metadata: namedRangePartitionMetadata("b"),
want: []int{0, 1, 2},
},
{
name: "dotted alias cannot be distinguished from dotted partition column",
filters: []*plan.Expr{makeNamedEqualExpr(1, "x.a.b", 1)},
metadata: namedRangePartitionMetadata("a.b"),
want: []int{0, 1, 2},
},
{
name: "empty filters",
filters: []*plan.Expr{},
Expand Down Expand Up @@ -478,6 +502,23 @@ func makeEqualExpr(colPos int32, value int64) *plan.Expr {
}
}

func makeNamedEqualExpr(colPos int32, name string, value int64) *plan.Expr {
expr := makeEqualExpr(colPos, value)
expr.GetF().Args[0].GetCol().Name = name
return expr
}

func namedRangePartitionMetadata(name string) partition.PartitionMetadata {
return partition.PartitionMetadata{
Method: partition.PartitionMethod_Range,
Partitions: []partition.Partition{
{Position: 0, Expr: newTestRangeExpr(name, 0)},
{Position: 1, Expr: newTestRangeExpr(name, 1)},
{Position: 2, Expr: newTestRangeExpr(name, 2)},
},
}
}

func makeEqualExprInt32(colPos int32, value int32) *plan.Expr {
return &plan.Expr{
Typ: plan.Type{Id: int32(types.T_bool)},
Expand Down
181 changes: 152 additions & 29 deletions pkg/sql/plan/expr_opt.go
Original file line number Diff line number Diff line change
Expand Up @@ -169,18 +169,21 @@ func (builder *QueryBuilder) appendCompositePartBlockFilters(filters map[int32][
for nodeID, candidates := range filters {
node := builder.qry.Nodes[nodeID]
for _, candidate := range candidates {
duplicate := false
for _, existing := range node.BlockFilterList {
if blockFilterEquivalent(existing, candidate.copy) {
duplicate = true
break
}
}
if !duplicate {
node.BlockFilterList = append(node.BlockFilterList, candidate.copy)
}
appendUniqueBlockFilter(node, candidate.copy, false)
}
}
}

func appendUniqueBlockFilter(node *plan.Node, filter *plan.Expr, copyFilter bool) {
for _, existing := range node.BlockFilterList {
if blockFilterEquivalent(existing, filter) {
return
}
}
if copyFilter {
filter = DeepCopyExpr(filter)
}
node.BlockFilterList = append(node.BlockFilterList, filter)
}

// retainConsumedCompositePartBlockFilters keeps only predicates removed by the
Expand Down Expand Up @@ -282,26 +285,145 @@ func (builder *QueryBuilder) appendCompoundKeyBlockFilters(nodeID int32) {
return
}
allowed := map[int32]struct{}{compoundPos: {}}
for _, filter := range node.FilterList {
if !ExprIsZonemappable(builder.GetContext(), filter) ||
!exprOnlyReferencesColumns(filter, node.BindingTags[0], allowed) {
continue
var leadingPos int32 = -1
if node.TableDef.ClusterBy != nil && util.JudgeIsCompositeClusterByColumn(node.TableDef.ClusterBy.Name) {
parts := util.SplitCompositeClusterByColumnName(node.TableDef.ClusterBy.Name)
if len(parts) > 0 {
if pos, found := node.TableDef.Name2ColIndex[parts[0]]; found {
leadingPos = pos
}
}
duplicate := false
for _, existing := range node.BlockFilterList {
if blockFilterEquivalent(existing, filter) {
duplicate = true
break
} else if node.TableDef.Pkey != nil && len(node.TableDef.Pkey.Names) > 1 {
if pos, found := node.TableDef.Name2ColIndex[node.TableDef.Pkey.Names[0]]; found {
leadingPos = pos
}
}
for _, filter := range node.FilterList {
zonemappable := ExprIsZonemappable(builder.GetContext(), filter)
if leadingPos >= 0 && zonemappable {
if prefix := builder.leadingCompositeRangeBlockFilter(filter, node.TableDef, node.BindingTags[0], leadingPos, compoundPos); prefix != nil {
appendUniqueBlockFilter(node, prefix, false)
}
}
if !duplicate {
node.BlockFilterList = append(node.BlockFilterList, DeepCopyExpr(filter))
if !zonemappable || !exprOnlyReferencesColumns(filter, node.BindingTags[0], allowed) {
continue
}
appendUniqueBlockFilter(node, filter, true)
}
}
visit(nodeID)
}

// leadingCompositeRangeBlockFilter adds an object-pruning predicate without
// replacing the SQL row predicate. In particular, NULL and unsupported bound
// types must continue through the ordinary row comparison.
func (builder *QueryBuilder) leadingCompositeRangeBlockFilter(filter *plan.Expr, tableDef *plan.TableDef, tag, leadingPos, compoundPos int32) *plan.Expr {
fn := filter.GetF()
if fn == nil || fn.Func == nil || len(fn.Args) < 2 || fn.Args[0] == nil {
return nil
}
op := fn.Func.ObjName
col := fn.Args[0]
if op == "<" || op == "<=" || op == ">" || op == ">=" {
op = canonicalRangeOp(fn)
if fn.Args[0].GetCol() == nil && len(fn.Args) == 2 {
col = fn.Args[1]
}
}
if col.GetCol() == nil || col.GetCol().RelPos != tag || col.GetCol().ColPos != leadingPos ||
!compositeRangeOrderPreserving(types.T(col.Typ.Id)) {
return nil
}
boundCompatible := func(bound *plan.Expr) bool {
if bound == nil {
return false
}
if lit := stripConstLiteralCasts(bound).GetLit(); lit != nil && lit.Isnull {
return false
}
return isRuntimeConstExpr(bound) && bound.Typ.Id == col.Typ.Id &&
bound.Typ.Scale == col.Typ.Scale
}
key := &plan.Expr{Typ: tableDef.Cols[compoundPos].Typ, Expr: &plan.Expr_Col{Col: &plan.ColRef{
RelPos: tag, ColPos: compoundPos, Name: tableDef.Cols[compoundPos].Name,
}}}
serial := func(bound *plan.Expr) *plan.Expr {
ret, ok := builder.bindCompositeKeySerial([]*plan.Expr{bound})
if !ok {
return nil
}
return ret
}
var name string
var args []*plan.Expr
switch op {
case "between", "in_range":
if (op == "between" && len(fn.Args) != 3) ||
(op == "in_range" && len(fn.Args) != 4) {
return nil
}
if !boundCompatible(fn.Args[1]) || !boundCompatible(fn.Args[2]) {
return nil
}
lower, upper := serial(fn.Args[1]), serial(fn.Args[2])
if lower == nil || upper == nil {
return nil
}
name = "prefix_between"
args = []*plan.Expr{key, lower, upper}
if op == "in_range" {
if fn.Args[3] == nil || !isRuntimeConstExpr(fn.Args[3]) {
return nil
}
name = "prefix_in_range"
args = append(args, fn.Args[3])
}
case "<", "<=", ">", ">=":
boundExpr := rangeFilterConstValue(fn)
if len(fn.Args) != 2 || !boundCompatible(boundExpr) {
return nil
}
bound := serial(boundExpr)
if bound == nil {
return nil
}
empty := makePlan2StringConstExprWithType("")
empty.Typ.Id = int32(types.T_varchar)
var flag byte
if op == "<" || op == "<=" {
args = []*plan.Expr{key, empty, bound}
if op == "<" {
flag = 2
}
} else {
args = []*plan.Expr{key, bound, empty}
if op == ">" {
flag = 1
}
}
name = "prefix_in_range"
args = append(args, makePlan2Uint8ConstExprWithType(flag))
default:
return nil
}
ret, ok := builder.bindCompositeKeyPredicate(name, args...)
if !ok {
return nil
}
return ret
}

func compositeRangeOrderPreserving(oid types.T) bool {
switch oid {
case types.T_int8, types.T_int16, types.T_int32, types.T_int64,
types.T_uint8, types.T_uint16, types.T_uint32, types.T_uint64,
types.T_date, types.T_time, types.T_datetime, types.T_timestamp,
types.T_decimal64, types.T_decimal128, types.T_decimal256:
return true
}
return false
}

func existingCompositeBlockFilters(node *plan.Node) []*plan.Expr {
if node.TableDef == nil || len(node.BindingTags) == 0 {
return nil
Expand Down Expand Up @@ -2239,6 +2361,11 @@ func (builder *QueryBuilder) doMergeFiltersOnCompositeKey(tableDef *plan.TableDe
if _, ok := sortKeyPartCols[col.ColPos]; !ok {
continue
}
// Keep first-component SQL bounds intact. Their optional compound-key
// object filter is added separately without casting or consuming them.
if col.ColPos == tableDef.Name2ColIndex[Parts[0]] {
continue
}
if isLower {
colLowerBounds[col.ColPos] = i
} else {
Expand Down Expand Up @@ -2360,6 +2487,10 @@ func (builder *QueryBuilder) doMergeFiltersOnCompositeKey(tableDef *plan.TableDe
return filters
}
lastFuncName := lastFn.Func.ObjName
if len(filterIdx) == 1 && (lastFuncName == "between" || lastFuncName == "in_range" ||
lastFuncName == "<" || lastFuncName == "<=" || lastFuncName == ">" || lastFuncName == ">=") {
return filters
}
if lastFuncName == "in" {
if !hasNonNilFunctionArgs(lastFn, 2) {
return filters
Expand Down Expand Up @@ -2448,10 +2579,6 @@ func (builder *QueryBuilder) doMergeFiltersOnCompositeKey(tableDef *plan.TableDe
serialArgs[i] = filters[filterIdx[i]].GetF().Args[1]
}

if len(filterIdx) < numParts && len(serialArgs) == 0 {
return filters
}

tmpSerialArgs := DeepCopyExprList(serialArgs)
tmpSerialArgs = append(tmpSerialArgs, lastFn.Args[1])
leftArg, ok := builder.bindCompositeKeySerial(tmpSerialArgs)
Expand Down Expand Up @@ -2488,10 +2615,6 @@ func (builder *QueryBuilder) doMergeFiltersOnCompositeKey(tableDef *plan.TableDe
serialArgs[i] = filters[filterIdx[i]].GetF().Args[1]
}

if len(filterIdx) < numParts && len(serialArgs) == 0 {
return filters
}

tmpSerialArgs := append(DeepCopyExprList(serialArgs), lastFn.Args[1])
boundArg, ok := builder.bindCompositeKeySerial(tmpSerialArgs)
if !ok {
Expand Down
Loading
Loading