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
Original file line number Diff line number Diff line change
Expand Up @@ -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<Long, Long> getBinlogColocatedReplicaBackendPathMap(Tablet baseTablet, Tablet rowBinlogTablet)
public static Multimap<Long, Long> getBinlogColocatedReplicaBackendPathMap(
Tablet baseTablet, Tablet rowBinlogTablet)
throws UserException {
Multimap<Long, Long> baseBePathsMap = baseTablet.getNormalReplicaBackendPathMap();
Multimap<Long, Long> binlogBePathsMap = rowBinlogTablet.getNormalReplicaBackendPathMap();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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();
}
Expand All @@ -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())));
}

Expand Down Expand Up @@ -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
Expand All @@ -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();
}
Expand All @@ -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())));
}

Expand Down Expand Up @@ -5214,12 +5222,13 @@ private static final class PartitionResultSnapshot {
private final Partition partition;
private final long partitionId;
private final TOlapTablePartition tPartition;
private final List<Tablet> tablets;
private final List<TabletLocationSnapshot> tablets;
private final int quorum;
private final boolean cacheLoadTabletIdx;

private PartitionResultSnapshot(Partition partition, long partitionId,
TOlapTablePartition tPartition, List<Tablet> tablets, int quorum, boolean cacheLoadTabletIdx) {
TOlapTablePartition tPartition, List<TabletLocationSnapshot> tablets, int quorum,
boolean cacheLoadTabletIdx) {
this.partition = partition;
this.partitionId = partitionId;
this.tPartition = tPartition;
Expand All @@ -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<Long> nodeIds) {
TTabletLocation location = new TTabletLocation(tablet.getId(), nodeIds);
if (isRowBinlog()) {
location.setBaseTabletId(rowBinlogBaseTablet.getId());
}
return location;
}
}

private static List<PartitionResultSnapshot> snapshotPartitionResultsByName(OlapTable olapTable,
Collection<String> partitionNames, boolean loadToSingleTablet, boolean enableAdaptiveRandomBucket,
String resultName) throws UserException {
Expand Down Expand Up @@ -5281,13 +5312,26 @@ private static PartitionResultSnapshot snapshotPartitionResult(PartitionInfo par
TOlapTablePartition tPartition = new TOlapTablePartition();
tPartition.setId(partitionId);
OlapTableSink.setPartitionKeys(tPartition, partitionItem, partColNum);
List<Tablet> partitionTabletSnapshot = new ArrayList<>();
List<TabletLocationSnapshot> partitionTabletSnapshot = new ArrayList<>();
for (MaterializedIndex index : partition.getMaterializedIndices(MaterializedIndex.IndexExtState.ALL, true)) {
List<Tablet> 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 =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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"
Expand Down
Original file line number Diff line number Diff line change
@@ -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

Original file line number Diff line number Diff line change
@@ -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
"""
}
Loading