PaimonRowLevelDmlTransform.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.TableIf;
import org.apache.doris.connector.spi.DorisConnectorException;
import org.apache.doris.connector.spi.handle.WriteOperation;
import org.apache.doris.connector.spi.pushdown.ConnectorPredicate;
import org.apache.doris.connector.spi.write.ConnectorRowChangeStyle;
import org.apache.doris.connector.spi.write.ConnectorRowLevelDmlRequest;
import org.apache.doris.datasource.plugin.PluginDrivenExternalTable;
import org.apache.doris.nereids.NereidsPlanner;
import org.apache.doris.nereids.analyzer.UnboundConnectorTableSink;
import org.apache.doris.nereids.analyzer.UnboundRelation;
import org.apache.doris.nereids.analyzer.UnboundSlot;
import org.apache.doris.nereids.exceptions.AnalysisException;
import org.apache.doris.nereids.parser.LogicalPlanBuilderAssistant;
import org.apache.doris.nereids.rules.exploration.join.JoinReorderContext;
import org.apache.doris.nereids.trees.expressions.EqualTo;
import org.apache.doris.nereids.trees.expressions.StatementScopeIdGenerator;
import org.apache.doris.nereids.trees.plans.JoinType;
import org.apache.doris.nereids.trees.plans.Plan;
import org.apache.doris.nereids.trees.plans.commands.info.ConnectorChangelogRowChangeSpec;
import org.apache.doris.nereids.trees.plans.commands.insert.BaseExternalTableInsertExecutor;
import org.apache.doris.nereids.trees.plans.commands.insert.PluginDrivenInsertExecutor;
import org.apache.doris.nereids.trees.plans.commands.merge.MergeMatchedClause;
import org.apache.doris.nereids.trees.plans.logical.LogicalJoin;
import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
import org.apache.doris.nereids.trees.plans.logical.LogicalSubQueryAlias;
import org.apache.doris.nereids.trees.plans.physical.PhysicalConnectorTableSink;
import org.apache.doris.nereids.trees.plans.physical.PhysicalSink;
import org.apache.doris.nereids.util.RelationUtil;
import org.apache.doris.planner.DataSink;
import org.apache.doris.planner.PlanFragment;
import org.apache.doris.qe.ConnectContext;
import com.google.common.collect.ImmutableList;
import java.util.List;
import java.util.Optional;
import java.util.Set;
import java.util.TreeSet;
/** Changelog row-level DML transform used by the Paimon connector. */
public class PaimonRowLevelDmlTransform implements RowLevelDmlTransform {
@Override
public boolean handles(TableIf table) {
return table instanceof PluginDrivenExternalTable
&& ((PluginDrivenExternalTable) table).getConnectorRowChangeStyle()
== ConnectorRowChangeStyle.CHANGELOG;
}
@Override
public void checkMode(TableIf table, RowLevelDmlOp op) {
// The statement-specific validation runs in synthesize, where assignments and MERGE clauses are present.
}
@Override
public LogicalPlan synthesize(ConnectContext ctx, RowLevelDmlArgs args, RowLevelDmlOp op) {
PluginDrivenExternalTable table = (PluginDrivenExternalTable) args.getTable();
if (op == RowLevelDmlOp.DELETE && (args.isTempPart() || !args.getPartitions().isEmpty())) {
throw new AnalysisException(
"Paimon DELETE does not support partition name lists; use a WHERE predicate");
}
validate(table, args, op);
switch (op) {
case DELETE:
return deletePlan(ctx, args);
case UPDATE:
return updatePlan(ctx, args);
default:
return mergePlan(ctx, args);
}
}
private LogicalPlan deletePlan(ConnectContext ctx, RowLevelDmlArgs args) {
List<String> target = args.getTableAlias() != null
? ImmutableList.of(args.getTableAlias())
: RelationUtil.getQualifierName(ctx, args.getNameParts());
return new UnboundConnectorTableSink<>(args.getNameParts(), args.getLogicalQuery(),
new ConnectorChangelogRowChangeSpec.Delete(target, args.shouldDeduplicateTargetRows()));
}
private LogicalPlan updatePlan(ConnectContext ctx, RowLevelDmlArgs args) {
for (EqualTo assignment : args.getAssignments()) {
UpdateCommand.checkAssignmentColumn(ctx,
((UnboundSlot) assignment.left()).getNameParts(),
args.getNameParts(), args.getTableAlias());
}
List<String> target = args.getTableAlias() != null
? ImmutableList.of(args.getTableAlias())
: RelationUtil.getQualifierName(ctx, args.getNameParts());
LogicalPlan sink = new UnboundConnectorTableSink<>(args.getNameParts(), args.getLogicalQuery(),
new ConnectorChangelogRowChangeSpec.Update(target, args.getAssignments()));
return args.getCte().isPresent() ? (LogicalPlan) args.getCte().get().withChildren(sink) : sink;
}
private LogicalPlan mergePlan(ConnectContext ctx, RowLevelDmlArgs args) {
for (MergeMatchedClause clause : args.getMatchedClauses()) {
for (EqualTo assignment : clause.getAssignments()) {
UpdateCommand.checkAssignmentColumn(ctx,
((UnboundSlot) assignment.left()).getNameParts(),
args.getTargetNameParts(), args.getTargetAlias().orElse(null));
}
}
List<String> targetName = args.getTargetAlias().isPresent()
? ImmutableList.of(args.getTargetAlias().get())
: RelationUtil.getQualifierName(ctx, args.getTargetNameParts());
ConnectorChangelogRowChangeSpec.Merge spec = new ConnectorChangelogRowChangeSpec.Merge(
targetName, args.getMatchedClauses(), args.getNotMatchedClauses());
LogicalPlan target = LogicalPlanBuilderAssistant.withCheckPolicy(
new UnboundRelation(StatementScopeIdGenerator.newRelationId(), args.getTargetNameParts()));
if (args.getTargetAlias().isPresent()) {
target = new LogicalSubQueryAlias<>(args.getTargetAlias().get(), target);
}
JoinType joinType = args.getNotMatchedClauses().isEmpty()
? JoinType.INNER_JOIN : JoinType.LEFT_OUTER_JOIN;
LogicalPlan join = new LogicalJoin<>(joinType, ImmutableList.of(),
ImmutableList.of(args.getOnClause()), args.getSource(), target, JoinReorderContext.EMPTY);
LogicalPlan sink = new UnboundConnectorTableSink<>(args.getTargetNameParts(), join, spec);
return args.getCte().isPresent() ? (LogicalPlan) args.getCte().get().withChildren(sink) : sink;
}
private void validate(PluginDrivenExternalTable table, RowLevelDmlArgs args, RowLevelDmlOp op) {
Set<String> updatedColumns = new TreeSet<>(String.CASE_INSENSITIVE_ORDER);
boolean containsUpdate = op == RowLevelDmlOp.UPDATE;
boolean containsDelete = op == RowLevelDmlOp.DELETE;
if (op == RowLevelDmlOp.UPDATE) {
addUpdatedColumns(updatedColumns, args.getAssignments());
} else if (op == RowLevelDmlOp.MERGE) {
for (MergeMatchedClause clause : args.getMatchedClauses()) {
containsDelete |= clause.isDelete();
containsUpdate |= !clause.isDelete();
addUpdatedColumns(updatedColumns, clause.getAssignments());
}
}
try {
table.validateConnectorRowLevelDml(new ConnectorRowLevelDmlRequest(
toWriteOperation(op), updatedColumns, containsUpdate, containsDelete));
} catch (DorisConnectorException e) {
throw new AnalysisException(e.getMessage(), e);
}
}
private void addUpdatedColumns(Set<String> columns, List<EqualTo> assignments) {
for (EqualTo assignment : assignments) {
List<String> parts = ((UnboundSlot) assignment.left()).getNameParts();
columns.add(parts.get(parts.size() - 1));
}
}
private WriteOperation toWriteOperation(RowLevelDmlOp op) {
switch (op) {
case DELETE:
return WriteOperation.DELETE;
case UPDATE:
return WriteOperation.UPDATE;
default:
return WriteOperation.MERGE;
}
}
@Override
public BaseExternalTableInsertExecutor newExecutor(ConnectContext ctx, TableIf table, String label,
NereidsPlanner planner, boolean emptyInsert, RowLevelDmlOp op) {
return new PluginDrivenInsertExecutor(ctx, (PluginDrivenExternalTable) table, label,
planner, Optional.empty(), emptyInsert, -1L);
}
@Override
public PhysicalSink<?> requirePhysicalSink(NereidsPlanner planner, RowLevelDmlOp op) {
return planner.getPhysicalPlan().<PhysicalSink<?>>collect(PhysicalSink.class::isInstance)
.stream().filter(PhysicalConnectorTableSink.class::isInstance).findAny()
.orElseThrow(() -> new AnalysisException(op + " plan must use connector table sink"));
}
@Override
public String labelPrefix(RowLevelDmlOp op) {
return "paimon_" + op.name().toLowerCase();
}
@Override
public void setupConflictDetection(BaseExternalTableInsertExecutor executor, Plan analyzedPlan,
TableIf table, RowLevelDmlOp op) {
// Paimon commits reconcile conflicts through the connector transaction.
}
@Override
public void finalizeSink(BaseExternalTableInsertExecutor executor, RowLevelDmlOp op,
PlanFragment fragment, DataSink sink, PhysicalSink<?> physicalSink) {
((PluginDrivenInsertExecutor) executor).finalizeRowLevelDmlSink(fragment, sink, physicalSink);
}
@Override
public Optional<ConnectorPredicate> extractWriteConstraint(Plan analyzedPlan, TableIf table) {
return Optional.empty();
}
}