UpdateMvByPartitionCommand.java

// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements.  See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership.  The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License.  You may obtain a copy of the License at
//
//   http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied.  See the License for the
// specific language governing permissions and limitations
// under the License.

package org.apache.doris.nereids.trees.plans.commands;

import org.apache.doris.catalog.Column;
import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.ListPartitionItem;
import org.apache.doris.catalog.MTMV;
import org.apache.doris.catalog.OlapTable;
import org.apache.doris.catalog.Partition;
import org.apache.doris.catalog.PartitionItem;
import org.apache.doris.catalog.PartitionKey;
import org.apache.doris.catalog.TableIf;
import org.apache.doris.catalog.Type;
import org.apache.doris.common.AnalysisException;
import org.apache.doris.common.UserException;
import org.apache.doris.datasource.ExternalTable;
import org.apache.doris.datasource.mvcc.MvccUtil;
import org.apache.doris.mtmv.BaseColInfo;
import org.apache.doris.mtmv.BaseTableInfo;
import org.apache.doris.mtmv.MTMVRelatedTableIf;
import org.apache.doris.mtmv.MTMVUtil;
import org.apache.doris.nereids.StatementContext;
import org.apache.doris.nereids.analyzer.UnboundRelation;
import org.apache.doris.nereids.analyzer.UnboundSlot;
import org.apache.doris.nereids.analyzer.UnboundTableSinkCreator;
import org.apache.doris.nereids.parser.NereidsParser;
import org.apache.doris.nereids.trees.expressions.EqualTo;
import org.apache.doris.nereids.trees.expressions.Expression;
import org.apache.doris.nereids.trees.expressions.GreaterThanEqual;
import org.apache.doris.nereids.trees.expressions.InPredicate;
import org.apache.doris.nereids.trees.expressions.IsNull;
import org.apache.doris.nereids.trees.expressions.LessThan;
import org.apache.doris.nereids.trees.expressions.Slot;
import org.apache.doris.nereids.trees.expressions.literal.BooleanLiteral;
import org.apache.doris.nereids.trees.expressions.literal.Literal;
import org.apache.doris.nereids.trees.expressions.literal.NullLiteral;
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.algebra.Sink;
import org.apache.doris.nereids.trees.plans.commands.insert.InsertOverwriteTableCommand;
import org.apache.doris.nereids.trees.plans.logical.LogicalCTE;
import org.apache.doris.nereids.trees.plans.logical.LogicalCatalogRelation;
import org.apache.doris.nereids.trees.plans.logical.LogicalFilter;
import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
import org.apache.doris.nereids.trees.plans.logical.LogicalSink;
import org.apache.doris.nereids.trees.plans.logical.LogicalSubQueryAlias;
import org.apache.doris.nereids.trees.plans.visitor.DefaultPlanRewriter;
import org.apache.doris.nereids.util.ExpressionUtils;
import org.apache.doris.nereids.util.RelationUtil;
import org.apache.doris.qe.ConnectContext;

import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.Lists;
import com.google.common.collect.Range;
import com.google.common.collect.Sets;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;

import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;

/**
 * Update mv by partition
 */
public class UpdateMvByPartitionCommand extends InsertOverwriteTableCommand {
    private static final Logger LOG = LogManager.getLogger(UpdateMvByPartitionCommand.class);

    private UpdateMvByPartitionCommand(LogicalPlan logicalQuery) {
        super(logicalQuery, Optional.empty(), Optional.empty(), Optional.empty());
    }

    @Override
    public boolean isForceDropPartition() {
        // After refreshing the data in MTMV, it will be synchronized with the base table
        // and there is no need to put it in the recycle bin
        return true;
    }

