This is an automated email from the ASF dual-hosted git repository. suxiaogang223 pushed a commit to branch codex/paimon-row-level-dml in repository https://gitbox.apache.org/repos/asf/doris.git
commit 8ff74c6abee342636831f9ac51178b5e560bb174 Author: suxiaogang <[email protected]> AuthorDate: Thu Aug 13 13:21:47 2026 +0800 fix: harden Paimon row-level DML planning --- .../analysis/PaimonRowChangeCapabilities.java | 28 +++- .../rules/analysis/PaimonRowChangePlanBuilder.java | 20 +-- .../trees/plans/commands/PaimonDeleteCommand.java | 13 +- .../trees/plans/commands/PaimonMergeCommand.java | 7 +- .../trees/plans/commands/PaimonUpdateCommand.java | 7 +- .../plans/commands/info/PaimonRowChangeSpec.java | 28 ++-- .../test_paimon_write_row_level_dml.groovy | 148 +++++++++++++++++++++ 7 files changed, 215 insertions(+), 36 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/PaimonRowChangeCapabilities.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/PaimonRowChangeCapabilities.java index 06b843f4110..e2a7e5d1e46 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/PaimonRowChangeCapabilities.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/PaimonRowChangeCapabilities.java @@ -34,6 +34,7 @@ import org.apache.doris.qe.ConnectContext; import org.apache.paimon.CoreOptions; import org.apache.paimon.options.Options; +import org.apache.paimon.table.BucketMode; import org.apache.paimon.table.FileStoreTable; import java.util.Collection; @@ -49,6 +50,7 @@ final class PaimonRowChangeCapabilities { static void check(PaimonWriteTarget target, PaimonRowChangeSpec spec, CascadesContext cascadesContext) { requireNoDataMask(target, spec, cascadesContext); + requireCompleteKeyDynamicIndex(target.getTable(), spec); if (spec instanceof PaimonRowChangeSpec.Update) { checkUpdate(target, updatedColumns(((PaimonRowChangeSpec.Update) spec).getAssignments())); @@ -108,10 +110,17 @@ final class PaimonRowChangeCapabilities { throw new AnalysisException("Paimon UPDATE cannot modify sequence-field column '" + column + "'"); } - if (partitionKeys.contains(column) && options.bucket() != -1) { - throw new AnalysisException("Paimon UPDATE cannot modify partition column '" - + column + "' unless bucket=-1 because the old partition row cannot be " - + "removed without an UPDATE_BEFORE record"); + if (partitionKeys.contains(column)) { + if (options.bucket() != -1) { + throw new AnalysisException("Paimon UPDATE cannot modify partition column '" + + column + "' unless bucket=-1 because the old partition row cannot be " + + "removed without an UPDATE_BEFORE record"); + } + if (options.ignoreDelete()) { + throw new AnalysisException("Paimon UPDATE cannot modify partition column '" + + column + "' when ignore-delete=true because removing the old " + + "partition row requires a delete record"); + } } } CoreOptions.MergeEngine engine = options.mergeEngine(); @@ -173,6 +182,17 @@ final class PaimonRowChangeCapabilities { } } + private static void requireCompleteKeyDynamicIndex( + FileStoreTable table, PaimonRowChangeSpec spec) { + CoreOptions options = CoreOptions.fromMap(table.options()); + if (table.bucketMode() == BucketMode.KEY_DYNAMIC + && options.crossPartitionUpsertIndexTtl() != null) { + throw new AnalysisException("Paimon " + spec.getDmlCommandType() + + " is not supported when cross-partition-upsert.index-ttl is configured " + + "because row-change DML requires a complete key-dynamic index"); + } + } + private static void requireNoRowKindField(CoreOptions options, String operation) { if (options.rowkindField().isPresent()) { throw new AnalysisException("Paimon " + operation diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/PaimonRowChangePlanBuilder.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/PaimonRowChangePlanBuilder.java index a339c229319..d51961fb2a7 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/PaimonRowChangePlanBuilder.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/PaimonRowChangePlanBuilder.java @@ -18,7 +18,6 @@ package org.apache.doris.nereids.rules.analysis; import org.apache.doris.catalog.Column; -import org.apache.doris.common.util.Util; import org.apache.doris.datasource.paimon.PaimonRowChangeOperation; import org.apache.doris.datasource.paimon.PaimonWriteTarget; import org.apache.doris.nereids.CascadesContext; @@ -80,16 +79,13 @@ final class PaimonRowChangePlanBuilder { } } - String targetName = update.getTableAlias() != null - ? update.getTableAlias() - : Util.getTempTableDisplayName(target.getDorisTable().getName()); ExpressionAnalyzer analyzer = expressionAnalyzer(child, cascadesContext); List<NamedExpression> projects = new ArrayList<>(); projects.add(operation(PaimonRowChangeOperation.UPDATE)); for (Column column : target.getSchema()) { Expression value = changes.remove(column.getName()); if (value == null) { - value = new UnboundSlot(targetName, column.getName()); + value = targetSlot(update.getTargetNameInPlan(), column.getName()); } projects.add(bindColumn(analyzer, value, column)); } @@ -103,17 +99,21 @@ final class PaimonRowChangePlanBuilder { private static LogicalProject<?> buildDelete(PaimonWriteTarget target, PaimonRowChangeSpec.Delete delete, LogicalPlan child, CascadesContext cascadesContext) { - String targetName = delete.getTableAlias() != null - ? delete.getTableAlias() - : Util.getTempTableDisplayName(target.getDorisTable().getName()); ExpressionAnalyzer analyzer = expressionAnalyzer(child, cascadesContext); List<NamedExpression> projects = new ArrayList<>(); projects.add(operation(PaimonRowChangeOperation.DELETE)); for (Column column : target.getSchema()) { projects.add(bindColumn( - analyzer, new UnboundSlot(targetName, column.getName()), column)); + analyzer, targetSlot(delete.getTargetNameInPlan(), column.getName()), column)); } - return new LogicalProject<>(projects, child); + // SQL DELETE changes each target row once even when a USING join matches it repeatedly. + return new LogicalProject<>(projects, true, child); + } + + private static UnboundSlot targetSlot(List<String> targetNameInPlan, String columnName) { + List<String> nameParts = new ArrayList<>(targetNameInPlan); + nameParts.add(columnName); + return new UnboundSlot(nameParts); } private static Alias operation(byte operation) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonDeleteCommand.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonDeleteCommand.java index 0bd50522d38..3129405a573 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonDeleteCommand.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonDeleteCommand.java @@ -26,9 +26,12 @@ import org.apache.doris.nereids.trees.plans.commands.info.PaimonRowChangeSpec; import org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTableCommand; import org.apache.doris.nereids.trees.plans.logical.LogicalPlan; import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor; +import org.apache.doris.nereids.util.RelationUtil; import org.apache.doris.qe.ConnectContext; import org.apache.doris.qe.StmtExecutor; +import com.google.common.collect.ImmutableList; + import java.util.List; import java.util.Optional; @@ -48,18 +51,20 @@ public class PaimonDeleteCommand extends Command implements ForwardWithSync, Exp @Override public void run(ConnectContext ctx, StmtExecutor executor) throws Exception { - new InsertIntoTableCommand(buildPlan(), Optional.empty(), Optional.empty(), Optional.empty()) + new InsertIntoTableCommand(buildPlan(ctx), Optional.empty(), Optional.empty(), Optional.empty()) .run(ctx, executor); } - private LogicalPlan buildPlan() { + private LogicalPlan buildPlan(ConnectContext ctx) { + List<String> targetNameInPlan = tableAlias != null + ? ImmutableList.of(tableAlias) : RelationUtil.getQualifierName(ctx, nameParts); return new UnboundPaimonTableSink<>(nameParts, logicalQuery, - new PaimonRowChangeSpec.Delete(tableAlias)); + new PaimonRowChangeSpec.Delete(targetNameInPlan)); } @Override public Plan getExplainPlan(ConnectContext ctx) { - return buildPlan(); + return buildPlan(ctx); } @Override diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonMergeCommand.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonMergeCommand.java index 8233f54346f..d85d68777ad 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonMergeCommand.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonMergeCommand.java @@ -38,6 +38,7 @@ 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.visitor.PlanVisitor; +import org.apache.doris.nereids.util.RelationUtil; import org.apache.doris.qe.ConnectContext; import org.apache.doris.qe.StmtExecutor; @@ -50,7 +51,6 @@ import java.util.Optional; public class PaimonMergeCommand extends Command implements ForwardWithSync, Explainable { private final List<String> targetNameParts; private final Optional<String> targetAlias; - private final List<String> targetNameInPlan; private final Optional<LogicalPlan> cte; private final LogicalPlan source; private final Expression onClause; @@ -65,8 +65,6 @@ public class PaimonMergeCommand extends Command implements ForwardWithSync, Expl super(PlanType.MERGE_INTO_COMMAND); this.targetNameParts = targetNameParts; this.targetAlias = targetAlias; - this.targetNameInPlan = targetAlias.isPresent() - ? ImmutableList.of(targetAlias.get()) : targetNameParts; this.cte = cte; this.source = source; this.onClause = onClause; @@ -88,6 +86,9 @@ public class PaimonMergeCommand extends Command implements ForwardWithSync, Expl targetNameParts, targetAlias.orElse(null)); } } + List<String> targetNameInPlan = targetAlias.isPresent() + ? ImmutableList.of(targetAlias.get()) + : RelationUtil.getQualifierName(ctx, targetNameParts); PaimonRowChangeSpec.Merge spec = new PaimonRowChangeSpec.Merge( targetNameInPlan, matchedClauses, notMatchedClauses); return new UnboundPaimonTableSink<>(targetNameParts, generateBasePlan(), spec); diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonUpdateCommand.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonUpdateCommand.java index d2eff4289ae..e5695839f08 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonUpdateCommand.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/PaimonUpdateCommand.java @@ -28,9 +28,12 @@ import org.apache.doris.nereids.trees.plans.commands.info.PaimonRowChangeSpec; import org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTableCommand; import org.apache.doris.nereids.trees.plans.logical.LogicalPlan; import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor; +import org.apache.doris.nereids.util.RelationUtil; import org.apache.doris.qe.ConnectContext; import org.apache.doris.qe.StmtExecutor; +import com.google.common.collect.ImmutableList; + import java.util.List; import java.util.Optional; @@ -64,8 +67,10 @@ public class PaimonUpdateCommand extends Command implements ForwardWithSync, Exp UpdateCommand.checkAssignmentColumn(ctx, ((UnboundSlot) assignment.left()).getNameParts(), nameParts, tableAlias); } + List<String> targetNameInPlan = tableAlias != null + ? ImmutableList.of(tableAlias) : RelationUtil.getQualifierName(ctx, nameParts); return new UnboundPaimonTableSink<>(nameParts, logicalQuery, - new PaimonRowChangeSpec.Update(tableAlias, assignments)); + new PaimonRowChangeSpec.Update(targetNameInPlan, assignments)); } @Override diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/PaimonRowChangeSpec.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/PaimonRowChangeSpec.java index 853b53ab298..18e9e42a57c 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/PaimonRowChangeSpec.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/PaimonRowChangeSpec.java @@ -36,16 +36,16 @@ public abstract class PaimonRowChangeSpec { /** UPDATE description. */ public static final class Update extends PaimonRowChangeSpec { - private final String tableAlias; + private final List<String> targetNameInPlan; private final List<EqualTo> assignments; - public Update(String tableAlias, List<EqualTo> assignments) { - this.tableAlias = tableAlias; + public Update(List<String> targetNameInPlan, List<EqualTo> assignments) { + this.targetNameInPlan = Utils.copyRequiredList(targetNameInPlan); this.assignments = Utils.copyRequiredList(assignments); } - public String getTableAlias() { - return tableAlias; + public List<String> getTargetNameInPlan() { + return targetNameInPlan; } public List<EqualTo> getAssignments() { @@ -71,26 +71,26 @@ public abstract class PaimonRowChangeSpec { return false; } Update update = (Update) other; - return Objects.equals(tableAlias, update.tableAlias) + return Objects.equals(targetNameInPlan, update.targetNameInPlan) && Objects.equals(assignments, update.assignments); } @Override public int hashCode() { - return Objects.hash(tableAlias, assignments); + return Objects.hash(targetNameInPlan, assignments); } } /** DELETE description. */ public static final class Delete extends PaimonRowChangeSpec { - private final String tableAlias; + private final List<String> targetNameInPlan; - public Delete(String tableAlias) { - this.tableAlias = tableAlias; + public Delete(List<String> targetNameInPlan) { + this.targetNameInPlan = Utils.copyRequiredList(targetNameInPlan); } - public String getTableAlias() { - return tableAlias; + public List<String> getTargetNameInPlan() { + return targetNameInPlan; } @Override @@ -106,12 +106,12 @@ public abstract class PaimonRowChangeSpec { @Override public boolean equals(Object other) { return other instanceof Delete - && Objects.equals(tableAlias, ((Delete) other).tableAlias); + && Objects.equals(targetNameInPlan, ((Delete) other).targetNameInPlan); } @Override public int hashCode() { - return Objects.hash(tableAlias); + return Objects.hash(targetNameInPlan); } } diff --git a/regression-test/suites/paimon_write/test_paimon_write_row_level_dml.groovy b/regression-test/suites/paimon_write/test_paimon_write_row_level_dml.groovy index de3e774891c..61ae3e97894 100644 --- a/regression-test/suites/paimon_write/test_paimon_write_row_level_dml.groovy +++ b/regression-test/suites/paimon_write/test_paimon_write_row_level_dml.groovy @@ -163,6 +163,37 @@ suite("test_paimon_write_row_level_dml", "p0,external,paimon") { 'bucket' = '2', 'bucket-key' = 'id' ); + + DROP TABLE IF EXISTS paimon.${dbName}.t_cross_partition_ttl; + CREATE TABLE paimon.${dbName}.t_cross_partition_ttl ( + pt STRING, id INT, name STRING + ) USING paimon + PARTITIONED BY (pt) + TBLPROPERTIES ( + 'primary-key' = 'id', + 'bucket' = '-1', + 'cross-partition-upsert.index-ttl' = '1 h' + ); + + DROP TABLE IF EXISTS paimon.${dbName}.t_cross_partition_ignore_delete; + CREATE TABLE paimon.${dbName}.t_cross_partition_ignore_delete ( + pt STRING, id INT, name STRING + ) USING paimon + PARTITIONED BY (pt) + TBLPROPERTIES ( + 'primary-key' = 'id', + 'bucket' = '-1', + 'ignore-delete' = 'true' + ); + + DROP TABLE IF EXISTS paimon.${dbName}.t_same_name; + CREATE TABLE paimon.${dbName}.t_same_name ( + id INT, name STRING + ) USING paimon + TBLPROPERTIES ( + 'primary-key' = 'id', + 'bucket' = '1' + ); """ sql """drop catalog if exists ${catalogName}""" @@ -195,8 +226,49 @@ suite("test_paimon_write_row_level_dml", "p0,external,paimon") { (2, 'ignored', 0, 'D'), (5, 'Eve', 50, 'I') """ + sql """drop table if exists internal.${dbName}.t_delete_source""" + sql """ + CREATE TABLE internal.${dbName}.t_delete_source ( + id INT + ) ENGINE=OLAP + DUPLICATE KEY(id) + DISTRIBUTED BY HASH(id) BUCKETS 1 + PROPERTIES('replication_num'='1') + """ + sql """INSERT INTO internal.${dbName}.t_delete_source VALUES (10), (10)""" + sql """drop table if exists internal.${dbName}.t_same_name""" + sql """ + CREATE TABLE internal.${dbName}.t_same_name ( + id INT, name STRING + ) ENGINE=OLAP + DUPLICATE KEY(id) + DISTRIBUTED BY HASH(id) BUCKETS 1 + PROPERTIES('replication_num'='1') + """ try { + def latestSnapshotId = { String tableName -> + def rows = spark_paimon """ + SELECT max(snapshot_id) + FROM paimon.${dbName}.`${tableName}\$snapshots` + """ + assertEquals(1, rows.size()) + assertTrue(rows[0][0] != null) + return rows[0][0].toString() + } + + def incrementalAuditLog = { tableName, columns, beforeSnapshot, afterSnapshot -> + return spark_paimon """ + SELECT ${columns} + FROM paimon_incremental_query( + 'paimon.${dbName}.`${tableName}\$audit_log`', + '${beforeSnapshot}', + '${afterSnapshot}' + ) + ORDER BY id, rowkind + """ + } + sql """INSERT INTO t_dml VALUES (1, 'Alice', 10, 'active'), (2, 'Bob', 20, 'active'), @@ -258,6 +330,80 @@ suite("test_paimon_write_row_level_dml", "p0,external,paimon") { """ exception "Paimon UPDATE cannot modify partition column 'pt' unless bucket=-1" } + test { + sql """UPDATE t_cross_partition_ttl SET name = 'updated' WHERE id = 1""" + exception "cross-partition-upsert.index-ttl is configured" + } + test { + sql """DELETE FROM t_cross_partition_ttl WHERE id = 1""" + exception "cross-partition-upsert.index-ttl is configured" + } + test { + sql """ + MERGE INTO t_cross_partition_ttl t + USING internal.${dbName}.t_merge_source s ON t.id = s.id + WHEN MATCHED THEN DELETE + """ + exception "cross-partition-upsert.index-ttl is configured" + } + test { + sql """ + UPDATE t_cross_partition_ignore_delete + SET pt = 'new_pt' WHERE id = 1 + """ + exception "cannot modify partition column 'pt' when ignore-delete=true" + } + test { + sql """ + MERGE INTO t_cross_partition_ignore_delete t + USING internal.${dbName}.t_merge_source s ON t.id = s.id + WHEN MATCHED THEN UPDATE SET pt = 'new_pt', name = s.name + """ + exception "cannot modify partition column 'pt' when ignore-delete=true" + } + + sql """INSERT INTO t_input_changelog VALUES (10, 'delete-me')""" + String deleteInputBefore = latestSnapshotId("t_input_changelog") + sql """ + DELETE FROM t_input_changelog t + USING internal.${dbName}.t_delete_source s + WHERE t.id = s.id + """ + String deleteInputAfter = latestSnapshotId("t_input_changelog") + assertEquals([ + ["-D", 10, "delete-me"] + ], incrementalAuditLog("t_input_changelog", "rowkind, id, name", + deleteInputBefore, deleteInputAfter)) + assertEquals([[0L]], sql("SELECT count(*) FROM t_input_changelog")) + + sql """INSERT INTO t_same_name VALUES (1, 'old')""" + sql """INSERT INTO internal.${dbName}.t_same_name VALUES (1, 'from-update')""" + sql """ + UPDATE ${catalogName}.${dbName}.t_same_name + SET name = internal.${dbName}.t_same_name.name + FROM internal.${dbName}.t_same_name + WHERE ${catalogName}.${dbName}.t_same_name.id + = internal.${dbName}.t_same_name.id + """ + assertEquals([[1, "from-update"]], sql("SELECT * FROM t_same_name")) + sql """TRUNCATE TABLE internal.${dbName}.t_same_name""" + sql """INSERT INTO internal.${dbName}.t_same_name VALUES (1, 'from-merge')""" + sql """ + MERGE INTO t_same_name + USING internal.${dbName}.t_same_name + ON ${catalogName}.${dbName}.t_same_name.id + = internal.${dbName}.t_same_name.id + WHEN MATCHED THEN UPDATE + SET name = internal.${dbName}.t_same_name.name + """ + assertEquals([[1, "from-merge"]], sql("SELECT * FROM t_same_name")) + sql """ + DELETE FROM ${catalogName}.${dbName}.t_same_name + USING internal.${dbName}.t_same_name + WHERE ${catalogName}.${dbName}.t_same_name.id + = internal.${dbName}.t_same_name.id + """ + assertEquals([[0L]], sql("SELECT count(*) FROM t_same_name")) sql """ UPDATE t_dml @@ -475,5 +621,7 @@ suite("test_paimon_write_row_level_dml", "p0,external,paimon") { } finally { sql """drop catalog if exists ${catalogName}""" sql """drop table if exists internal.${dbName}.t_merge_source""" + sql """drop table if exists internal.${dbName}.t_delete_source""" + sql """drop table if exists internal.${dbName}.t_same_name""" } } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
