This is an automated email from the ASF dual-hosted git repository.
yujun777 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new e885b4634df [refactor](ivm) Rename IvmInfo.refreshVersion to
sequencePrefix (#68336)
e885b4634df is described below
commit e885b4634dfa2348bdcc288d2e83e443a13ea9bd
Author: yujun <[email protected]>
AuthorDate: Tue Sep 22 14:11:22 2026 +0800
[refactor](ivm) Rename IvmInfo.refreshVersion to sequencePrefix (#68336)
`IvmInfo.refreshVersion` 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.
The name reads like an epoch or like the MV's own version, which is what
the per-partition `refreshEpoch` being added alongside it is not, and
the two sitting in the same code base is a trap for the next reader.
This renames it to `sequencePrefix`: the prefix of the `(sequence
prefix, delta index, op)` triple that `IvmSequenceCalculator` encodes
into the sequence column. Nothing about the value changes -- it still
counts committing IVM transactions and still prefixes the sequence
column.
- field and accessors: `sequencePrefix`, `getSequencePrefix()`,
`advanceSequencePrefix()`
- `MTMV.getNextRefreshVersion()` -> `MTMV.getNextSequencePrefix()`
- `IvmSequenceCalculator` identifiers, including
`LARGEINT_SEQUENCE_PREFIX_SHIFT` and the range-check messages
- the persisted name changes outright from `"rv"` to `"sp"`: IVM is not
released, so there is no image or journal in the wild that writes the
old name
Trace: https://github.com/apache/doris/issues/65418
---
.../main/java/org/apache/doris/catalog/MTMV.java | 4 +--
.../main/java/org/apache/doris/mtmv/ivm/AGENTS.md | 2 +-
.../doris/mtmv/ivm/IvmDeltaRewriteState.java | 4 +--
.../apache/doris/mtmv/ivm/IvmDeltaRewriter.java | 8 ++---
.../java/org/apache/doris/mtmv/ivm/IvmInfo.java | 15 +++++----
.../doris/mtmv/ivm/IvmSequenceCalculator.java | 38 +++++++++++-----------
.../doris/transaction/DatabaseTransactionMgr.java | 6 ++--
.../doris/mtmv/ivm/IvmAggDeltaHandlerTest.java | 8 ++---
.../doris/mtmv/ivm/IvmDeltaRewriteStateTest.java | 4 +--
.../org/apache/doris/mtmv/ivm/IvmInfoTest.java | 28 +++++++++++-----
.../doris/mtmv/ivm/IvmSequenceCalculatorTest.java | 14 ++++----
.../transaction/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 53efb6b4c62..4050426bab6 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 class MTMV extends OlapTable {
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 81c2f35a0e9..41f75dc6f26 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 50ca07d0a39..a6d493b51fa 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<OlapTable, OlapTableStream> streams,
- boolean includeExhaustedStreams, long refreshVersion, DataType
sequenceType,
+ boolean includeExhaustedStreams, long sequencePrefix, DataType
sequenceType,
Map<OlapTable, List<Long>> 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 1d97e87796c..63918dad4f9 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 class IvmDeltaRewriter {
rewriteContext.isIncludeExhaustedStreams());
Pair<Plan, List<LogicalProject<?>>> 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<IvmDeltaRewriteResult> deltaResult = rewriteDelta(rootPlan,
refreshContext, rewriteState);
if (!deltaResult.isPresent()) {
@@ -156,7 +156,7 @@ public class IvmDeltaRewriter {
return IvmDeltaRewriteHelper.INSTANCE.freshPlan(rewritten);
}
- private IvmDeltaRewriteState createDeltaRewriteState(Plan plan,
IvmIncrRefreshContext ctx, long refreshVersion,
+ private IvmDeltaRewriteState createDeltaRewriteState(Plan plan,
IvmIncrRefreshContext ctx, long sequencePrefix,
Map<BaseTableInfo, Set<Long>> scopePartitionIds) {
Map<OlapTable, OlapTableStream> streams = new HashMap<>();
// Window limits apply to every base table in the plan, including
excluded
@@ -186,7 +186,7 @@ public class IvmDeltaRewriter {
}
}
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 1c985e7e736..2708810770e 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 class IvmInfo {
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 class IvmInfo {
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 74d3549ed3a..b1eb1d1b7cb 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 @@ import java.math.BigInteger;
/**
* Encodes IVM delta ordering into the MTMV sequence column.
*
- * <p>BIGINT encodes {@code (refresh version, delta index, op)}. LARGEINT
additionally encodes
- * a 64-bit binlog sequence: {@code (refresh version, delta index, binlog
sequence, op)}.
+ * <p>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 @@ abstract class IvmSequenceCalculator {
}
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 @@ abstract class IvmSequenceCalculator {
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 @@ abstract class IvmSequenceCalculator {
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 43b6a073c56..4dfb0b075d2 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 @@ public class DatabaseTransactionMgr {
// 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 44b2597e3c3..579af9f450a 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<? extends Plan> 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 @@ class IvmAggDeltaHandlerTest extends IvmDeltaTestBase {
PlanBundle bundle = normalizeAggPlan(agg);
bundle.rewriteResult.setIdentityKeySlots(identityKeys);
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);
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriteStateTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriteStateTest.java
index 35aa2de91ca..990b660cde5 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriteStateTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmDeltaRewriteStateTest.java
@@ -43,7 +43,7 @@ import java.util.Optional;
class IvmDeltaRewriteStateTest extends IvmDeltaTestBase {
@Test
- void testSequenceEncodesRefreshVersionAndDeltaIndex() {
+ void testSequenceEncodesSequencePrefixAndDeltaIndex() {
IvmDeltaRewriteState state = new IvmDeltaRewriteState(
ImmutableMap.of(), false, 7L, BigIntType.INSTANCE,
ImmutableMap.of());
@@ -172,7 +172,7 @@ class IvmDeltaRewriteStateTest extends IvmDeltaTestBase {
}
@Test
- void testLargeIntSequenceEncodesRefreshVersionAndDeltaIndex() {
+ void testLargeIntSequenceEncodesSequencePrefixAndDeltaIndex() {
IvmDeltaRewriteState state = new IvmDeltaRewriteState(
ImmutableMap.of(), false, 7L, LargeIntType.INSTANCE,
ImmutableMap.of());
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmInfoTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmInfoTest.java
index 20bbd3f6c86..fdd1c899bce 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmInfoTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmInfoTest.java
@@ -26,21 +26,31 @@ import java.util.Collections;
class IvmInfoTest {
@Test
- void testRefreshVersionAdvancesAfterCommittedRefresh() {
+ void testSequencePrefixAdvancesAfterCommittedRefresh() {
IvmInfo info = new IvmInfo();
- Assertions.assertEquals(0, info.getRefreshVersion());
- info.advanceRefreshVersion();
- Assertions.assertEquals(1, info.getRefreshVersion());
+ Assertions.assertEquals(0, info.getSequencePrefix());
+ info.advanceSequencePrefix();
+ Assertions.assertEquals(1, info.getSequencePrefix());
}
@Test
- void testRefreshVersionPersistsThroughGson() {
+ void testSequencePrefixPersistsThroughGson() {
IvmInfo info = new IvmInfo();
- info.advanceRefreshVersion();
+ info.advanceSequencePrefix();
IvmInfo recovered =
GsonUtils.GSON.fromJson(GsonUtils.GSON.toJson(info), IvmInfo.class);
- Assertions.assertEquals(1, recovered.getRefreshVersion());
+ Assertions.assertEquals(1, recovered.getSequencePrefix());
+ }
+
+ @Test
+ void testSequencePrefixIsPersistedAsSp() {
+ IvmInfo info = new IvmInfo();
+ info.advanceSequencePrefix();
+
+ String json = GsonUtils.GSON.toJson(info);
+ Assertions.assertTrue(json.contains("\"sp\":1"), json);
+ Assertions.assertFalse(json.contains("\"rv\""), json);
}
@Test
@@ -59,7 +69,7 @@ class IvmInfoTest {
info.requireCompleteBaselineRebuild();
info.setUseFullKeys(true);
info.setPlanSignature("abc123");
- info.advanceRefreshVersion();
+ info.advanceSequencePrefix();
IvmInfo copy = new IvmInfo(info);
info.clearBaselineRebuild();
@@ -68,7 +78,7 @@ class IvmInfoTest {
Assertions.assertTrue(copy.isBaselineRebuildRequired());
Assertions.assertTrue(copy.isUseFullKeys());
Assertions.assertEquals("abc123", copy.getPlanSignature());
- Assertions.assertEquals(1, copy.getRefreshVersion());
+ Assertions.assertEquals(1, copy.getSequencePrefix());
}
@Test
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmSequenceCalculatorTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmSequenceCalculatorTest.java
index 12557c9e86e..13a5ee01e1a 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmSequenceCalculatorTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmSequenceCalculatorTest.java
@@ -57,9 +57,9 @@ class IvmSequenceCalculatorTest {
@Test
void testLargeIntSequenceUsesFullEncodingRange() {
- long maxRefreshVersion = Long.MAX_VALUE >>> 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 @@ class IvmSequenceCalculatorTest {
@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 cfe9d6658d7..5e336047e99 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 class DatabaseTransactionMgrTest {
}
@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 class DatabaseTransactionMgrTest {
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 class DatabaseTransactionMgrTest {
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());
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]