    /**
     * Construct command
     *
     * @param mv materialize view
     * @param partitionNames update partitions in mv and tables
     * @param tableWithPartKey the partitions key for different table
     * @param statementContext statementContext
     * @param readableBasePartitions the base partitions each base table may be read from, by base table,
     *                               named with the partition names of the base table. A base table named
     *                               with an empty set is not read at all; a base table absent from the map,
     *                               or a null map, keeps the MV partition's own key range. The tables named
     *                               are olap ones, and the partitions are looked up on them
     * @return command
     */
    public static UpdateMvByPartitionCommand from(MTMV mv, Set<String> partitionNames,
            Map<TableIf, String> tableWithPartKey, StatementContext statementContext,
            Map<BaseTableInfo, Set<String>> readableBasePartitions) throws UserException {
        NereidsParser parser = new NereidsParser();
        Map<TableIf, Set<Expression>> predicates =
                constructTableWithPredicates(mv, partitionNames, tableWithPartKey, readableBasePartitions);
        List<String> parts = constructPartsForMv(partitionNames);
        Plan plan = parser.parseSingle(mv.getQuerySql());
        if (plan instanceof Sink) {
            plan = plan.child(0);
        }
        List<String> sinkColumns = mv.isIvm() ? mv.getInsertedColumnNames() : ImmutableList.of();
        LogicalSink<? extends Plan> sink = UnboundTableSinkCreator.createUnboundTableSink(mv.getFullQualifiers(),
                sinkColumns, ImmutableList.of(), parts, plan);
        if (LOG.isDebugEnabled()) {
            LOG.debug("MTMVTask plan for mvName: {}, partitionNames: {}, plan: {}", mv.getName(), partitionNames,
                    sink.treeString());
        }
        statementContext.setMvRefreshPredicates(predicates);
        return new UpdateMvByPartitionCommand(sink);
    }

    private static List<String> constructPartsForMv(Set<String> partitionNames) {
        return Lists.newArrayList(partitionNames);
    }

    /**
     * The predicate every base table of the MV definition is read through.
     *
     * <p>A table the caller scopes is read from exactly the base partitions it named. Those are the ones
     * the refresh is about to record as this MV partition's, and the read is what has to match the record:
     * reading the MV partition's own key range instead also reads base partitions no snapshot describes,
     * and a later silent change to one of them -- dropped, with the base partition set back to what it
     * was -- leaves the rows it put in this MV partition behind while the partition is still judged
     * synchronized, so the transparent rewrite serves them and no refresh plans it again.
     *
     * <p>Every other table keeps the MV partition's own key range, which is what the tables the caller
     * does not scope were always read through. Scoped tables are olap ones; the partition names are
     * looked up on one, see the caller.
     */
    private static Map<TableIf, Set<Expression>> constructTableWithPredicates(MTMV mv,
            Set<String> partitionNames, Map<TableIf, String> tableWithPartKey,
            Map<BaseTableInfo, Set<String>> readableBasePartitions) throws AnalysisException {
        Set<PartitionItem> mvItems = Sets.newHashSet();
        for (String partitionName : partitionNames) {
            mvItems.add(mv.getPartitionItemOrAnalysisException(partitionName));
        }
        ImmutableMap.Builder<TableIf, Set<Expression>> builder = new ImmutableMap.Builder<>();
        for (Map.Entry<TableIf, String> entry : tableWithPartKey.entrySet()) {
            TableIf table = entry.getKey();
            String colName = entry.getValue();
            Set<String> readable = readableBasePartitions == null ? null
                    : readableBasePartitions.get(new BaseTableInfo(table));
            if (readable == null) {
                builder.put(table, constructPredicates(mvItems, colName));
                continue;
            }
            OlapTable olapTable = (OlapTable) table;
            Set<PartitionItem> items = Sets.newHashSet();
            for (String partitionName : readable) {
                items.add(olapTable.getPartitionItemOrAnalysisException(partitionName));
            }
            if (items.stream().anyMatch(PartitionItem::isDefaultPartition)) {
                // One of the partitions this MV partition is recorded with is a list partitioned table's
                // default partition, which takes the rows no other partition of it claims. Those rows are
                // the ones the MV partition's own key range names, wherever the base table put them, and a
                // partition of the MV takes them by that key rather than by the partition they were placed
                // in. So a table whose mapped partitions include one is read the way an unscoped one is:
                // the MV partition's key range, at the partition column's own type. That read can be seen to
                // be too wide -- it is the one this scope exists to narrow -- rather than one that drops
                // rows belonging to the MV partition being refreshed. The mapping names the default
                // partition in every MV partition that reads the table, so this is reached for each of them
                // and not only for the one the sentinel key maps to.
                builder.put(table, constructPredicates(mvItems, colName,
                        Optional.of(partitionColumnType(olapTable, colName))));
                continue;
            }
            if (readable.isEmpty()) {
                // No partition of this table feeds the MV partitions being refreshed, which is "no row"
                // rather than "every row": constructPredicates answers the other way for an empty set,
                // and that answer would put every row of the table into each of them.
                builder.put(table, Sets.newHashSet(BooleanLiteral.FALSE));
                continue;
            }
            builder.put(table, constructPredicatesOfBasePartitions(items, olapTable, colName));
        }
        return builder.build();
    }

