From 73d08f9ee22ba4163e3165249d84c92022238a18 Mon Sep 17 00:00:00 2001 From: yujun Date: Mon, 21 Sep 2026 18:18:35 +0800 Subject: [PATCH] [refactor](ivm) Rename IvmInfo.refreshVersion to sequencePrefix The field is the high part of the sequence values an IVM MV's rows are stamped with -- the low part next to it is a delta index -- but "refreshVersion" reads like an epoch or like the MV's own version, which is exactly what the per-partition refreshEpoch being added alongside it is not. Nothing about the value changes: it still counts committing IVM transactions and still prefixes the sequence column. Key changes: - Rename the field and its accessors to sequencePrefix / getSequencePrefix / advanceSequencePrefix, and MTMV.getNextRefreshVersion to getNextSequencePrefix - Rename the identifiers in IvmSequenceCalculator to match, including LARGEINT_SEQUENCE_PREFIX_SHIFT and the range-check messages - Persist it as "sp" instead of "rv": IVM is not released, so there is no image or journal that writes the old name Unit Test: - IvmInfoTest.testSequencePrefixIsPersistedAsSp pins the persisted name and that "rv" is gone - IvmInfoTest / IvmSequenceCalculatorTest / IvmAggDeltaHandlerTest / IvmDeltaRewriteStateTest / DatabaseTransactionMgrTest cover the renamed API --- .../java/org/apache/doris/catalog/MTMV.java | 4 +- .../java/org/apache/doris/mtmv/ivm/AGENTS.md | 2 +- .../doris/mtmv/ivm/IvmDeltaRewriteState.java | 4 +- .../doris/mtmv/ivm/IvmDeltaRewriter.java | 8 ++-- .../org/apache/doris/mtmv/ivm/IvmInfo.java | 15 ++++---- .../doris/mtmv/ivm/IvmSequenceCalculator.java | 38 +++++++++---------- .../transaction/DatabaseTransactionMgr.java | 6 +-- .../mtmv/ivm/IvmAggDeltaHandlerTest.java | 8 ++-- .../mtmv/ivm/IvmDeltaRewriteStateTest.java | 4 +- .../apache/doris/mtmv/ivm/IvmInfoTest.java | 28 +++++++++----- .../mtmv/ivm/IvmSequenceCalculatorTest.java | 14 +++---- .../DatabaseTransactionMgrTest.java | 6 +-- 12 files changed, 74 insertions(+), 63 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java index 53efb6b4c62508..4050426bab621e 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java @@ -255,8 +255,8 @@ public boolean isIvm() { return getIvmInfo().isEnableIvm(); } - public long getNextRefreshVersion() { - return Config.isCloudMode() ? getNextVersion() : getIvmInfo().getRefreshVersion() + 1; + public long getNextSequencePrefix() { + return Config.isCloudMode() ? getNextVersion() : getIvmInfo().getSequencePrefix() + 1; } public boolean addTaskResult(AlterMTMV alterMTMV, boolean isReplay) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/AGENTS.md b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/AGENTS.md index 81c2f35a0e9b31..41f75dc6f26e19 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/AGENTS.md +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/AGENTS.md @@ -155,7 +155,7 @@ __DORIS_IVM_ROW_ID_COL__ | k1 | cnt | sum_v1 | __DORIS_IVM_DML_FACTOR_COL__ | __ ### Semantics -- **Read-only.** No insert transaction is built. Stream offsets, refresh version, and MV metadata +- **Read-only.** No insert transaction is built. Stream offsets, sequence prefix, and MV metadata are never modified. - **Idempotent.** Repeating the same dry run returns identical rows as long as the base table data has not changed. diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriteState.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriteState.java index 50ca07d0a395fd..a6d493b51fa506 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriteState.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriteState.java @@ -57,11 +57,11 @@ class IvmDeltaRewriteState { private int nextDeltaScanIndex; IvmDeltaRewriteState(Map streams, - boolean includeExhaustedStreams, long refreshVersion, DataType sequenceType, + boolean includeExhaustedStreams, long sequencePrefix, DataType sequenceType, Map> windowPartitionIdsByTable) { this.streams = new HashMap<>(streams); this.includeExhaustedStreams = includeExhaustedStreams; - this.sequenceCalculator = IvmSequenceCalculator.create(refreshVersion, sequenceType); + this.sequenceCalculator = IvmSequenceCalculator.create(sequencePrefix, sequenceType); this.windowPartitionIdsByTable = windowPartitionIdsByTable; } diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriter.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriter.java index 1d97e87796cf97..63918dad4f934b 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriter.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriter.java @@ -67,8 +67,8 @@ public Plan generateIncrRefreshPlan(Plan sinkChild, IvmRewriteResult rewriteResu rewriteContext.isIncludeExhaustedStreams()); Pair>> prefixChain = helper.detachAdaptProjectChain(sinkChild); Plan rootPlan = prefixChain.first; - long refreshVersion = refreshContext.getMtmv().getNextRefreshVersion(); - IvmDeltaRewriteState rewriteState = createDeltaRewriteState(rootPlan, refreshContext, refreshVersion, + long sequencePrefix = refreshContext.getMtmv().getNextSequencePrefix(); + IvmDeltaRewriteState rewriteState = createDeltaRewriteState(rootPlan, refreshContext, sequencePrefix, rewriteContext.getIncrementalScopePartitionIds()); Optional deltaResult = rewriteDelta(rootPlan, refreshContext, rewriteState); if (!deltaResult.isPresent()) { @@ -156,7 +156,7 @@ private static Pair> rewriteSnapshot(Plan plan, IvmDeltaRe return IvmDeltaRewriteHelper.INSTANCE.freshPlan(rewritten); } - private IvmDeltaRewriteState createDeltaRewriteState(Plan plan, IvmIncrRefreshContext ctx, long refreshVersion, + private IvmDeltaRewriteState createDeltaRewriteState(Plan plan, IvmIncrRefreshContext ctx, long sequencePrefix, Map> scopePartitionIds) { Map streams = new HashMap<>(); // Window limits apply to every base table in the plan, including excluded @@ -186,7 +186,7 @@ private IvmDeltaRewriteState createDeltaRewriteState(Plan plan, IvmIncrRefreshCo } } applyScopePartitionIds(windowPartitionIdsByTable, planTables, scopePartitionIds); - return new IvmDeltaRewriteState(streams, ctx.isIncludeExhaustedStreams(), refreshVersion, + return new IvmDeltaRewriteState(streams, ctx.isIncludeExhaustedStreams(), sequencePrefix, DataType.fromCatalogType(ctx.getMtmv().getColumn(Column.SEQUENCE_COL).getType()), windowPartitionIdsByTable); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmInfo.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmInfo.java index 1c985e7e736962..2708810770e1a6 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmInfo.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmInfo.java @@ -53,8 +53,9 @@ public class IvmInfo { @SerializedName("ps") private String planSignature; - @SerializedName("rv") - private long refreshVersion; + /** The prefix of the sequence values this MV's rows are stamped with; see IvmSequenceCalculator. */ + @SerializedName("sp") + private long sequencePrefix; public IvmInfo() { } @@ -65,7 +66,7 @@ public IvmInfo(IvmInfo other) { this.pendingBaselineRebuildPartitions = new HashSet<>(other.pendingBaselineRebuildPartitions); this.useFullKeys = other.useFullKeys; this.planSignature = other.planSignature; - this.refreshVersion = other.refreshVersion; + this.sequencePrefix = other.sequencePrefix; } public boolean isEnableIvm() { @@ -121,12 +122,12 @@ public void setPlanSignature(String planSignature) { this.planSignature = planSignature; } - public long getRefreshVersion() { - return refreshVersion; + public long getSequencePrefix() { + return sequencePrefix; } - public void advanceRefreshVersion() { - refreshVersion++; + public void advanceSequencePrefix() { + sequencePrefix++; } @Override diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmSequenceCalculator.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmSequenceCalculator.java index 74d3549ed3a210..b1eb1d1b7cbf71 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmSequenceCalculator.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmSequenceCalculator.java @@ -33,32 +33,32 @@ /** * Encodes IVM delta ordering into the MTMV sequence column. * - *

BIGINT encodes {@code (refresh version, delta index, op)}. LARGEINT additionally encodes - * a 64-bit binlog sequence: {@code (refresh version, delta index, binlog sequence, op)}. + *

BIGINT encodes {@code (sequence prefix, delta index, op)}. LARGEINT additionally encodes + * a 64-bit binlog sequence: {@code (sequence prefix, delta index, binlog sequence, op)}. */ abstract class IvmSequenceCalculator { static final int DELTA_INDEX_BITS = 10; private static final int BIGINT_LOW_BITS = DELTA_INDEX_BITS + 1; private static final int LARGEINT_BINLOG_BITS = 64; private static final int LARGEINT_DELTA_INDEX_SHIFT = LARGEINT_BINLOG_BITS + 1; - private static final int LARGEINT_REFRESH_VERSION_SHIFT = + private static final int LARGEINT_SEQUENCE_PREFIX_SHIFT = LARGEINT_DELTA_INDEX_SHIFT + DELTA_INDEX_BITS; private static final int MAX_DELTA_INDEX = 1 << DELTA_INDEX_BITS; private static final BigInteger MAX_BINLOG_SEQUENCE = BigInteger.ONE.shiftLeft(LARGEINT_BINLOG_BITS).subtract(BigInteger.ONE); - final long refreshVersion; + final long sequencePrefix; - private IvmSequenceCalculator(long refreshVersion) { - this.refreshVersion = refreshVersion; + private IvmSequenceCalculator(long sequencePrefix) { + this.sequencePrefix = sequencePrefix; } - static IvmSequenceCalculator create(long refreshVersion, DataType sequenceType) { + static IvmSequenceCalculator create(long sequencePrefix, DataType sequenceType) { if (sequenceType.equals(BigIntType.INSTANCE)) { - return new BigIntSequenceCalculator(refreshVersion); + return new BigIntSequenceCalculator(sequencePrefix); } if (sequenceType.equals(LargeIntType.INSTANCE)) { - return new LargeIntSequenceCalculator(refreshVersion); + return new LargeIntSequenceCalculator(sequencePrefix); } throw new IvmException(IvmFailureReason.PLAN_REWRITE_FAILED, "unsupported IVM sequence type: " + sequenceType.simpleString()); @@ -87,11 +87,11 @@ final void checkBinlogSequence(BigInteger binlogSequence) { } private static class BigIntSequenceCalculator extends IvmSequenceCalculator { - private BigIntSequenceCalculator(long refreshVersion) { - super(refreshVersion); - if (refreshVersion < 0 || refreshVersion > (Long.MAX_VALUE >>> BIGINT_LOW_BITS)) { + private BigIntSequenceCalculator(long sequencePrefix) { + super(sequencePrefix); + if (sequencePrefix < 0 || sequencePrefix > (Long.MAX_VALUE >>> BIGINT_LOW_BITS)) { throw new IvmException(IvmFailureReason.PLAN_REWRITE_FAILED, - "IVM refresh version exceeds the BIGINT sequence encoding range: " + refreshVersion); + "IVM sequence prefix exceeds the BIGINT sequence encoding range: " + sequencePrefix); } } @@ -102,18 +102,18 @@ Literal encode(int deltaIndex, BigInteger binlogSequence, boolean positive) { throw new IvmException(IvmFailureReason.PLAN_REWRITE_FAILED, "BIGINT IVM sequence does not support binlog sequence"); } - long sequence = (refreshVersion << BIGINT_LOW_BITS) + long sequence = (sequencePrefix << BIGINT_LOW_BITS) | ((long) deltaIndex << 1) | (positive ? 1 : 0); return new BigIntLiteral(sequence); } } private static class LargeIntSequenceCalculator extends IvmSequenceCalculator { - private LargeIntSequenceCalculator(long refreshVersion) { - super(refreshVersion); - if (refreshVersion < 0 || refreshVersion > (Long.MAX_VALUE >>> BIGINT_LOW_BITS)) { + private LargeIntSequenceCalculator(long sequencePrefix) { + super(sequencePrefix); + if (sequencePrefix < 0 || sequencePrefix > (Long.MAX_VALUE >>> BIGINT_LOW_BITS)) { throw new IvmException(IvmFailureReason.PLAN_REWRITE_FAILED, - "IVM refresh version exceeds the LARGEINT sequence encoding range: " + refreshVersion); + "IVM sequence prefix exceeds the LARGEINT sequence encoding range: " + sequencePrefix); } } @@ -121,7 +121,7 @@ private LargeIntSequenceCalculator(long refreshVersion) { Literal encode(int deltaIndex, BigInteger binlogSequence, boolean positive) { checkDeltaIndex(deltaIndex); checkBinlogSequence(binlogSequence); - BigInteger sequence = BigInteger.valueOf(refreshVersion).shiftLeft(LARGEINT_REFRESH_VERSION_SHIFT) + BigInteger sequence = BigInteger.valueOf(sequencePrefix).shiftLeft(LARGEINT_SEQUENCE_PREFIX_SHIFT) .or(BigInteger.valueOf(deltaIndex).shiftLeft(LARGEINT_DELTA_INDEX_SHIFT)) .or(binlogSequence.shiftLeft(1)); if (positive) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/transaction/DatabaseTransactionMgr.java b/fe/fe-core/src/main/java/org/apache/doris/transaction/DatabaseTransactionMgr.java index 43b6a073c5694e..4dfb0b075d2da0 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/transaction/DatabaseTransactionMgr.java +++ b/fe/fe-core/src/main/java/org/apache/doris/transaction/DatabaseTransactionMgr.java @@ -2418,15 +2418,15 @@ private void updateCatalogAfterCommitted(TransactionState transactionState, Data // update table stream offset if necessary if (!CollectionUtils.isEmpty(transactionState.getStreamUpdateInfos())) { updateStreamOffset(transactionState, transactionState.getCommitTime()); - updateIvmRefreshVersion(transactionState, db); + updateIvmSequencePrefix(transactionState, db); } } - private void updateIvmRefreshVersion(TransactionState transactionState, Database db) { + private void updateIvmSequencePrefix(TransactionState transactionState, Database db) { for (Long tableId : transactionState.getTableIdList()) { Table table = db.getTableNullable(tableId); if (table instanceof MTMV && ((MTMV) table).isIvm()) { - ((MTMV) table).getIvmInfo().advanceRefreshVersion(); + ((MTMV) table).getIvmInfo().advanceSequencePrefix(); } } } diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmAggDeltaHandlerTest.java b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmAggDeltaHandlerTest.java index 44b2597e3c3167..579af9f450a10b 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmAggDeltaHandlerTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmAggDeltaHandlerTest.java @@ -70,8 +70,8 @@ class IvmAggDeltaHandlerTest extends IvmDeltaTestBase { private AggRewriteResult rewriteAgg(LogicalAggregate agg) { PlanBundle bundle = normalizeAggPlan(agg); MTMV mtmv = buildMtmvFromPlan(bundle.normalizedPlan.getOutput()); - mtmv.getIvmInfo().advanceRefreshVersion(); - mtmv.getIvmInfo().advanceRefreshVersion(); + mtmv.getIvmInfo().advanceSequencePrefix(); + mtmv.getIvmInfo().advanceSequencePrefix(); Plan rewritten = new IvmDeltaRewriter().generateIncrRefreshPlan( bundle.normalizedPlan, bundle.rewriteResult, IvmRewriteContext.incremental(mtmv), bundle.connectContext); @@ -89,8 +89,8 @@ private AggRewriteResult rewriteAggWithIdentityKeys(LogicalAggregate>> 11; + long maxSequencePrefix = Long.MAX_VALUE >>> 11; IvmSequenceCalculator calculator = IvmSequenceCalculator.create( - maxRefreshVersion, LargeIntType.INSTANCE); + maxSequencePrefix, LargeIntType.INSTANCE); LargeIntLiteral sequence = (LargeIntLiteral) calculator.encode( 1023, BigInteger.ONE.shiftLeft(64).subtract(BigInteger.ONE), true); @@ -70,12 +70,12 @@ void testLargeIntSequenceUsesFullEncodingRange() { @Test void testSequenceRejectsValuesOutsideEncodingRanges() { - long maxRefreshVersion = Long.MAX_VALUE >>> 11; - IvmSequenceCalculator.create(maxRefreshVersion, LargeIntType.INSTANCE); + long maxSequencePrefix = Long.MAX_VALUE >>> 11; + IvmSequenceCalculator.create(maxSequencePrefix, LargeIntType.INSTANCE); - IvmException refreshVersionException = Assertions.assertThrows(IvmException.class, - () -> IvmSequenceCalculator.create(maxRefreshVersion + 1, LargeIntType.INSTANCE)); - Assertions.assertTrue(refreshVersionException.getMessage().contains("refresh version")); + IvmException sequencePrefixException = Assertions.assertThrows(IvmException.class, + () -> IvmSequenceCalculator.create(maxSequencePrefix + 1, LargeIntType.INSTANCE)); + Assertions.assertTrue(sequencePrefixException.getMessage().contains("sequence prefix")); IvmSequenceCalculator calculator = IvmSequenceCalculator.create(1, LargeIntType.INSTANCE); IvmException deltaIndexException = Assertions.assertThrows(IvmException.class, diff --git a/fe/fe-core/src/test/java/org/apache/doris/transaction/DatabaseTransactionMgrTest.java b/fe/fe-core/src/test/java/org/apache/doris/transaction/DatabaseTransactionMgrTest.java index cfe9d6658d746d..5e336047e99f0e 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/transaction/DatabaseTransactionMgrTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/transaction/DatabaseTransactionMgrTest.java @@ -503,7 +503,7 @@ public void testAbortTransactionWithNotFoundException() throws UserException { } @Test - public void testUpdateCatalogAfterCommittedAdvancesIvmRefreshVersionForNormalCommitAndReplay() + public void testUpdateCatalogAfterCommittedAdvancesIvmSequencePrefixForNormalCommitAndReplay() throws Exception { DatabaseTransactionMgr masterDbTransMgr = masterTransMgr.getDatabaseTransactionMgr(CatalogTestUtil.testDbId1); Database masterDb = masterEnv.getInternalCatalog().getDbOrMetaException(CatalogTestUtil.testDbId1); @@ -520,7 +520,7 @@ public void testUpdateCatalogAfterCommittedAdvancesIvmRefreshVersionForNormalCom method.setAccessible(true); method.invoke(masterDbTransMgr, normalCommitTxn, masterDb, false); - Assertions.assertEquals(1L, normalCommitIvmInfo.getRefreshVersion()); + Assertions.assertEquals(1L, normalCommitIvmInfo.getSequencePrefix()); Mockito.verify(normalCommitStream).unprotectedUpdateStreamUpdate( normalCommitTxn.getStreamUpdateInfos().get(0).getUpdate(), normalCommitTxn.getCommitTime()); @@ -535,7 +535,7 @@ public void testUpdateCatalogAfterCommittedAdvancesIvmRefreshVersionForNormalCom slaveDb.getId(), 9001L, 10002L, replayStreamId, 456L); method.invoke(slaveDbTransMgr, replayTxn, slaveDb, true); - Assertions.assertEquals(1L, replayIvmInfo.getRefreshVersion()); + Assertions.assertEquals(1L, replayIvmInfo.getSequencePrefix()); Mockito.verify(replayStream).unprotectedUpdateStreamUpdate( replayTxn.getStreamUpdateInfos().get(0).getUpdate(), replayTxn.getCommitTime()); }