diff --git a/fe/fe-core/src/main/java/org/apache/doris/planner/OlapTableSink.java b/fe/fe-core/src/main/java/org/apache/doris/planner/OlapTableSink.java index 7a868f896d418a..161323945cdc2c 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/planner/OlapTableSink.java +++ b/fe/fe-core/src/main/java/org/apache/doris/planner/OlapTableSink.java @@ -1168,7 +1168,8 @@ public TOlapTableLocationParam createDummyLocation(OlapTable table) throws UserE // In non-cloud mode the binlog tablet must be on the same disk as its base tablet (cloud mode // does not need this), so keep only the (backend, pathHash) entries shared by both tablets. - private Multimap getBinlogColocatedReplicaBackendPathMap(Tablet baseTablet, Tablet rowBinlogTablet) + public static Multimap getBinlogColocatedReplicaBackendPathMap( + Tablet baseTablet, Tablet rowBinlogTablet) throws UserException { Multimap baseBePathsMap = baseTablet.getNormalReplicaBackendPathMap(); Multimap binlogBePathsMap = rowBinlogTablet.getNormalReplicaBackendPathMap(); diff --git a/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java b/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java index 052dfbac3def35..89daeaed6d9a02 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java +++ b/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java @@ -4786,7 +4786,8 @@ public TCreatePartitionResult createPartition(TCreatePartitionRequest request) t } } int quorum = partitionSnapshot.quorum; - for (Tablet tablet : partitionSnapshot.tablets) { + for (TabletLocationSnapshot tabletSnapshot : partitionSnapshot.tablets) { + Tablet tablet = tabletSnapshot.tablet; // we should ensure the replica backend is alive // otherwise, there will be a 'unknown node id, id=xxx' error for stream load // BE id -> path hash @@ -4803,6 +4804,9 @@ public TCreatePartitionResult createPartition(TCreatePartitionRequest request) t } bePathsMap = cloudTablet.getNormalReplicaBackendPathMapByClusterId(cachedClusterId); } + } else if (tabletSnapshot.isRowBinlog()) { + bePathsMap = OlapTableSink.getBinlogColocatedReplicaBackendPathMap( + tabletSnapshot.rowBinlogBaseTablet, tablet); } else { bePathsMap = tablet.getNormalReplicaBackendPathMap(); } @@ -4815,7 +4819,7 @@ public TCreatePartitionResult createPartition(TCreatePartitionRequest request) t if (bePathsMap.keySet().size() < quorum) { LOG.warn("auto go quorum exception"); } - partitionTablets.add(new TTabletLocation(tablet.getId(), + partitionTablets.add(tabletSnapshot.createLocation( Lists.newArrayList(bePathsMap.keySet()))); } @@ -5115,7 +5119,8 @@ public TReplacePartitionResult replacePartition(TReplacePartitionRequest request } } int quorum = partitionSnapshot.quorum; - for (Tablet tablet : partitionSnapshot.tablets) { + for (TabletLocationSnapshot tabletSnapshot : partitionSnapshot.tablets) { + Tablet tablet = tabletSnapshot.tablet; // we should ensure the replica backend is alive // otherwise, there will be a 'unknown node id, id=xxx' error for stream load // BE id -> path hash @@ -5133,6 +5138,9 @@ public TReplacePartitionResult replacePartition(TReplacePartitionRequest request bePathsMap = cloudTablet .getNormalReplicaBackendPathMapByClusterId(replaceCachedClusterId); } + } else if (tabletSnapshot.isRowBinlog()) { + bePathsMap = OlapTableSink.getBinlogColocatedReplicaBackendPathMap( + tabletSnapshot.rowBinlogBaseTablet, tablet); } else { bePathsMap = tablet.getNormalReplicaBackendPathMap(); } @@ -5145,7 +5153,7 @@ public TReplacePartitionResult replacePartition(TReplacePartitionRequest request if (bePathsMap.keySet().size() < quorum) { LOG.warn("auto go quorum exception"); } - partitionTablets.add(new TTabletLocation(tablet.getId(), + partitionTablets.add(tabletSnapshot.createLocation( Lists.newArrayList(bePathsMap.keySet()))); } @@ -5214,12 +5222,13 @@ private static final class PartitionResultSnapshot { private final Partition partition; private final long partitionId; private final TOlapTablePartition tPartition; - private final List tablets; + private final List tablets; private final int quorum; private final boolean cacheLoadTabletIdx; private PartitionResultSnapshot(Partition partition, long partitionId, - TOlapTablePartition tPartition, List tablets, int quorum, boolean cacheLoadTabletIdx) { + TOlapTablePartition tPartition, List tablets, int quorum, + boolean cacheLoadTabletIdx) { this.partition = partition; this.partitionId = partitionId; this.tPartition = tPartition; @@ -5229,6 +5238,28 @@ private PartitionResultSnapshot(Partition partition, long partitionId, } } + private static final class TabletLocationSnapshot { + private final Tablet tablet; + private final Tablet rowBinlogBaseTablet; + + private TabletLocationSnapshot(Tablet tablet, Tablet rowBinlogBaseTablet) { + this.tablet = tablet; + this.rowBinlogBaseTablet = rowBinlogBaseTablet; + } + + private boolean isRowBinlog() { + return rowBinlogBaseTablet != null; + } + + private TTabletLocation createLocation(List nodeIds) { + TTabletLocation location = new TTabletLocation(tablet.getId(), nodeIds); + if (isRowBinlog()) { + location.setBaseTabletId(rowBinlogBaseTablet.getId()); + } + return location; + } + } + private static List snapshotPartitionResultsByName(OlapTable olapTable, Collection partitionNames, boolean loadToSingleTablet, boolean enableAdaptiveRandomBucket, String resultName) throws UserException { @@ -5281,13 +5312,26 @@ private static PartitionResultSnapshot snapshotPartitionResult(PartitionInfo par TOlapTablePartition tPartition = new TOlapTablePartition(); tPartition.setId(partitionId); OlapTableSink.setPartitionKeys(tPartition, partitionItem, partColNum); - List partitionTabletSnapshot = new ArrayList<>(); + List partitionTabletSnapshot = new ArrayList<>(); for (MaterializedIndex index : partition.getMaterializedIndices(MaterializedIndex.IndexExtState.ALL, true)) { List indexTablets = new ArrayList<>(index.getTablets()); - tPartition.addToIndexes(new TOlapTableIndexTablets(index.getId(), Lists.newArrayList( - indexTablets.stream().map(Tablet::getId).collect(Collectors.toList())))); - tPartition.setNumBuckets(indexTablets.size()); - partitionTabletSnapshot.addAll(indexTablets); + if (index.isRowBinlog()) { + for (Tablet tablet : indexTablets) { + long baseTabletId = tablet.getRowBinlogBaseTabletId(); + Tablet baseTablet = partition.getBaseIndex().getTablet(baseTabletId); + Preconditions.checkNotNull(baseTablet, + "row binlog tablet %s's base tablet %s can not be found in partition %s", + tablet.getId(), baseTabletId, partitionId); + partitionTabletSnapshot.add(new TabletLocationSnapshot(tablet, baseTablet)); + } + } else { + tPartition.addToIndexes(new TOlapTableIndexTablets(index.getId(), Lists.newArrayList( + indexTablets.stream().map(Tablet::getId).collect(Collectors.toList())))); + tPartition.setNumBuckets(indexTablets.size()); + for (Tablet tablet : indexTablets) { + partitionTabletSnapshot.add(new TabletLocationSnapshot(tablet, null)); + } + } } tPartition.setIsMutable(partitionInfo.getIsMutable(partitionId)); boolean randomDistribution = diff --git a/fe/fe-core/src/test/java/org/apache/doris/service/FrontendServiceImplTest.java b/fe/fe-core/src/test/java/org/apache/doris/service/FrontendServiceImplTest.java index 6c97272f5cb214..3114b87e7e39b8 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/service/FrontendServiceImplTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/service/FrontendServiceImplTest.java @@ -20,9 +20,11 @@ import org.apache.doris.analysis.UserIdentity; import org.apache.doris.catalog.Database; import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.MaterializedIndex; import org.apache.doris.catalog.OlapTable; import org.apache.doris.catalog.Partition; import org.apache.doris.catalog.TableIf; +import org.apache.doris.catalog.Tablet; import org.apache.doris.common.AuthenticationException; import org.apache.doris.common.Config; import org.apache.doris.common.ErrorCode; @@ -60,6 +62,7 @@ import org.apache.doris.thrift.TShowUserResult; import org.apache.doris.thrift.TStatusCode; import org.apache.doris.thrift.TTableStatus; +import org.apache.doris.thrift.TTabletLocation; import org.apache.doris.transaction.GlobalTransactionMgrIface; import org.apache.doris.transaction.TransactionState; import org.apache.doris.transaction.WriteBlockAllocatingTransaction; @@ -342,6 +345,54 @@ public void testCreatePartitionRange() throws Exception { Assertions.assertNotNull(p20230807); } + @Test + public void testCreatePartitionWithRowBinlog() throws Exception { + String createOlapTblStmt = "CREATE TABLE test.partition_range_with_row_binlog(\n" + + " event_day DATETIME NOT NULL,\n" + + " site_id INT,\n" + + " city_code VARCHAR(100)\n" + + ")\n" + + "DUPLICATE KEY(event_day, site_id, city_code)\n" + + "AUTO PARTITION BY RANGE (date_trunc(event_day, 'day')) ()\n" + + "DISTRIBUTED BY HASH(event_day, site_id) BUCKETS 2\n" + + "PROPERTIES(\"replication_num\" = \"1\", \"binlog.enable\" = \"true\", " + + "\"binlog.format\" = \"ROW\");"; + createTable(createOlapTblStmt); + + Database db = Env.getCurrentInternalCatalog().getDbOrAnalysisException("test"); + OlapTable table = (OlapTable) db.getTableOrAnalysisException("partition_range_with_row_binlog"); + TNullableStringLiteral start = new TNullableStringLiteral(); + start.setValue("2023-08-09 00:00:00"); + TCreatePartitionRequest request = new TCreatePartitionRequest(); + request.setDbId(db.getId()); + request.setTableId(table.getId()); + request.setPartitionValues(Collections.singletonList(Collections.singletonList(start))); + + TCreatePartitionResult result = new FrontendServiceImpl(exeEnv).createPartition(request); + + Assertions.assertEquals(TStatusCode.OK, result.getStatus().getStatusCode()); + Assertions.assertEquals(1, result.getPartitionsSize()); + Assertions.assertEquals(table.getIndexNumber(), result.getPartitions().get(0).getIndexesSize()); + + Partition createdPartition = table.getPartition("p20230809000000"); + Assertions.assertNotNull(createdPartition); + MaterializedIndex rowBinlogIndex = createdPartition + .getMaterializedIndices(MaterializedIndex.IndexExtState.ALL, true).stream() + .filter(MaterializedIndex::isRowBinlog) + .findFirst() + .orElseThrow(); + Assertions.assertEquals(createdPartition.getBaseIndex().getTablets().size() * 2, + result.getTabletsSize()); + for (Tablet rowBinlogTablet : rowBinlogIndex.getTablets()) { + TTabletLocation location = result.getTablets().stream() + .filter(tablet -> tablet.getTabletId() == rowBinlogTablet.getId()) + .findFirst() + .orElseThrow(); + Assertions.assertTrue(location.isSetBaseTabletId()); + Assertions.assertEquals(rowBinlogTablet.getRowBinlogBaseTabletId(), location.getBaseTabletId()); + } + } + @Test public void testCreatePartitionReturnsRetryErrorWhenResultPartitionIsMissing() throws Exception { String createOlapTblStmt = "CREATE TABLE test.partition_dropped_before_result_snapshot(\n" diff --git a/regression-test/data/row_binlog_p0/test_row_binlog_auto_partition.out b/regression-test/data/row_binlog_p0/test_row_binlog_auto_partition.out new file mode 100644 index 00000000000000..e78e7a6bc514ac --- /dev/null +++ b/regression-test/data/row_binlog_p0/test_row_binlog_auto_partition.out @@ -0,0 +1,9 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !auto_partition_data -- +2026-08-10 1 first +2026-08-11 2 second + +-- !auto_partition_binlog -- +0 2026-08-10 1 first +0 2026-08-11 2 second + diff --git a/regression-test/suites/row_binlog_p0/test_row_binlog_auto_partition.groovy b/regression-test/suites/row_binlog_p0/test_row_binlog_auto_partition.groovy new file mode 100644 index 00000000000000..c11030e15b7d59 --- /dev/null +++ b/regression-test/suites/row_binlog_p0/test_row_binlog_auto_partition.groovy @@ -0,0 +1,54 @@ +// 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. + +suite("test_row_binlog_auto_partition", "nonConcurrent") { + sql "DROP TABLE IF EXISTS test_row_binlog_auto_partition FORCE" + + sql """ + CREATE TABLE test_row_binlog_auto_partition ( + event_date DATE NOT NULL, + id INT, + value STRING + ) + DUPLICATE KEY(event_date, id) + AUTO PARTITION BY RANGE (date_trunc(event_date, 'day')) () + DISTRIBUTED BY HASH(id) BUCKETS 1 + PROPERTIES ( + "replication_num" = "1", + "binlog.enable" = "true", + "binlog.format" = "ROW" + ) + """ + + sql """ + INSERT INTO test_row_binlog_auto_partition VALUES + ('2026-08-10', 1, 'first'), + ('2026-08-11', 2, 'second') + """ + + order_qt_auto_partition_data """ + SELECT event_date, id, value + FROM test_row_binlog_auto_partition + ORDER BY event_date, id + """ + + order_qt_auto_partition_binlog """ + SELECT __DORIS_BINLOG_OP__, event_date, id, value + FROM binlog("table" = "test_row_binlog_auto_partition") + ORDER BY event_date, id + """ +}