    /**
     * construct predicates for partition items, the min key is the min key of range items.
     * For list partition or less than partition items, the min key is null.
     */
    @VisibleForTesting
    public static Set<Expression> constructPredicates(Set<PartitionItem> partitions, String colName) {
        return constructPredicates(partitions, colName, Optional.empty());
    }

    /** The same, for the callers that know the partition column the predicate is written against. */
    private static Set<Expression> constructPredicates(Set<PartitionItem> partitions, String colName,
            Optional<Type> columnType) {
        return constructPredicates(partitions, new UnboundSlot(colName), columnType);
    }

    /**
     * construct predicates for partition items, the min key is the min key of range items.
     * For list partition or less than partition items, the min key is null.
     */
    @VisibleForTesting
    public static Set<Expression> constructPredicates(Set<PartitionItem> partitions, Slot colSlot) {
        return constructPredicates(partitions, colSlot, Optional.empty());
    }

    private static Set<Expression> constructPredicates(Set<PartitionItem> partitions, Slot colSlot,
            Optional<Type> columnType) {
        Set<Expression> predicates = new HashSet<>();
        if (partitions.isEmpty()) {
            return Sets.newHashSet(BooleanLiteral.TRUE);
        }
        if (partitions.iterator().next() instanceof ListPartitionItem) {
            for (PartitionItem item : partitions) {
                predicates.add(convertListPartitionToIn(item, colSlot, columnType));
            }
        } else {
            for (PartitionItem item : partitions) {
                predicates.add(convertRangePartitionToCompare(item, colSlot, columnType));
            }
        }
        return predicates;
    }

    /**
     * The predicate a base table is read through when the refresh is to read exactly these partitions of it.
     *
     * <p>A partition of a list partitioned table holds one key per partition column, and the column the MV
     * partition is named by is only one of them. A predicate on that column alone also reaches the
     * partitions whose other keys differ -- a table partitioned by (d, region) has one partition of
     * (d0, 'US') and one of (d0, 'EU'), and `d = d0` reaches both, while only the second is a partition
     * this refresh is to read; a later drop of the first would then leave its rows in the MV partition
     * while the snapshot, which names only the second, still calls it synchronized. So a list partition is
     * pinned to its whole key. A range partition is pinned to its bounds, which is the same thing: a base
     * table partitioned by range has a single partition column, see
     * {@code RangePartitionItem#toPartitionKeyDesc(int)}.
     *
     * <p>The partitions are never empty: a table the caller scopes with no partition is read as nothing
     * before this is reached, see {@code constructTableWithPredicates}.
     */
    private static Set<Expression> constructPredicatesOfBasePartitions(Set<PartitionItem> partitions,
            OlapTable baseTable, String colName) {
        List<Column> partitionColumns = baseTable.getPartitionColumns();
        List<Type> partitionColumnTypes = Lists.transform(partitionColumns, Column::getType);
        if (!(partitions.iterator().next() instanceof ListPartitionItem)) {
            Set<Expression> predicates = new HashSet<>();
            for (PartitionItem item : partitions) {
                predicates.add(convertRangePartitionToCompare(item, new UnboundSlot(colName),
                        Optional.of(partitionColumnTypes.get(0))));
            }
            return predicates;
        }
        List<Slot> partitionSlots = Lists.newArrayList();
        for (Column partitionColumn : partitionColumns) {
            partitionSlots.add(new UnboundSlot(partitionColumn.getName()));
        }
        Set<Expression> predicates = new HashSet<>();
        for (PartitionItem item : partitions) {
            predicates.add(convertListPartitionToKey(item, partitionSlots, partitionColumnTypes));
        }
        return predicates;
    }

    /** The type of the partition column of this table the MV partition is named by. */
    private static Type partitionColumnType(OlapTable table, String colName) {
        for (Column partitionColumn : table.getPartitionColumns()) {
            if (partitionColumn.getName().equalsIgnoreCase(colName)) {
                return partitionColumn.getType();
            }
        }
        throw new IllegalStateException(
                "no partition column " + colName + " on table " + table.getName());
    }

    /**
     * One partition of a list partitioned table, pinned to the whole of each key it holds: the keys are
     * what tells it apart from a partition that shares a key with it, and the value of a key a row does not
     * have is asked for as {@code IS NULL}, since no comparison to it is ever true.
     *
     * <p>A list partitioned table's default partition is not one of these: its key is the sentinel the rows
     * no other partition claims are placed by rather than a value, and a partition of the MV takes the rows
     * whose own key falls in it wherever the base table put them. What it holds cannot be said with a
     * predicate on the partition columns, so a table that has one is read the way an unscoped one is, see
     * {@code constructTableWithPredicates}.
     */
    private static Expression convertListPartitionToKey(PartitionItem item, List<Slot> partitionSlots,
            List<Type> partitionColumnTypes) {
        List<Expression> keys = new ArrayList<>();
        for (PartitionKey key : ((ListPartitionItem) item).getItems()) {
            List<Expression> oneKey = new ArrayList<>();
            for (int pos = 0; pos < partitionSlots.size(); pos++) {
                Expression value = convertPartitionKeyToLiteral(key, pos,
                        Optional.of(partitionColumnTypes.get(pos)));
                oneKey.add(value instanceof NullLiteral ? new IsNull(partitionSlots.get(pos))
                        : new EqualTo(partitionSlots.get(pos), value));
            }
            keys.add(ExpressionUtils.and(oneKey));
        }
        Preconditions.checkState(!keys.isEmpty(), "a list partition holds at least one key: %s", item);
        return ExpressionUtils.or(keys);
    }

    /**
     * A partition key is a value of the partition column it is written against, and what tells two of them
     * apart can be a scale the key's primitive type does not carry: a literal rounded to a coarser one is
     * one no row of the partition compares equal to, and the rows of the partition are then read as none.
     * The callers that have the column pass its type; the ones that do not leave the key's own primitive
     * type, which is what a reader of these predicates was given before.
     */
    private static Expression convertPartitionKeyToLiteral(PartitionKey key, int keyPos,
            Optional<Type> columnType) {
        return Literal.fromLegacyLiteral(key.getKeys().get(keyPos),
                columnType.orElseGet(() -> Type.fromPrimitiveType(key.getTypes().get(keyPos))));
    }

    private static Expression convertListPartitionToIn(PartitionItem item, Slot col, Optional<Type> columnType) {
        List<Expression> inValues = ((ListPartitionItem) item).getItems().stream()
                .map(key -> convertPartitionKeyToLiteral(key, 0, columnType))
                .collect(ImmutableList.toImmutableList());
        List<Expression> predicates = new ArrayList<>();
        if (inValues.stream().anyMatch(NullLiteral.class::isInstance)) {
            inValues = inValues.stream()
                    .filter(e -> !(e instanceof NullLiteral))
                    .collect(Collectors.toList());
            Expression isNullPredicate = new IsNull(col);
            predicates.add(isNullPredicate);
        }
        if (!inValues.isEmpty()) {
            predicates.add(new InPredicate(col, inValues));
        }
        if (predicates.isEmpty()) {
            return BooleanLiteral.of(true);
        }
        return ExpressionUtils.or(predicates);
    }

    private static Expression convertRangePartitionToCompare(PartitionItem item, Slot col,
            Optional<Type> columnType) {
        Range<PartitionKey> range = item.getItems();
        List<Expression> expressions = new ArrayList<>();
        if (range.hasLowerBound() && !range.lowerEndpoint().isMinValue()) {
            PartitionKey key = range.lowerEndpoint();
            expressions.add(new GreaterThanEqual(col, convertPartitionKeyToLiteral(key, 0, columnType)));
        }
        if (range.hasUpperBound() && !range.upperEndpoint().isMaxValue()) {
            PartitionKey key = range.upperEndpoint();
            expressions.add(new LessThan(col, convertPartitionKeyToLiteral(key, 0, columnType)));
        }
        if (expressions.isEmpty()) {
            return BooleanLiteral.of(true);
        }
        Expression predicate = ExpressionUtils.and(expressions);
        // The partition without can be the first partition of LESS THAN PARTITIONS
        // The null value can insert into this partition, so we need to add or is null condition
        if (!range.hasLowerBound() || range.lowerEndpoint().isMinValue()) {
            predicate = ExpressionUtils.or(predicate, new IsNull(col));
        }
        return predicate;
    }

    /**
     * Add predicates on base table when mv can partition update, Also support plan that contain cte and view
     */
    public static class PredicateAdder extends DefaultPlanRewriter<PredicateAddContext> {

        // record view and cte name parts, these should be ignored and visit it's actual plan
        public Set<List<String>> virtualRelationNamePartSet = new HashSet<>();

        @Override
        public Plan visitUnboundRelation(UnboundRelation unboundRelation, PredicateAddContext predicates) {

            if (predicates.getPredicates() == null || predicates.getPredicates().isEmpty()) {
                return unboundRelation;
            }
            if (virtualRelationNamePartSet.contains(unboundRelation.getNameParts())) {
                return unboundRelation;
            }
            List<String> tableQualifier = RelationUtil.getQualifierName(ConnectContext.get(),
                    unboundRelation.getNameParts());
            TableIf table = RelationUtil.getTable(tableQualifier, Env.getCurrentEnv(), Optional.empty());
            if (predicates.getPredicates().containsKey(table)) {
                return new LogicalFilter<>(
                        ExpressionUtils.extractConjunctionToSet(
                                ExpressionUtils.or(predicates.getPredicates().get(table))
                        ),
                        unboundRelation
                );
            }
            return unboundRelation;
        }

        @Override
        public Plan visitLogicalCTE(LogicalCTE<? extends Plan> cte, PredicateAddContext predicates) {
            if (predicates.isEmpty()) {
                return cte;
            }
            List<LogicalSubQueryAlias<Plan>> rewrittenSubQueryAlias = new ArrayList<>();
            for (LogicalSubQueryAlias<Plan> subQueryAlias : cte.getAliasQueries()) {
                List<Plan> subQueryAliasChildren = new ArrayList<>();
                this.virtualRelationNamePartSet.add(subQueryAlias.getQualifier());
                subQueryAlias.children().forEach(subQuery ->
                        subQueryAliasChildren.add(subQuery.accept(this, predicates))
                );
                rewrittenSubQueryAlias.add(subQueryAlias.withChildren(subQueryAliasChildren));
            }
            return super.visitLogicalCTE(new LogicalCTE<>(cte.isRecursive(),
                    rewrittenSubQueryAlias, cte.child()), predicates);
        }

        @Override
        public Plan visitLogicalSubQueryAlias(LogicalSubQueryAlias<? extends Plan> subQueryAlias,
                PredicateAddContext predicates) {
            if (predicates.isEmpty()) {
                return subQueryAlias;
            }
            this.virtualRelationNamePartSet.add(subQueryAlias.getQualifier());
            return super.visitLogicalSubQueryAlias(subQueryAlias, predicates);
        }

        /**
         * The ranges of the MV partitions this compensation removes, as a predicate on the base table's
         * partition column, or nothing when they cannot be written on one column.
         *
         * <p>A default partition's rows belong to whichever MV partition their own key falls in, so this is
         * what they are read through: pinning them to the partition's own key -- the sentinel those rows were
         * placed by -- reads none of them, and reading them whole would add the rows of the MV partitions the
         * plan still has. Only an MV partitioned by this one column can be written that way here.
         */
        private static Set<Expression> mvPartitionsToReadThrough(PredicateAddContext predicates,
                BaseColInfo relatedTableColumnInfo, Slot partitionSlot) {
            Set<Expression> res = Sets.newHashSet();
            if (predicates.getMvPartitionsToRemove().isEmpty()) {
                return res;
            }
            for (Map.Entry<BaseTableInfo, Set<String>> entry
                    : predicates.getMvPartitionsToRemove().entrySet()) {
                try {
                    TableIf table = MTMVUtil.getTable(entry.getKey());
                    if (!(table instanceof MTMV)) {
                        continue;
                    }
                    MTMV mtmv = (MTMV) table;
                    if (mtmv.getPartitionColumns().size() != 1
                            || !mtmv.getPartitionColumns().get(0).getName()
                                    .equalsIgnoreCase(relatedTableColumnInfo.getColName())) {
                        continue;
                    }
                    Type columnType = mtmv.getPartitionColumns().get(0).getType();
                    for (String partitionName : entry.getValue()) {
                        PartitionItem item = mtmv.getPartitionItemOrAnalysisException(partitionName);
                        if (!(item instanceof ListPartitionItem)) {
                            continue;
                        }
                        List<Expression> values = ((ListPartitionItem) item).getItems().stream()
                                .map(key -> convertPartitionKeyToLiteral(key, 0, Optional.of(columnType)))
                                .collect(Collectors.toList());
                        res.add(new InPredicate(partitionSlot, values));
                    }
                } catch (Exception e) {
                    // The ranges are what narrows this read; if they cannot be read, the read is left as it
                    // was rather than narrowed to a guess.
                    LOG.warn("Failed to read the removed MV partitions for the base union filter", e);
                    return Sets.newHashSet();
                }
            }
            return res;
        }

        @Override
        public Plan visitLogicalCatalogRelation(LogicalCatalogRelation catalogRelation,
                PredicateAddContext predicates) {
            if (predicates.isEmpty()) {
                return catalogRelation;
            }
            TableIf table = catalogRelation.getTable();
            if (predicates.getPredicates() != null) {
                if (predicates.getPredicates().containsKey(table)) {
                    return new LogicalFilter<>(
                            ExpressionUtils.extractConjunctionToSet(
                                    ExpressionUtils.or(predicates.getPredicates().get(table))
                            ),
                            catalogRelation);
                }
            }
            if (predicates.getPartitions() != null) {
                if (!(table instanceof MTMVRelatedTableIf)) {
                    return catalogRelation;
                }
                for (Map.Entry<BaseColInfo, Set<String>> filterTableEntry : predicates.getPartitions().entrySet()) {
                    BaseColInfo relatedTableColumnInfo = filterTableEntry.getKey();
                    if (!Objects.equals(new BaseTableInfo(table), relatedTableColumnInfo.getTableInfo())) {
                        continue;
                    }
                    Slot partitionSlot = null;
                    for (Slot slot : catalogRelation.getOutput()) {
                        if (slot.getName().equals(relatedTableColumnInfo.getColName())) {
                            partitionSlot = slot;
                            break;
                        }
                    }
                    if (partitionSlot == null) {
                        predicates.setHandleSuccess(false);
                        return catalogRelation;
                    }
                    // if partition has no data, doesn't add filter
                    Set<PartitionItem> partitionHasDataItems = new HashSet<>();
                    MTMVRelatedTableIf targetTable = (MTMVRelatedTableIf) table;
                    for (String partitionName : filterTableEntry.getValue()) {
                        if (targetTable instanceof OlapTable) {
                            Partition partition = targetTable.getPartition(partitionName);
                            if (partition == null) {
                                // partition maybe deleted, skip it
                                continue;
                            }
                            if (!((OlapTable) targetTable).selectNonEmptyPartitionIds(
                                    Lists.newArrayList(partition.getId()), Optional.empty()).isEmpty()) {
                                // Add filter only when partition has data when olap table
                                partitionHasDataItems.add(
                                        ((OlapTable) targetTable).getPartitionInfo().getItem(partition.getId()));
                            }
                        }
                        if (targetTable instanceof ExternalTable) {
                            PartitionItem partitionItem = ((ExternalTable) targetTable).getNameToPartitionItems(
                                    MvccUtil.getSnapshotFromContext(targetTable)).get(partitionName);
                            // Add filter only when partition has data when external table
                            if (partitionItem != null) {
                                partitionHasDataItems.add(partitionItem);
                            }
                        }
                    }
                    if (partitionHasDataItems.isEmpty()) {
                        predicates.setNeedAddFilter(false);
                    }
                    if (!partitionHasDataItems.isEmpty()) {
                        // The partitions are pinned the way a refresh pins them: the whole key of each, at
                        // the partition column's own type. A predicate on the MV's partition column alone
                        // reads the partitions that differ in the other keys too, and those rows are ones
                        // this branch must not add -- the MV branch of the union already supplies them, or
                        // the compensation would count them twice -- while a key written with a scale has to
                        // be compared at that scale or its rows are read as none.
                        // A list partitioned table's default partition is not pinned the way the others
                        // are: its own key is the sentinel the rows no other partition claims were placed
                        // by, so pinning it to that key reads none of them. It is read the way it was
                        // before the whole key was pinned -- on the MV's partition column alone -- which
                        // can be seen to be too wide rather than one that drops its rows.
                        boolean hasDefaultPartition = partitionHasDataItems.stream()
                                .anyMatch(PartitionItem::isDefaultPartition);
                        Set<Expression> mvPartitionPredicates = hasDefaultPartition
                                ? mvPartitionsToReadThrough(predicates, relatedTableColumnInfo, partitionSlot)
                                : Sets.newHashSet();
                        Set<Expression> preds;
                        if (!mvPartitionPredicates.isEmpty()) {
                            preds = mvPartitionPredicates;
                        } else if (targetTable instanceof OlapTable && !hasDefaultPartition) {
                            preds = constructPredicatesOfBasePartitions(partitionHasDataItems,
                                    (OlapTable) targetTable, relatedTableColumnInfo.getColName());
                        } else {
                            preds = constructPredicates(partitionHasDataItems, partitionSlot);
                        }
                        return new LogicalFilter<>(
                                ExpressionUtils.extractConjunctionToSet(ExpressionUtils.or(preds)),
                                catalogRelation
                        );
                    }
                }
            }
            return catalogRelation;
        }
    }

    /**
     * Predicate context, which support add predicate by expression or by partition name
     * Add by predicates has high priority
     */
    public static class PredicateAddContext {

        private final Map<TableIf, Set<Expression>> predicates;
        private final Map<BaseColInfo, Set<String>> partitions;
        // The MV partitions this compensation takes out of the rewritten plan, by MV. A base partition that
        // takes the rows no other partition of it claims cannot be pinned through its own key -- that key is
        // the sentinel those rows were placed by -- so it is read through the ranges of these MV partitions.
        private final Map<BaseTableInfo, Set<String>> mvPartitionsToRemove;
        private boolean handleSuccess = true;
        // when add filter by partition, if partition has no data, doesn't need to add filter. should be false
        private boolean needAddFilter = true;

        public PredicateAddContext(Map<TableIf, Set<Expression>> predicates,
                Map<BaseColInfo, Set<String>> partitions,
                Map<BaseTableInfo, Set<String>> mvPartitionsToRemove) {
            this.predicates = predicates;
            this.partitions = partitions;
            this.mvPartitionsToRemove = mvPartitionsToRemove;
        }

        public Map<TableIf, Set<Expression>> getPredicates() {
            return predicates;
        }

        public Map<BaseColInfo, Set<String>> getPartitions() {
            return partitions;
        }

        public Map<BaseTableInfo, Set<String>> getMvPartitionsToRemove() {
            return mvPartitionsToRemove == null ? ImmutableMap.of() : mvPartitionsToRemove;
        }

        public boolean isEmpty() {
            return predicates == null && partitions == null;
        }

        public boolean isHandleSuccess() {
            return handleSuccess;
        }

        public void setHandleSuccess(boolean handleSuccess) {
            this.handleSuccess = handleSuccess;
        }

        public boolean isNeedAddFilter() {
            return needAddFilter;
        }

        public void setNeedAddFilter(boolean needAddFilter) {
            this.needAddFilter = needAddFilter;
        }
    }
}