This is an automated email from the ASF dual-hosted git repository.
morningman pushed a commit to branch branch-incremental-computation
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-incremental-computation
by this push:
new c3f92f5b864 [chore](fe) Pick merged incremental computation PRs
(#68491)
c3f92f5b864 is described below
commit c3f92f5b864382a810d07f02b40525929f70cb23
Author: Mingyu Chen (Rayner) <[email protected]>
AuthorDate: Thu Sep 24 17:42:18 2026 +0800
[chore](fe) Pick merged incremental computation PRs (#68491)
### What problem does this PR solve?
Issue Number: None (branch pick)
Related PR: #67490, #68141, #66963, #68269
Problem Summary: Pick the four merged `incremental-computation` PRs that
do not have the `incremental-computation-picked` label into
`branch-incremental-computation`, in their master merge order:
1. #67490: prevent a local exchange under a serial parent pipeline.
2. #68141: manage MTMV planning caches through a bounded global LRU
cache and expose cache statistics.
3. #66963: restrict HLL, QUANTILE_STATE, and AGG_STATE columns to
Aggregate Key tables by default, with a temporary compatibility setting.
4. #68269: preserve the origin statement for internal MTMV refreshes and
parse `SET_VAR` hints with the refresh context installed.
Dependency check: #33262, referenced by #68141, is already an ancestor
of this branch. #66664 is an open related row-binlog fix rather than a
code prerequisite for #66963. #68272 is the corresponding branch-4.1 fix
for #68269, not a master prerequisite. No additional PR needed picking.
### Release note
Bring the four fixes and MTMV cache management change above to the
incremental computation branch. In particular, non-Aggregate Key tables
now reject new HLL, QUANTILE_STATE, and AGG_STATE columns by default;
incremental MTMV refreshes and dry runs can parse MV queries with
`SET_VAR` hints.
### Check List (For Author)
- Test:
- `./build.sh --fe` (passed, including Checkstyle with 0 violations)
- Focused `./run-fe-ut.sh --run ...` across 13 changed and related test
classes (356 tests passed, 0 failures/errors/skips)
- SQL regression tests were not run; this worktree has no BE cluster.
- Behavior changed: Yes, as described in the source PRs and release
note.
- Does this need documentation: No; the source PRs do not require a
documentation PR.
---------
Co-authored-by: Mryange <[email protected]>
Co-authored-by: xy720 <[email protected]>
Co-authored-by: Luwei <[email protected]>
Co-authored-by: morrySnow <[email protected]>
---
.../main/java/org/apache/doris/common/Config.java | 41 ++++
.../java/org/apache/doris/common/ConfigTest.java | 26 ++
.../src/main/java/org/apache/doris/DorisFE.java | 1 +
.../main/java/org/apache/doris/catalog/Env.java | 8 +
.../main/java/org/apache/doris/catalog/MTMV.java | 114 +++++----
.../doris/common/proc/MTMVCacheHotProcNode.java | 110 +++++++++
.../apache/doris/common/proc/MTMVCacheProcDir.java | 61 +++++
.../doris/common/proc/MTMVCacheStatProcNode.java | 49 ++++
.../org/apache/doris/common/proc/ProcService.java | 1 +
.../apache/doris/job/extensions/mtmv/MTMVTask.java | 8 +-
.../org/apache/doris/mtmv/MTMVCacheManager.java | 210 ++++++++++++++++
.../java/org/apache/doris/mtmv/MTMVPlanUtil.java | 9 +-
.../doris/mtmv/ivm/IvmIncrRefreshManager.java | 2 +
.../org/apache/doris/nereids/StatementContext.java | 13 +
.../glue/translator/PlanTranslatorContext.java | 14 ++
.../trees/plans/commands/RefreshMTMVCommand.java | 3 +
.../plans/commands/info/ColumnDefinition.java | 24 +-
.../trees/plans/commands/info/CreateTableInfo.java | 4 +-
.../org/apache/doris/planner/AddLocalExchange.java | 1 +
.../java/org/apache/doris/planner/PlanNode.java | 16 +-
.../doris/alter/InternalSchemaAlterTest.java | 10 +
.../org/apache/doris/catalog/CreateTableTest.java | 28 +++
.../CreateTableWithBloomFilterIndexTest.java | 31 ++-
.../apache/doris/mtmv/MTMVCacheManagerTest.java | 268 +++++++++++++++++++++
.../org/apache/doris/mtmv/MTMVPlanUtilTest.java | 16 ++
.../java/org/apache/doris/mtmv/MTMVTaskTest.java | 5 +
.../test/java/org/apache/doris/mtmv/MTMVTest.java | 233 +++++++++++++++++-
.../doris/mtmv/ivm/IvmIncrRefreshManagerTest.java | 27 +++
.../plans/commands/RefreshMTMVCommandTest.java | 18 ++
.../commands/UpdateMvByPartitionCommandTest.java | 9 +-
.../plans/commands/info/ColumnDefinitionTest.java | 93 +++++++
.../data/mtmv_p0/test_mtmv_cache_proc.out | 5 +
.../suites/correctness_p0/test_default_hll.groovy | 6 +-
.../duplicate/storage/test_duplicate_hll.groovy | 4 +
.../storage/test_duplicate_quantile_state.groovy | 4 +
...test_state_types_only_in_aggregate_table.groovy | 113 +++++++++
.../data_model_p0/unique/test_unique_hll.groovy | 4 +
.../unique/test_unique_quantile_state.groovy | 4 +
.../test_remote_doris_unique_table_select.groovy | 4 +
.../mtmv_p0/ivm/test_ivm_refresh_dry_run.groovy | 2 +-
.../suites/mtmv_p0/test_mtmv_cache_proc.groovy | 83 +++++++
.../mv_p0/mv_negative/dup_negative_test.groovy | 4 +
.../mv_p0/mv_negative/mor_negative_test.groovy | 4 +
.../mv_p0/mv_negative/mow_negative_test.groovy | 4 +
.../test_local_shuffle_rqg_bugs.groovy | 35 +++
.../support_type/any_value/any_value.groovy | 6 +-
.../suites/query_p0/join/test_join_on.groovy | 4 +
47 files changed, 1658 insertions(+), 81 deletions(-)
diff --git a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
index 7cb32c1e9a1..52668110bc6 100644
--- a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
+++ b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
@@ -1857,6 +1857,11 @@ public class Config extends ConfigBase {
@ConfField(mutable = true, masterOnly = true)
public static boolean enable_quantile_state_type = true;
+ @ConfField(mutable = true, masterOnly = true, description = "Temporary
compatibility switch that allows HLL, "
+ + "QUANTILE_STATE, and AGG_STATE columns in non-aggregate key
tables. Disabled by default. This switch "
+ + "is intended only for migration and will be removed after the
compatibility transition period.")
+ public static boolean enable_non_aggregate_table_state_types = false;
+
/*---------------------- JOB CONFIG START------------------------*/
/**
* The number of threads used to dispatch timer job.
@@ -2294,6 +2299,42 @@ public class Config extends ConfigBase {
+ "pruning.")
public static int cache_partition_meta_table_manage_num = 100;
+ @ConfField(
+ mutable = true,
+ callback = NonNegativeMtmvCacheNumConfHandler.class,
+ callbackClassString =
"org.apache.doris.mtmv.MTMVCacheManager$UpdateConfig",
+ description = "Max mtmv plan cache entries kept by
MTMVCacheManager. 0 disables the cache, "
+ + "negative values are rejected. Default 3000.")
+ public static int mtmv_cache_manage_num = 3000;
+
+ public static class NonNegativeMtmvCacheNumConfHandler implements
ConfHandler {
+ @Override
+ public void handle(Field field, String value) throws Exception {
+ int parsed = Integer.parseInt(value.trim());
+ if (parsed < 0) {
+ throw new ConfigException(field.getName() + " must not be
negative, 0 disables the cache");
+ }
+ field.setInt(null, parsed);
+ }
+ }
+
+ public static void validateMtmvCacheConfig() throws ConfigException {
+ if (mtmv_cache_manage_num < 0) {
+ throw new ConfigException("mtmv_cache_manage_num must not be
negative, 0 disables the cache");
+ }
+ }
+
+ @ConfField(
+ mutable = true,
+ callbackClassString =
"org.apache.doris.mtmv.MTMVCacheManager$UpdateConfig",
+ description = "Idle expiration in seconds for entries in
MTMVCacheManager. Default 86400.")
+ public static long expire_mtmv_cache_in_fe_second = 86400;
+
+ @ConfField(
+ mutable = true,
+ description = "Row cap for SHOW PROC '/mtmv_cache/hot'. Default
500.")
+ public static int mtmv_cache_hot_show_num = 500;
+
/**
* HBO plan stats. cache number which can be reused for the next query.
*/
diff --git a/fe/fe-common/src/test/java/org/apache/doris/common/ConfigTest.java
b/fe/fe-common/src/test/java/org/apache/doris/common/ConfigTest.java
index 64778f41091..439cbd3fea3 100644
--- a/fe/fe-common/src/test/java/org/apache/doris/common/ConfigTest.java
+++ b/fe/fe-common/src/test/java/org/apache/doris/common/ConfigTest.java
@@ -165,6 +165,32 @@ public class ConfigTest {
}
}
+ @Test
+ public void testMtmvCacheManageNumRejectsNegative() throws Exception {
+ int original = Config.mtmv_cache_manage_num;
+ try {
+ Config.mtmv_cache_manage_num = 100;
+ // ADMIN SET FRONTEND CONFIG runs the annotation callback before
the cache-reload handler,
+ // so a negative maximum must be refused there and leave the field
untouched.
+ ConfigException negative =
Assertions.assertThrows(ConfigException.class,
+ () -> ConfigBase.setMutableConfig("mtmv_cache_manage_num",
"-1"));
+ Assertions.assertTrue(negative.getMessage().contains("must not be
negative"));
+ Assertions.assertEquals(100, Config.mtmv_cache_manage_num);
+
+ // 0 is the documented way to disable the cache.
+ new Config.NonNegativeMtmvCacheNumConfHandler()
+ .handle(ConfigBase.getField("mtmv_cache_manage_num"), " 0
");
+ Assertions.assertEquals(0, Config.mtmv_cache_manage_num);
+ Assertions.assertDoesNotThrow(Config::validateMtmvCacheConfig);
+
+ // fe.conf assigns the field without running any callback, so
startup validates it too.
+ Config.mtmv_cache_manage_num = -1;
+ Assertions.assertThrows(ConfigException.class,
Config::validateMtmvCacheConfig);
+ } finally {
+ Config.mtmv_cache_manage_num = original;
+ }
+ }
+
@Test
public void testValidateWebSqlStartupConfig() throws ConfigException {
int originalIdleTimeout = Config.web_sql_session_idle_timeout_seconds;
diff --git a/fe/fe-core/src/main/java/org/apache/doris/DorisFE.java
b/fe/fe-core/src/main/java/org/apache/doris/DorisFE.java
index 56d90909b13..f341a486ff9 100755
--- a/fe/fe-core/src/main/java/org/apache/doris/DorisFE.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/DorisFE.java
@@ -151,6 +151,7 @@ public class DorisFE {
// Because the path of custom config file is defined in fe.conf
config.initCustom(Config.custom_config_dir + "/fe_custom.conf");
Config.validateWebSqlConfig();
+ Config.validateMtmvCacheConfig();
// inverted_index_storage_format's runtime callback is not invoked
while parsing
// fe.conf/fe_custom.conf, so validate the loaded value here after
both files are loaded
// and merged, to reject a "V1" left over in the config files at
startup.
diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java
b/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java
index 51833705acf..878f1c1ddbc 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java
@@ -155,6 +155,7 @@ import org.apache.doris.meta.MetaContext;
import org.apache.doris.metric.MetricRepo;
import org.apache.doris.mtmv.BaseTableInfo;
import org.apache.doris.mtmv.MTMVAlterOpType;
+import org.apache.doris.mtmv.MTMVCacheManager;
import org.apache.doris.mtmv.MTMVPartitionExprFactory;
import org.apache.doris.mtmv.MTMVPartitionInfo;
import org.apache.doris.mtmv.MTMVPartitionInfo.MTMVPartitionType;
@@ -588,6 +589,8 @@ public class Env {
private final NereidsSortedPartitionsCacheManager
sortedPartitionsCacheManager;
+ private final MTMVCacheManager mtmvCacheManager;
+
private final SplitSourceManager splitSourceManager;
private final GlobalExternalTransactionInfoMgr
globalExternalTransactionInfoMgr;
@@ -884,6 +887,7 @@ public class Env {
this.dnsCache = new DNSCache();
this.sqlCacheManager = new NereidsSqlCacheManager();
this.sortedPartitionsCacheManager = new
NereidsSortedPartitionsCacheManager();
+ this.mtmvCacheManager = new MTMVCacheManager();
this.splitSourceManager = new SplitSourceManager();
this.globalExternalTransactionInfoMgr = new
GlobalExternalTransactionInfoMgr();
this.tokenManager = new TokenManager();
@@ -7637,6 +7641,10 @@ public class Env {
return sqlCacheManager;
}
+ public MTMVCacheManager getMtmvCacheManager() {
+ return mtmvCacheManager;
+ }
+
public NereidsSortedPartitionsCacheManager
getSortedPartitionsCacheManager() {
return sortedPartitionsCacheManager;
}
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 88cc6979940..2727cc0d2a1 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
@@ -33,6 +33,7 @@ import org.apache.doris.mtmv.BaseTableInfo;
import org.apache.doris.mtmv.EnvInfo;
import org.apache.doris.mtmv.MTMVAlterOpType;
import org.apache.doris.mtmv.MTMVCache;
+import org.apache.doris.mtmv.MTMVCacheManager;
import org.apache.doris.mtmv.MTMVJobInfo;
import org.apache.doris.mtmv.MTMVJobManager;
import org.apache.doris.mtmv.MTMVPartitionExpander;
@@ -54,6 +55,7 @@ import org.apache.doris.mtmv.MTMVStatus;
import org.apache.doris.mtmv.MTMVUtil;
import org.apache.doris.mtmv.ivm.IvmInfo;
import org.apache.doris.mtmv.ivm.IvmUtil;
+import org.apache.doris.nereids.StatementContext;
import org.apache.doris.nereids.rules.analysis.SessionVarGuardRewriter;
import
org.apache.doris.nereids.trees.plans.commands.info.RefreshMTMVInfo.RefreshMode;
import org.apache.doris.persist.AlterMTMV;
@@ -123,11 +125,6 @@ public class MTMV extends OlapTable {
*/
@SerializedName("pst")
private Map<String, MTMVPartitionState> partitionStates;
- // Should update after every fresh, not persist
- // Cache with SessionVarGuardExpr: used when query session variables
differ from MV creation variables
- private MTMVCache cacheWithGuard;
- // Cache without SessionVarGuardExpr: used when query session variables
match MV creation variables
- private MTMVCache cacheWithoutGuard;
// Increased every time rewrite cache is invalidated to prevent publishing
stale in-flight cache builds.
private transient long rewriteCacheGeneration;
private long schemaChangeVersion;
@@ -281,8 +278,8 @@ public class MTMV extends OlapTable {
}
try {
// The replay thread may not have initialized the catalog yet
to avoid getting stuck due
- // to connection issues such as S3, so it is directly set to
null
- if (!isReplay) {
+ // to connection issues such as S3, so it is directly set to
null.
+ if (!isReplay &&
Env.getCurrentEnv().getMtmvCacheManager().isEnabled()) {
ConnectContext currentContext = ConnectContext.get();
// shouldn't do this while holding mvWriteLock
// TODO: these two cache compute share something same, can
be simplified in future
@@ -327,12 +324,19 @@ public class MTMV extends OlapTable {
}
ivmInfo.clearBaselineRebuild();
}
+ // The refresh publishes a new plan, so every cache built
before this commit is stale.
+ // Bump before publishing so an in-flight build cannot pass
its generation check later.
+ boolean publishCache = needUpdateCache && cacheGeneration ==
rewriteCacheGeneration && !isDropped;
+ rewriteCacheGeneration++;
if (needUpdateCache) {
- if (cacheGeneration == rewriteCacheGeneration) {
- // Initialize cacheWithGuard, cacheWithoutGuard will
be lazily generated when needed
- this.cacheWithGuard = mtmvCacheWithGuard;
- // Clear the other cache to ensure consistency
- this.cacheWithoutGuard = mtmvCacheWithoutGuard;
+ MTMVCacheManager manager =
Env.getCurrentEnv().getMtmvCacheManager();
+ if (publishCache && mtmvCacheWithGuard != null) {
+ manager.put(this.id, true, mtmvCacheWithGuard);
+ } else {
+ manager.invalidate(this.id);
+ }
+ if (publishCache && mtmvCacheWithoutGuard != null) {
+ manager.put(this.id, false, mtmvCacheWithoutGuard);
}
}
} else {
@@ -543,51 +547,56 @@ public class MTMV extends OlapTable {
*/
public MTMVCache getOrGenerateCache(ConnectContext connectionContext)
throws
org.apache.doris.nereids.exceptions.AnalysisException {
- // store two MTMVCaches: one is a cache where SessionVariables differ
from those at creation time,
- // and the MTMV plan includes a guardexpr;
- // the other is a cache where SessionVariables are the same as at
creation time, and the MTMV plan
- // does not include a guardexpr;
- // This way, when sessionVariables are the same, rewriting is possible;
- // When sessionVariables are different, there are two cases:
- // 1. If a guardexpr is present, rewriting is not possible;
- // 2. If no guardexpr is present, rewriting is possible.
- // Determine if current session variables match MV creation session
variables
Map<String, String> currentSessionVars =
connectionContext.getSessionVariable().getAffectQueryResultInPlanVariables();
boolean sessionVarsMatch =
SessionVarGuardRewriter.checkSessionVariablesMatch(
currentSessionVars, this.sessionVariables);
+ boolean guarded = !sessionVarsMatch;
+ MTMVCacheManager manager = Env.getCurrentEnv().getMtmvCacheManager();
+ StatementContext statementContext =
connectionContext.getStatementContext();
while (true) {
long cacheGeneration;
- // Select appropriate cache based on session variable match
+ MTMVCache cached;
readMvLock();
try {
- MTMVCache cache = getCache(sessionVarsMatch);
- if (cache != null) {
- return cache;
+ cached = manager.isEnabled() ? manager.getIfPresent(this.id,
guarded) : null;
+ if (cached == null && statementContext != null) {
+ cached = statementContext.getQueryLocalMtmvCache(this.id,
guarded);
}
cacheGeneration = rewriteCacheGeneration;
} finally {
readMvUnlock();
}
-
- // Generate cache if not exists
- // Concurrent situations may result in duplicate cache generation,
- // but we tolerate this in order to prevent nested use of readLock
and write MvLock for the table
- MTMVCache mtmvCache = createRewriteCache(connectionContext, false,
!sessionVarsMatch);
- writeMvLock();
+ if (cached != null) {
+ return cached;
+ }
+ MTMVCache generated = createRewriteCache(connectionContext, false,
guarded);
+ readMvLock();
try {
- MTMVCache cache = getCache(sessionVarsMatch);
- if (cache != null) {
- return cache;
- }
if (cacheGeneration != rewriteCacheGeneration) {
+ // Someone invalidated between our snapshot and now; drop
the stale build and retry.
continue;
}
- setCache(sessionVarsMatch, mtmvCache);
- return mtmvCache;
+ if (manager.isEnabled()) {
+ MTMVCache existing = manager.getIfPresent(this.id,
guarded);
+ if (existing != null) {
+ return existing;
+ }
+ if (!isDropped) {
+ manager.put(this.id, guarded, generated);
+ }
+ } else if (statementContext != null && !isDropped) {
+ // Global cache is disabled (maximumSize=0); keep one copy
for this statement only.
+ MTMVCache existing =
statementContext.getQueryLocalMtmvCache(this.id, guarded);
+ if (existing != null) {
+ return existing;
+ }
+ statementContext.putQueryLocalMtmvCache(this.id, guarded,
generated);
+ }
+ return generated;
} finally {
- writeMvUnlock();
+ readMvUnlock();
}
}
}
@@ -1047,8 +1056,7 @@ public class MTMV extends OlapTable {
writeMvLock();
try {
rewriteCacheGeneration++;
- cacheWithGuard = null;
- cacheWithoutGuard = null;
+ Env.getCurrentEnv().getMtmvCacheManager().invalidate(this.id);
} finally {
writeMvUnlock();
}
@@ -1198,18 +1206,6 @@ public class MTMV extends OlapTable {
this.mvRwLock.writeLock().unlock();
}
- private MTMVCache getCache(boolean sessionVarsMatch) {
- return sessionVarsMatch ? cacheWithoutGuard : cacheWithGuard;
- }
-
- private void setCache(boolean sessionVarsMatch, MTMVCache cache) {
- if (sessionVarsMatch) {
- this.cacheWithoutGuard = cache;
- } else {
- this.cacheWithGuard = cache;
- }
- }
-
// toString() is not easy to find where to call the method
public String toInfoString() {
final StringBuilder sb = new StringBuilder("MTMV{");
@@ -1287,6 +1283,20 @@ public class MTMV extends OlapTable {
compatiblePctSnapshot(partitionSnapshots);
}
+ @Override
+ public void markDropped() {
+ super.markDropped();
+ // A refresh or query building a cache outside the MV lock must not
+ // be able to republish it after the drop.
+ writeMvLock();
+ try {
+ rewriteCacheGeneration++;
+ Env.getCurrentEnv().getMtmvCacheManager().invalidate(this.id);
+ } finally {
+ writeMvUnlock();
+ }
+ }
+
private void compatiblePctSnapshot(Map<String,
MTMVRefreshPartitionSnapshot> partitionSnapshots) {
BaseTableInfo relatedTableInfo = mvPartitionInfo.getRelatedTableInfo();
if (relatedTableInfo == null) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/common/proc/MTMVCacheHotProcNode.java
b/fe/fe-core/src/main/java/org/apache/doris/common/proc/MTMVCacheHotProcNode.java
new file mode 100644
index 00000000000..45f0e32f7f5
--- /dev/null
+++
b/fe/fe-core/src/main/java/org/apache/doris/common/proc/MTMVCacheHotProcNode.java
@@ -0,0 +1,110 @@
+// 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.
+
+package org.apache.doris.common.proc;
+
+import org.apache.doris.catalog.Database;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.Table;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.Config;
+import org.apache.doris.datasource.InternalCatalog;
+import org.apache.doris.mtmv.MTMVCacheManager;
+import org.apache.doris.mtmv.MTMVCacheManager.HotEntry;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.Lists;
+
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+public class MTMVCacheHotProcNode implements ProcNodeInterface {
+ public static final ImmutableList<String> TITLE_NAMES = new
ImmutableList.Builder<String>()
+
.add("MtmvId").add("DbName").add("MvName").add("Guarded").add("IdleMs").build();
+
+ private static final String UNKNOWN_DB = "<unknown>";
+ private static final String DROPPED_MV = "<dropped>";
+
+ private final MTMVCacheManager manager;
+
+ public MTMVCacheHotProcNode(MTMVCacheManager manager) {
+ this.manager = manager;
+ }
+
+ @Override
+ public ProcResult fetchResult() throws AnalysisException {
+ BaseProcResult result = new BaseProcResult();
+ result.setNames(TITLE_NAMES);
+ List<HotEntry> entries =
manager.hotEntries(Config.mtmv_cache_hot_show_num);
+ if (entries.isEmpty()) {
+ return result;
+ }
+ Set<Long> wanted = new HashSet<>();
+ for (HotEntry e : entries) {
+ wanted.add(e.mtmvId);
+ }
+ Map<Long, Table> idToTable = resolveTables(wanted);
+ for (HotEntry entry : entries) {
+ String dbName = UNKNOWN_DB;
+ String mvName = DROPPED_MV;
+ Table table = idToTable.get(entry.mtmvId);
+ if (table != null) {
+ mvName = table.getName();
+ String qualified = table.getQualifiedDbName();
+ if (qualified != null && !qualified.isEmpty()) {
+ dbName = qualified;
+ }
+ }
+ result.addRow(Lists.newArrayList(
+ String.valueOf(entry.mtmvId),
+ dbName,
+ mvName,
+ entry.guarded ? "Yes" : "No",
+ String.valueOf(entry.idleMs)));
+ }
+ return result;
+ }
+
+ private static Map<Long, Table> resolveTables(Set<Long> ids) {
+ Map<Long, Table> out = new HashMap<>();
+ if (Env.getCurrentEnv() == null) {
+ return out;
+ }
+ InternalCatalog catalog = Env.getCurrentInternalCatalog();
+ if (catalog == null) {
+ return out;
+ }
+ for (Database db : catalog.getDbs()) {
+ for (Long id : ids) {
+ if (out.containsKey(id)) {
+ continue;
+ }
+ Table t = db.getTableNullable(id);
+ if (t != null) {
+ out.put(id, t);
+ }
+ }
+ if (out.size() == ids.size()) {
+ break;
+ }
+ }
+ return out;
+ }
+}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/common/proc/MTMVCacheProcDir.java
b/fe/fe-core/src/main/java/org/apache/doris/common/proc/MTMVCacheProcDir.java
new file mode 100644
index 00000000000..c1cc208e88f
--- /dev/null
+++
b/fe/fe-core/src/main/java/org/apache/doris/common/proc/MTMVCacheProcDir.java
@@ -0,0 +1,61 @@
+// 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.
+
+package org.apache.doris.common.proc;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.mtmv.MTMVCacheManager;
+
+import com.google.common.base.Strings;
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.Lists;
+
+/** Two-level proc dir for '/mtmv_cache': "stat" and "hot" child nodes. */
+public class MTMVCacheProcDir implements ProcDirInterface {
+ public static final ImmutableList<String> TITLE_NAMES = new
ImmutableList.Builder<String>()
+ .add("Name").add("Info").build();
+
+ @Override
+ public ProcResult fetchResult() throws AnalysisException {
+ BaseProcResult result = new BaseProcResult();
+ result.setNames(TITLE_NAMES);
+ result.addRow(Lists.newArrayList("stat", "Global cache stats"));
+ result.addRow(Lists.newArrayList("hot", "Top hot mtmv cache entries"));
+ return result;
+ }
+
+ @Override
+ public boolean register(String name, ProcNodeInterface node) {
+ return false;
+ }
+
+ @Override
+ public ProcNodeInterface lookup(String name) throws AnalysisException {
+ if (Strings.isNullOrEmpty(name)) {
+ throw new AnalysisException("mtmv_cache child name is empty");
+ }
+ MTMVCacheManager manager = Env.getCurrentEnv().getMtmvCacheManager();
+ if (name.equalsIgnoreCase("stat")) {
+ return new MTMVCacheStatProcNode(manager);
+ }
+ if (name.equalsIgnoreCase("hot")) {
+ return new MTMVCacheHotProcNode(manager);
+ }
+ throw new AnalysisException("unknown mtmv_cache child: " + name);
+ }
+}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/common/proc/MTMVCacheStatProcNode.java
b/fe/fe-core/src/main/java/org/apache/doris/common/proc/MTMVCacheStatProcNode.java
new file mode 100644
index 00000000000..ed735ab71e8
--- /dev/null
+++
b/fe/fe-core/src/main/java/org/apache/doris/common/proc/MTMVCacheStatProcNode.java
@@ -0,0 +1,49 @@
+// 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.
+
+package org.apache.doris.common.proc;
+
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.mtmv.MTMVCacheManager;
+import org.apache.doris.mtmv.MTMVCacheManager.Snapshot;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.Lists;
+
+public class MTMVCacheStatProcNode implements ProcNodeInterface {
+ public static final ImmutableList<String> TITLE_NAMES = new
ImmutableList.Builder<String>()
+ .add("Name").add("Value").build();
+
+ private final MTMVCacheManager manager;
+
+ public MTMVCacheStatProcNode(MTMVCacheManager manager) {
+ this.manager = manager;
+ }
+
+ @Override
+ public ProcResult fetchResult() throws AnalysisException {
+ BaseProcResult result = new BaseProcResult();
+ result.setNames(TITLE_NAMES);
+ Snapshot s = manager.snapshot();
+ result.addRow(Lists.newArrayList("size", String.valueOf(s.size)));
+ result.addRow(Lists.newArrayList("hitCount",
String.valueOf(s.hitCount)));
+ result.addRow(Lists.newArrayList("missCount",
String.valueOf(s.missCount)));
+ result.addRow(Lists.newArrayList("evictionCount",
String.valueOf(s.evictionCount)));
+ result.addRow(Lists.newArrayList("hitRate", String.format("%.4f",
s.hitRate)));
+ return result;
+ }
+}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/common/proc/ProcService.java
b/fe/fe-core/src/main/java/org/apache/doris/common/proc/ProcService.java
index a1f54901bde..63b7c0d96e7 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/common/proc/ProcService.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/common/proc/ProcService.java
@@ -59,6 +59,7 @@ public final class ProcService {
root.register("bdbje", new BDBJEProcDir());
root.register("diagnose", new DiagnoseProcDir());
root.register("binlog", new BinlogProcDir());
+ root.register("mtmv_cache", new MTMVCacheProcDir());
}
// 通过指定的路径获得对应的PROC Node
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java
index 45a613e5e61..b736d7b7372 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java
@@ -78,6 +78,7 @@ import
org.apache.doris.nereids.trees.plans.commands.CreateMTMVCommand;
import
org.apache.doris.nereids.trees.plans.commands.UpdateMvByPartitionCommand;
import
org.apache.doris.nereids.trees.plans.commands.info.RefreshMTMVInfo.RefreshMode;
import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.OriginStatement;
import org.apache.doris.qe.StmtExecutor;
import org.apache.doris.rpc.RpcException;
import org.apache.doris.system.SystemInfoService;
@@ -1139,10 +1140,11 @@ public class MTMVTask extends AbstractTask {
Map<TableIf, String> tableWithPartKey,
Optional<IvmRewriteContext> rewriteContext, RefreshMode
refreshMode)
throws Exception {
- // Create MTMV context first so that new StatementContext() captures
the
- // correct thread-local ConnectContext (with MTMV disabled rules,
etc.).
+ // Create the MTMV context before parsing the MV definition SQL so
SET_VAR hints
+ // resolve against the internal session (with MTMV disabled rules,
etc.).
ConnectContext mtmvCtx = MTMVPlanUtil.createMTMVContext(mtmv,
MTMVPlanUtil.DISABLE_RULES_WHEN_RUN_MTMV_TASK);
- StatementContext statementContext = new StatementContext();
+ StatementContext statementContext = new StatementContext(
+ mtmvCtx, new OriginStatement(mtmv.getQuerySql(), 0));
// Install the StatementContext on the ConnectContext before parsing
// the MV definition SQL. UpdateMvByPartitionCommand.from() calls
// NereidsParser.parseSingle() which, for SQL containing SET_VAR hints,
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVCacheManager.java
b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVCacheManager.java
new file mode 100644
index 00000000000..a6e35938318
--- /dev/null
+++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVCacheManager.java
@@ -0,0 +1,210 @@
+// 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.
+
+package org.apache.doris.mtmv;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.ConfigBase.DefaultConfHandler;
+
+import com.github.benmanes.caffeine.cache.Cache;
+import com.github.benmanes.caffeine.cache.Caffeine;
+import com.github.benmanes.caffeine.cache.stats.CacheStats;
+import com.google.common.annotations.VisibleForTesting;
+
+import java.lang.reflect.Field;
+import java.time.Duration;
+import java.util.Collections;
+import java.util.List;
+import java.util.Objects;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+/**
+ * FE-local cache manager for materialized view cache.
+ */
+public class MTMVCacheManager {
+
+ private final Object swapLock = new Object();
+ private volatile Cache<Key, MTMVCache> caches;
+
+ public MTMVCacheManager() {
+ caches = build(Config.mtmv_cache_manage_num,
Config.expire_mtmv_cache_in_fe_second);
+ }
+
+ public MTMVCache getIfPresent(long mtmvId, boolean guarded) {
+ return caches.getIfPresent(new Key(mtmvId, guarded));
+ }
+
+ public void put(long mtmvId, boolean guarded, MTMVCache cache) {
+ Objects.requireNonNull(cache, "mtmv cache to publish must not be
null");
+ synchronized (swapLock) {
+ caches.put(new Key(mtmvId, guarded), cache);
+ }
+ }
+
+ public void invalidate(long mtmvId) {
+ synchronized (swapLock) {
+ caches.invalidate(new Key(mtmvId, true));
+ caches.invalidate(new Key(mtmvId, false));
+ }
+ }
+
+ public void invalidateAll() {
+ synchronized (swapLock) {
+ caches.invalidateAll();
+ }
+ }
+
+ public long size() {
+ return caches.estimatedSize();
+ }
+
+ /** False when the live maximum is 0, i.e. every put would be discarded
immediately. */
+ public boolean isEnabled() {
+ return caches.policy().eviction().map(eviction ->
eviction.getMaximum() > 0).orElse(true);
+ }
+
+ public Snapshot snapshot() {
+ Cache<Key, MTMVCache> current = caches;
+ CacheStats s = current.stats();
+ return new Snapshot(current.estimatedSize(), s.hitCount(),
s.missCount(),
+ s.evictionCount(), s.hitRate());
+ }
+
+ /**
+ * Snapshot for SHOW PROC '/mtmv_cache/hot'. Ordered by
most-recently-accessed first when
+ * expireAfterAccess is enabled; falls back to iteration order with
idleMs=-1 otherwise.
+ */
+ public List<HotEntry> hotEntries(int limit) {
+ if (limit <= 0) {
+ return Collections.emptyList();
+ }
+ Cache<Key, MTMVCache> current = caches;
+ return current.policy().expireAfterAccess()
+ .map(exp -> exp.youngest(stream -> stream
+ .limit(limit)
+ .map(entry -> {
+ Key k = entry.getKey();
+ long expireMs =
exp.getExpiresAfter(TimeUnit.MILLISECONDS);
+ long idleMs = Math.max(expireMs -
entry.expiresAfter().toMillis(), 0L);
+ return new HotEntry(k.mtmvId, k.guarded, idleMs);
+ })
+ .collect(Collectors.toList())))
+ .orElseGet(() -> current.asMap().keySet().stream()
+ .limit(limit)
+ .map(k -> new HotEntry(k.mtmvId, k.guarded, -1L))
+ .collect(Collectors.toList()));
+ }
+
+ public void updateConfig() {
+ Cache<Key, MTMVCache> fresh = build(Config.mtmv_cache_manage_num,
Config.expire_mtmv_cache_in_fe_second);
+ synchronized (swapLock) {
+ fresh.putAll(caches.asMap());
+ fresh.cleanUp();
+ caches = fresh;
+ }
+ }
+
+ public static synchronized void reloadConfig() {
+ Env env = Env.getCurrentEnv();
+ if (env == null) {
+ return;
+ }
+ env.getMtmvCacheManager().updateConfig();
+ }
+
+ private static Cache<Key, MTMVCache> build(int maxSize, long
expireAfterAccessSeconds) {
+ Caffeine<Object, Object> builder =
Caffeine.newBuilder().softValues().recordStats()
+ .maximumSize(Math.max(maxSize, 0));
+ if (expireAfterAccessSeconds > 0) {
+
builder.expireAfterAccess(Duration.ofSeconds(expireAfterAccessSeconds));
+ }
+ return builder.build();
+ }
+
+ // NOTE: referenced by Config.mtmv_cache_manage_num.callbackClassString and
+ // Config.expire_mtmv_cache_in_fe_second.callbackClassString.
+ public static class UpdateConfig extends DefaultConfHandler {
+ @Override
+ public void handle(Field field, String confVal) throws Exception {
+ super.handle(field, confVal);
+ MTMVCacheManager.reloadConfig();
+ }
+ }
+
+ /** Stable composite key so it is immune to BaseTableInfo hashCode drift.
*/
+ public static final class Key {
+ public final long mtmvId;
+ public final boolean guarded;
+
+ public Key(long mtmvId, boolean guarded) {
+ this.mtmvId = mtmvId;
+ this.guarded = guarded;
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (!(o instanceof Key)) {
+ return false;
+ }
+ Key that = (Key) o;
+ return mtmvId == that.mtmvId && guarded == that.guarded;
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(mtmvId, guarded);
+ }
+ }
+
+ public static final class HotEntry {
+ public final long mtmvId;
+ public final boolean guarded;
+ public final long idleMs;
+
+ public HotEntry(long mtmvId, boolean guarded, long idleMs) {
+ this.mtmvId = mtmvId;
+ this.guarded = guarded;
+ this.idleMs = idleMs;
+ }
+ }
+
+ public static final class Snapshot {
+ public final long size;
+ public final long hitCount;
+ public final long missCount;
+ public final long evictionCount;
+ public final double hitRate;
+
+ public Snapshot(long size, long hitCount, long missCount, long
evictionCount, double hitRate) {
+ this.size = size;
+ this.hitCount = hitCount;
+ this.missCount = missCount;
+ this.evictionCount = evictionCount;
+ this.hitRate = hitRate;
+ }
+ }
+
+ @VisibleForTesting
+ public Cache<Key, MTMVCache> getCachesForTest() {
+ return caches;
+ }
+}
diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPlanUtil.java
b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPlanUtil.java
index 11f50909ccc..df7d7bf99ed 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPlanUtil.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPlanUtil.java
@@ -170,6 +170,8 @@ public class MTMVPlanUtil {
* executing {@link StmtExecutor} through {@code executorConsumer} before
the command
* runs and clearing it (with {@code null}) after the command finishes, so
task
* cancellation can interrupt the running statement.
+ *
+ * <p>The supplied statement context must contain the originating SQL
statement.
*/
public static void executeCommand(ConnectContext ctx, Command command,
StatementContext stmtCtx, @Nullable String auditStmt,
@@ -178,7 +180,10 @@ public class MTMVPlanUtil {
ctx.getState().setNereids(true);
ctx.getSessionVariable().setEnableMaterializedViewRewrite(false);
ctx.getSessionVariable().setEnableDmlMaterializedViewRewrite(false);
- StmtExecutor executor = new StmtExecutor(ctx, new
LogicalPlanAdapter(command, stmtCtx));
+ LogicalPlanAdapter adapter = new LogicalPlanAdapter(command, stmtCtx);
+
adapter.setOrigStmt(Preconditions.checkNotNull(stmtCtx.getOriginStatement(),
+ "MTMV command origin statement must not be null"));
+ StmtExecutor executor = new StmtExecutor(ctx, adapter);
ctx.setExecutor(executor);
ctx.setQueryId(AbstractTask.generateQueryId());
if (executorConsumer != null) {
@@ -696,7 +701,7 @@ public class MTMVPlanUtil {
if (col.getType().isVarBinaryType()) {
throw new AnalysisException("MTMV do not support varbinary
type : " + col.getName());
}
- col.validate(true, keysSet, Sets.newHashSet(),
finalEnableMergeOnWrite, keysType);
+ col.validate(true, keysSet, Sets.newHashSet(),
finalEnableMergeOnWrite, keysType, true);
}
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmIncrRefreshManager.java
b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmIncrRefreshManager.java
index 4559fd54937..cfbdff60d18 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmIncrRefreshManager.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/mtmv/ivm/IvmIncrRefreshManager.java
@@ -97,6 +97,8 @@ public class IvmIncrRefreshManager {
MTMV mtmv = context.getMtmv();
StatementContext statementContext = new StatementContext(
context.getConnectContext(), new
OriginStatement(mtmv.getQuerySql(), 0));
+ // SET_VAR hints are applied while parsing the MV query, before
executeCommand runs.
+ context.getConnectContext().setStatementContext(statementContext);
// The delta may only read the base partitions the MV's partition
definition keeps. A base
// partition outside that set, expired by partition_sync_limit, would
otherwise still be
// read through the delta and the join-opposite snapshot, and its rows
would have no MV
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
index 0ea2e7cb025..953e98751ef 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java
@@ -40,6 +40,7 @@ import org.apache.doris.datasource.mvcc.MvccTable;
import org.apache.doris.datasource.mvcc.MvccTableInfo;
import org.apache.doris.foundation.format.FormatOptions;
import org.apache.doris.mtmv.BaseTableInfo;
+import org.apache.doris.mtmv.MTMVCache;
import org.apache.doris.mtmv.ivm.IvmRewriteContext;
import org.apache.doris.nereids.analyzer.UnboundRelation;
import org.apache.doris.nereids.exceptions.AnalysisException;
@@ -312,6 +313,10 @@ public class StatementContext implements Closeable {
// Record mtmv and valid partitions map because this is time-consuming
behavior
private final Map<BaseTableInfo, Collection<Partition>>
mvCanRewritePartitionsMap = new HashMap<>();
+ // When the Env-wide MTMVCacheManager is disabled
(mtmv_cache_manage_num=0), reuse rewrite plans
+ // in the same statement so multiple rewrite paths do not rebuild the same
MV plan.
+ private final Map<Pair<Long, Boolean>, MTMVCache> queryLocalMtmvCaches =
new HashMap<>();
+
/// for dictionary sink.
private List<Backend> usedBackendsDistributing; // report used backends
after done distribute planning.
private long dictionaryUsedSrcVersion; // base table data version used in
this refreshing.
@@ -1484,6 +1489,14 @@ public class StatementContext implements Closeable {
this.materializationRewrittenSuccessSet.add(materializationQualifier);
}
+ public MTMVCache getQueryLocalMtmvCache(long mtmvId, boolean guarded) {
+ return queryLocalMtmvCaches.get(Pair.of(mtmvId, guarded));
+ }
+
+ public void putQueryLocalMtmvCache(long mtmvId, boolean guarded, MTMVCache
cache) {
+ queryLocalMtmvCaches.put(Pair.of(mtmvId, guarded), cache);
+ }
+
public Multimap<List<String>, Pair<RelationId, Set<String>>>
getTableUsedPartitionNameMap() {
return tableUsedPartitionNameMap;
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PlanTranslatorContext.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PlanTranslatorContext.java
index 0a09cd0670b..476579d2f12 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PlanTranslatorContext.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PlanTranslatorContext.java
@@ -133,6 +133,12 @@ public class PlanTranslatorContext {
// root and across pipeline boundaries (see shouldResetSerialFlagForChild).
private final Map<PlanNodeId, Boolean> serialAncestorInPipelineMap =
Maps.newHashMap();
+ // Per-node "does this pipeline have a serial parent pipeline" flag.
Mirrors BE's
+ // Pipeline::num_tasks_of_parent() gate in _add_local_exchange: a pipeline
whose parent
+ // has one task must not be split by another local exchange, because that
would raise only
+ // the new source side to N tasks while the paired sink/source operators
stay one-to-one.
+ private final Map<PlanNodeId, Boolean> serialParentPipelineMap =
Maps.newHashMap();
+
// Per-node "is there a downstream operator that depends on hash
distribution for
// correctness, with HASH/NOOP path connecting it to me" flag. Mirrors
BE's
// _followed_by_shuffled_operator propagation in
pipeline_fragment_context.cpp.
@@ -292,6 +298,14 @@ public class PlanTranslatorContext {
return serialAncestorInPipelineMap.getOrDefault(node.getId(), false);
}
+ public void setHasSerialParentPipeline(PlanNode node, boolean value) {
+ serialParentPipelineMap.put(node.getId(), value);
+ }
+
+ public boolean hasSerialParentPipeline(PlanNode node) {
+ return serialParentPipelineMap.getOrDefault(node.getId(), false);
+ }
+
public void setHasShuffleForCorrectnessAncestor(PlanNode node, boolean
value) {
shuffledAncestorMap.put(node.getId(), value);
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/RefreshMTMVCommand.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/RefreshMTMVCommand.java
index 583b5b2551a..04256becc73 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/RefreshMTMVCommand.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/RefreshMTMVCommand.java
@@ -152,6 +152,9 @@ public class RefreshMTMVCommand extends Command implements
Forward, Explainable
stmtCtx.setIvmRewriteContext(Optional.of(IvmRewriteContext.incrementalDryRun(mtmv,
dryRunLimit)));
// Excluded trigger tables must not be validated for binlog / key-type
support.
stmtCtx.setExcludedTriggerTables(mtmv.getExcludedTriggerTables());
+ // The MV query is parsed before the internal executor is created.
SET_VAR hints
+ // need this context already installed on the internal session during
parsing.
+ internalCtx.setStatementContext(stmtCtx);
return stmtCtx;
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/ColumnDefinition.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/ColumnDefinition.java
index 069bffb82d4..22cdbcc1f63 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/ColumnDefinition.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/ColumnDefinition.java
@@ -24,6 +24,7 @@ import org.apache.doris.catalog.AggregateType;
import org.apache.doris.catalog.Column;
import org.apache.doris.catalog.KeysType;
import org.apache.doris.common.CaseSensibility;
+import org.apache.doris.common.Config;
import org.apache.doris.common.FeNameFormat;
import org.apache.doris.common.util.SqlUtils;
import org.apache.doris.nereids.exceptions.AnalysisException;
@@ -311,6 +312,10 @@ public class ColumnDefinition {
return sb.toString();
}
+ private boolean isAggregateTableOnlyType() {
+ return type.isHllType() || type.isQuantileStateType() ||
type.isAggStateType();
+ }
+
private DataType updateCharacterTypeLength(DataType dataType) {
if (dataType instanceof ArrayType) {
return ArrayType.of(updateCharacterTypeLength(((ArrayType)
dataType).getItemType()));
@@ -384,7 +389,13 @@ public class ColumnDefinition {
*/
public void validate(boolean isOlap, Set<String> keysSet, Set<String>
clusterKeySet, boolean isEnableMergeOnWrite,
KeysType keysType) {
- validateInternal(isOlap, keysSet, clusterKeySet, isEnableMergeOnWrite,
keysType, false);
+ validate(isOlap, keysSet, clusterKeySet, isEnableMergeOnWrite,
keysType, false);
+ }
+
+ public void validate(boolean isOlap, Set<String> keysSet, Set<String>
clusterKeySet, boolean isEnableMergeOnWrite,
+ KeysType keysType, boolean isSystemGeneratedTable) {
+ validateInternal(isOlap, keysSet, clusterKeySet, isEnableMergeOnWrite,
keysType, false,
+ isSystemGeneratedTable);
}
/**
@@ -392,11 +403,11 @@ public class ColumnDefinition {
*/
public void validateNestedColumn(boolean isOlap, Set<String> keysSet,
Set<String> clusterKeySet,
boolean isEnableMergeOnWrite, KeysType keysType) {
- validateInternal(isOlap, keysSet, clusterKeySet, isEnableMergeOnWrite,
keysType, true);
+ validateInternal(isOlap, keysSet, clusterKeySet, isEnableMergeOnWrite,
keysType, true, false);
}
private void validateInternal(boolean isOlap, Set<String> keysSet,
Set<String> clusterKeySet,
- boolean isEnableMergeOnWrite, KeysType keysType, boolean
nestedColumn) {
+ boolean isEnableMergeOnWrite, KeysType keysType, boolean
nestedColumn, boolean isSystemGeneratedTable) {
try {
// if enableAddHiddenColumn is true, can add hidden column.
// So does not check if the column name starts with __DORIS_
@@ -414,6 +425,13 @@ public class ColumnDefinition {
}
type.validateDataType();
type = updateCharacterTypeLength(type);
+ if (!isSystemGeneratedTable && isOlap && keysType != KeysType.AGG_KEYS
&& isAggregateTableOnlyType()
+ && !Config.enable_non_aggregate_table_state_types) {
+ throw new AnalysisException(String.format(
+ "%s type is only supported in aggregate key tables,
column: %s. "
+ + "Set FE config
'enable_non_aggregate_table_state_types' to true to temporarily allow it",
+ type.toSql(), name));
+ }
if (type.isArrayType()) {
int depth = 0;
DataType curType = type;
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/CreateTableInfo.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/CreateTableInfo.java
index c23a95b85c6..314cee7a744 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/CreateTableInfo.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/info/CreateTableInfo.java
@@ -725,8 +725,10 @@ public class CreateTableInfo {
keysSet.addAll(keys);
Set<String> orderKeySet =
Sets.newTreeSet(String.CASE_INSENSITIVE_ORDER);
orderKeySet.addAll(sortOrderFields.stream().map(SortFieldInfo::getColumnName).collect(Collectors.toSet()));
+ // Internal statistics tables need state columns. The internal-query
flag can also be set by user SHOWs.
+ boolean isSystemGeneratedTable = targetIsInternalCatalog &&
FeConstants.INTERNAL_DB_NAME.equals(dbName);
columns.forEach(c -> c.validate(targetIsInternalCatalog, keysSet,
orderKeySet, finalEnableMergeOnWrite,
- keysType));
+ keysType, isSystemGeneratedTable));
try {
invertedIndexFileStorageFormat =
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/planner/AddLocalExchange.java
b/fe/fe-core/src/main/java/org/apache/doris/planner/AddLocalExchange.java
index 1adea56c13e..b5cf3f49766 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/planner/AddLocalExchange.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/planner/AddLocalExchange.java
@@ -197,6 +197,7 @@ public class AddLocalExchange {
? LocalExchangeTypeRequire.noRequire() :
sink.getLocalExchangeTypeRequire();
PlanNode root = fragment.getPlanRoot();
context.setHasSerialAncestorInPipeline(root, false);
+ context.setHasSerialParentPipeline(root, false);
Pair<PlanNode, LocalExchangeType> output = root
.enforceAndDeriveLocalExchange(context, null, require);
PlanNode newRoot = output.first;
diff --git a/fe/fe-core/src/main/java/org/apache/doris/planner/PlanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/planner/PlanNode.java
index a635ca26730..b592c15f627 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/planner/PlanNode.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/planner/PlanNode.java
@@ -50,6 +50,7 @@ import org.apache.doris.thrift.TPushAggOp;
import com.google.common.base.Joiner;
import com.google.common.base.Preconditions;
+import com.google.common.base.Suppliers;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
import org.apache.commons.collections4.CollectionUtils;
@@ -65,6 +66,7 @@ import java.util.Map.Entry;
import java.util.Set;
import java.util.function.Consumer;
import java.util.function.Predicate;
+import java.util.function.Supplier;
import java.util.stream.Collectors;
/**
@@ -1069,7 +1071,12 @@ public abstract class PlanNode extends
TreeNode<PlanNode> {
// node's sink, e.g. Exchange is in AGG_Sink pipeline).
// For non-splitting operators (shouldReset=false, e.g. streaming AGG):
// Inherit parent's serial flag + this node's own.
- boolean inheritedSerial = shouldResetSerialFlagForChild(childIndex)
+ boolean startsNewPipeline = shouldResetSerialFlagForChild(childIndex);
+ Supplier<Boolean> currentNodeSerialOnBe = Suppliers.memoize(
+ () ->
isSerialOperatorOnBe(translatorContext.getConnectContext()));
+ boolean currentPipelineSerial =
translatorContext.hasSerialAncestorInPipeline(this)
+ || currentNodeSerialOnBe.get();
+ boolean inheritedSerial = startsNewPipeline
? false : translatorContext.hasSerialAncestorInPipeline(this);
// Use isSerialOperatorOnBe (= isSerialNode &&
fragment.useSerialSource) instead of the
// raw isSerialNode(). BE's OperatorBase reads the Thrift
`is_serial_operator` flag —
@@ -1078,8 +1085,10 @@ public abstract class PlanNode extends
TreeNode<PlanNode> {
// Using isSerialNode here would set the child's serial-ancestor flag
wider than BE's
// view and over-skip required LocalExchanges downstream.
boolean childHasSerialAncestor = inheritedSerial
- || isSerialOperatorOnBe(translatorContext.getConnectContext());
+ || currentNodeSerialOnBe.get();
translatorContext.setHasSerialAncestorInPipeline(child,
childHasSerialAncestor);
+ translatorContext.setHasSerialParentPipeline(child, startsNewPipeline
+ ? currentPipelineSerial :
translatorContext.hasSerialParentPipeline(this));
// 1b. Propagate shuffle-for-correctness-ancestor flag to child.
// Mirrors BE's _followed_by_shuffled_operator: a downstream operator
needs hash
@@ -1140,7 +1149,8 @@ public abstract class PlanNode extends TreeNode<PlanNode>
{
// Use isSerialOperatorOnBe (not isSerialNode) because BE's
Pipeline::need_to_local_exchange
// checks op->is_serial_operator() which reads the Thrift flag set
from isSerialOperatorOnBe;
// when fragment.useSerialSource is false, BE treats this node as
non-serial.
- if (translatorContext.hasSerialAncestorInPipeline(this)
+ if (translatorContext.hasSerialParentPipeline(this)
+ || translatorContext.hasSerialAncestorInPipeline(this)
||
isSerialOperatorOnBe(translatorContext.getConnectContext())) {
return childOutput;
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/alter/InternalSchemaAlterTest.java
b/fe/fe-core/src/test/java/org/apache/doris/alter/InternalSchemaAlterTest.java
index 3ad88c49c40..af4901864d9 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/alter/InternalSchemaAlterTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/alter/InternalSchemaAlterTest.java
@@ -25,6 +25,7 @@ import org.apache.doris.catalog.InternalSchemaInitializer;
import org.apache.doris.catalog.OlapTable;
import org.apache.doris.catalog.Partition;
import org.apache.doris.catalog.PartitionInfo;
+import org.apache.doris.catalog.PrimitiveType;
import org.apache.doris.common.AnalysisException;
import org.apache.doris.common.Config;
import org.apache.doris.common.FeConstants;
@@ -89,4 +90,13 @@ public class InternalSchemaAlterTest extends
TestWithFeService {
Assertions.assertNotNull(table.getColumn(def.getName()));
}
}
+
+ @Test
+ public void testCheckPartitionStatisticsTable() throws AnalysisException {
+ Database db = Env.getCurrentEnv().getCatalogMgr()
+
.getInternalCatalog().getDbNullable(FeConstants.INTERNAL_DB_NAME);
+ Assertions.assertNotNull(db);
+ OlapTable table =
db.getOlapTableOrAnalysisException(StatisticConstants.PARTITION_STATISTIC_TBL_NAME);
+ Assertions.assertEquals(PrimitiveType.HLL,
table.getColumn("ndv").getType().getPrimitiveType());
+ }
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableTest.java
b/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableTest.java
index 5cc2899958c..37b7afd5961 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableTest.java
@@ -52,6 +52,34 @@ public class CreateTableTest extends TestWithFeService {
createDatabase("test");
}
+ @Test
+ public void testInternalQueryStateDoesNotExemptUserTable() {
+ boolean originalAllowStateTypes =
Config.enable_non_aggregate_table_state_types;
+ boolean originalInternal = connectContext.getState().isInternal();
+ boolean originalEnableAggState =
connectContext.getSessionVariable().enableAggState;
+ Config.enable_non_aggregate_table_state_types = false;
+ connectContext.getSessionVariable().enableAggState = true;
+ // An ordinary SHOW can leave the internal-query flag set on a user
connection.
+ connectContext.getState().setInternal(true);
+ try {
+ for (String keysType : new String[] {"DUPLICATE", "UNIQUE"}) {
+ for (String type : new String[] {"HLL NOT NULL",
"QUANTILE_STATE NOT NULL",
+ "AGG_STATE<sum(INT NOT NULL)>"}) {
+ AnalysisException exception =
Assertions.assertThrows(AnalysisException.class,
+ () -> createTable("CREATE TABLE
test.user_state_type (k INT, v " + type + ") "
+ + keysType + " KEY(k) DISTRIBUTED BY
HASH(k) BUCKETS 1 "
+ + "PROPERTIES('replication_num'='1')"));
+ Assertions.assertTrue(exception.getMessage().contains(
+ "type is only supported in aggregate key tables"));
+ }
+ }
+ } finally {
+ Config.enable_non_aggregate_table_state_types =
originalAllowStateTypes;
+ connectContext.getState().setInternal(originalInternal);
+ connectContext.getSessionVariable().enableAggState =
originalEnableAggState;
+ }
+ }
+
@Test
public void testDuplicateCreateTable() throws Exception {
// test
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableWithBloomFilterIndexTest.java
b/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableWithBloomFilterIndexTest.java
index c12843cf21a..21436bb6da7 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableWithBloomFilterIndexTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableWithBloomFilterIndexTest.java
@@ -19,6 +19,7 @@ package org.apache.doris.catalog;
import org.apache.doris.alter.AlterJobV2;
import org.apache.doris.catalog.info.IndexType;
+import org.apache.doris.common.Config;
import org.apache.doris.common.DdlException;
import org.apache.doris.common.ExceptionChecker;
import org.apache.doris.common.FeConstants;
@@ -446,18 +447,24 @@ public class CreateTableWithBloomFilterIndexTest extends
TestWithFeService {
@Test
public void testCreateTableWithHllBloomFilterIndex() {
- ExceptionChecker.expectThrowsWithMsg(DdlException.class,
- " HLL is not supported in bloom filter index. invalid column:
k1",
- () -> createTable("CREATE TABLE test.tbl_hll_bf (\n"
- + "v1 INT,\n"
- + "k1 HLL\n"
- + ") ENGINE=OLAP\n"
- + "DUPLICATE KEY(v1)\n"
- + "DISTRIBUTED BY HASH(v1) BUCKETS 1\n"
- + "PROPERTIES (\n"
- + "\"bloom_filter_columns\" = \"k1\",\n"
- + "\"replication_num\" = \"1\"\n"
- + ");"));
+ boolean originalValue = Config.enable_non_aggregate_table_state_types;
+ Config.enable_non_aggregate_table_state_types = true;
+ try {
+ ExceptionChecker.expectThrowsWithMsg(DdlException.class,
+ " HLL is not supported in bloom filter index. invalid
column: k1",
+ () -> createTable("CREATE TABLE test.tbl_hll_bf (\n"
+ + "v1 INT,\n"
+ + "k1 HLL\n"
+ + ") ENGINE=OLAP\n"
+ + "DUPLICATE KEY(v1)\n"
+ + "DISTRIBUTED BY HASH(v1) BUCKETS 1\n"
+ + "PROPERTIES (\n"
+ + "\"bloom_filter_columns\" = \"k1\",\n"
+ + "\"replication_num\" = \"1\"\n"
+ + ");"));
+ } finally {
+ Config.enable_non_aggregate_table_state_types = originalValue;
+ }
}
@Test
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVCacheManagerTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVCacheManagerTest.java
new file mode 100644
index 00000000000..6b94fa7c5da
--- /dev/null
+++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVCacheManagerTest.java
@@ -0,0 +1,268 @@
+// 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.
+
+package org.apache.doris.mtmv;
+
+import org.apache.doris.common.Config;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.mtmv.MTMVCacheManager.HotEntry;
+import org.apache.doris.mtmv.MTMVCacheManager.Key;
+import org.apache.doris.mtmv.MTMVCacheManager.Snapshot;
+
+import com.github.benmanes.caffeine.cache.Cache;
+import com.github.benmanes.caffeine.cache.Policy;
+import com.github.benmanes.caffeine.cache.stats.CacheStats;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import java.lang.ref.Reference;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.ConcurrentHashMap;
+
+public class MTMVCacheManagerTest {
+
+ @Test
+ public void testPutGetInvalidate() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ MTMVCache cacheGuarded = Mockito.mock(MTMVCache.class);
+ MTMVCache cacheUnguarded = Mockito.mock(MTMVCache.class);
+ manager.put(1L, true, cacheGuarded);
+ manager.put(1L, false, cacheUnguarded);
+ Assertions.assertSame(cacheGuarded, manager.getIfPresent(1L, true));
+ Assertions.assertSame(cacheUnguarded, manager.getIfPresent(1L, false));
+ Assertions.assertEquals(2L, manager.size());
+
+ manager.invalidate(1L);
+ Assertions.assertNull(manager.getIfPresent(1L, true));
+ Assertions.assertNull(manager.getIfPresent(1L, false));
+ Assertions.assertEquals(0L, manager.size());
+ }
+
+ @Test
+ public void testPutRejectsNull() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ Assertions.assertThrows(NullPointerException.class, () ->
manager.put(1L, true, null));
+ }
+
+ @Test
+ public void testDifferentMtmvsAreIndependent() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ MTMVCache c1 = Mockito.mock(MTMVCache.class);
+ MTMVCache c2 = Mockito.mock(MTMVCache.class);
+ manager.put(1L, true, c1);
+ manager.put(2L, true, c2);
+ manager.invalidate(1L);
+ Assertions.assertNull(manager.getIfPresent(1L, true));
+ Assertions.assertSame(c2, manager.getIfPresent(2L, true));
+ }
+
+ @Test
+ public void testSnapshotReportsHitAndMiss() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ MTMVCache c1 = Mockito.mock(MTMVCache.class);
+ manager.put(1L, true, c1);
+ manager.getIfPresent(1L, true);
+ manager.getIfPresent(1L, false);
+ Snapshot snap = manager.snapshot();
+ Assertions.assertEquals(1L, snap.size);
+ Assertions.assertTrue(snap.hitCount >= 1);
+ Assertions.assertTrue(snap.missCount >= 1);
+ Reference.reachabilityFence(c1);
+ }
+
+ @Test
+ public void testHotEntriesHonorsLimit() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ List<MTMVCache> values = new ArrayList<>();
+ for (int i = 0; i < 5; i++) {
+ MTMVCache c = Mockito.mock(MTMVCache.class);
+ values.add(c);
+ manager.put(i, true, c);
+ }
+ List<HotEntry> hot = manager.hotEntries(3);
+ Assertions.assertEquals(3, hot.size());
+ for (HotEntry e : hot) {
+ Assertions.assertTrue(e.idleMs >= 0,
+ "idleMs should be >= 0 when expireAfterAccess is set, got
" + e.idleMs);
+ }
+ Reference.reachabilityFence(values);
+ }
+
+ @Test
+ public void testHotEntriesEmptyForZeroOrNegativeLimit() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ MTMVCache c = Mockito.mock(MTMVCache.class);
+ manager.put(1L, true, c);
+ Assertions.assertTrue(manager.hotEntries(0).isEmpty());
+ Assertions.assertTrue(manager.hotEntries(-1).isEmpty());
+ }
+
+ @Test
+ public void testInvalidateAll() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ MTMVCache c = Mockito.mock(MTMVCache.class);
+ manager.put(1L, true, c);
+ manager.put(2L, false, c);
+ manager.invalidateAll();
+ Assertions.assertEquals(0L, manager.size());
+ }
+
+ // updateConfig() can swap the field between the two reads.
+ @Test
+ public void testSnapshotReadsOneCacheInstance() {
+ int originalMaxSize = Config.mtmv_cache_manage_num;
+ try {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ Cache<Key, MTMVCache> original = mockCache();
+ Mockito.when(original.estimatedSize()).thenReturn(7L);
+ Mockito.when(original.asMap()).thenReturn(new
ConcurrentHashMap<>());
+ Mockito.when(original.stats()).thenAnswer(invocation -> {
+ // The swap lands while snapshot() is between its reads.
+ Config.mtmv_cache_manage_num = 0;
+ manager.updateConfig();
+ return CacheStats.of(3L, 1L, 0L, 0L, 0L, 2L, 0L);
+ });
+ Deencapsulation.setField(manager, "caches", original);
+
+ Snapshot snap = manager.snapshot();
+
+ Assertions.assertEquals(7L, snap.size);
+ Assertions.assertEquals(3L, snap.hitCount);
+ Assertions.assertEquals(1L, snap.missCount);
+ Assertions.assertEquals(2L, snap.evictionCount);
+ } finally {
+ Config.mtmv_cache_manage_num = originalMaxSize;
+ }
+ }
+
+ @Test
+ public void testHotEntriesReadsOneCacheInstance() {
+ int originalMaxSize = Config.mtmv_cache_manage_num;
+ try {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ Cache<Key, MTMVCache> original = mockCache();
+ Policy<Key, MTMVCache> policy = mockPolicy();
+
Mockito.when(policy.expireAfterAccess()).thenReturn(Optional.empty());
+ Mockito.when(original.asMap()).thenReturn(
+ new ConcurrentHashMap<>(Collections.singletonMap(new
Key(1L, true),
+ Mockito.mock(MTMVCache.class))));
+ Mockito.when(original.policy()).thenAnswer(invocation -> {
+ // Swapping to a disabled cache would leave the fresh instance
empty.
+ Config.mtmv_cache_manage_num = 0;
+ manager.updateConfig();
+ return policy;
+ });
+ Deencapsulation.setField(manager, "caches", original);
+
+ List<HotEntry> hot = manager.hotEntries(10);
+
+ Assertions.assertEquals(1, hot.size());
+ Assertions.assertEquals(1L, hot.get(0).mtmvId);
+ } finally {
+ Config.mtmv_cache_manage_num = originalMaxSize;
+ }
+ }
+
+ @Test
+ public void testIsEnabledFollowsLiveMaxSize() {
+ int originalMaxSize = Config.mtmv_cache_manage_num;
+ try {
+ Config.mtmv_cache_manage_num = 10;
+ MTMVCacheManager manager = new MTMVCacheManager();
+ Assertions.assertTrue(manager.isEnabled());
+
+ Config.mtmv_cache_manage_num = 0;
+ manager.updateConfig();
+ Assertions.assertFalse(manager.isEnabled());
+ } finally {
+ Config.mtmv_cache_manage_num = originalMaxSize;
+ }
+ }
+
+ @SuppressWarnings("unchecked")
+ private static Cache<Key, MTMVCache> mockCache() {
+ return Mockito.mock(Cache.class);
+ }
+
+ @SuppressWarnings("unchecked")
+ private static Policy<Key, MTMVCache> mockPolicy() {
+ return Mockito.mock(Policy.class);
+ }
+
+ @Test
+ public void testZeroMaxSizeDisablesCacheInsteadOfUnbounding() {
+ int originalMaxSize = Config.mtmv_cache_manage_num;
+ try {
+ Config.mtmv_cache_manage_num = 0;
+ MTMVCacheManager manager = new MTMVCacheManager();
+ for (int i = 0; i < 5; i++) {
+ manager.put(i, true, Mockito.mock(MTMVCache.class));
+ }
+ manager.getCachesForTest().cleanUp();
+ Assertions.assertEquals(0L, manager.size());
+ } finally {
+ Config.mtmv_cache_manage_num = originalMaxSize;
+ }
+ }
+
+ @Test
+ public void testUpdateConfigShrinksToNewMaxSize() {
+ int originalMaxSize = Config.mtmv_cache_manage_num;
+ try {
+ Config.mtmv_cache_manage_num = 10;
+ MTMVCacheManager manager = new MTMVCacheManager();
+ List<MTMVCache> values = new ArrayList<>();
+ for (int i = 0; i < 10; i++) {
+ MTMVCache c = Mockito.mock(MTMVCache.class);
+ values.add(c);
+ manager.put(i, true, c);
+ }
+ manager.getCachesForTest().cleanUp();
+ Assertions.assertEquals(10L, manager.size());
+
+ Config.mtmv_cache_manage_num = 2;
+ manager.updateConfig();
+ Assertions.assertEquals(2L, manager.size());
+
+ Config.mtmv_cache_manage_num = 0;
+ manager.updateConfig();
+ Assertions.assertEquals(0L, manager.size());
+ MTMVCache discarded = Mockito.mock(MTMVCache.class);
+ manager.put(99L, true, discarded);
+ manager.getCachesForTest().cleanUp();
+ Assertions.assertEquals(0L, manager.size());
+ Reference.reachabilityFence(values);
+ Reference.reachabilityFence(discarded);
+ } finally {
+ Config.mtmv_cache_manage_num = originalMaxSize;
+ }
+ }
+
+ @Test
+ public void testKeyEqualityAndHash() {
+ Key k1 = new Key(42L, true);
+ Key k2 = new Key(42L, true);
+ Key k3 = new Key(42L, false);
+ Assertions.assertEquals(k1, k2);
+ Assertions.assertEquals(k1.hashCode(), k2.hashCode());
+ Assertions.assertNotEquals(k1, k3);
+ }
+}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVPlanUtilTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVPlanUtilTest.java
index c2b0f83d44f..d661d2df7f3 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVPlanUtilTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVPlanUtilTest.java
@@ -352,6 +352,22 @@ public class MTMVPlanUtilTest extends SqlTestBase {
Assertions.assertTrue(mtmvAnalyzeQueryInfo.getColumnDefinitions().size() == 2);
}
+ @Test
+ public void testCreateMTMVWithAggStateColumn() throws Exception {
+ boolean originalEnableAggState =
connectContext.getSessionVariable().enableAggState;
+ connectContext.getSessionVariable().enableAggState = true;
+ connectContext.setThreadLocalInfo();
+ try {
+ Assertions.assertDoesNotThrow(() -> createMvByNereids(
+ "create materialized view mv_with_agg_state BUILD DEFERRED
REFRESH COMPLETE ON MANUAL\n"
+ + "DISTRIBUTED BY RANDOM BUCKETS 1\n"
+ + "PROPERTIES ('replication_num' = '1')\n"
+ + "as select id, sum_union(sum_state(score)) from
test.T1 group by id"));
+ } finally {
+ connectContext.getSessionVariable().enableAggState =
originalEnableAggState;
+ }
+ }
+
@Test
public void testEnsureMTMVQueryUsable() throws Exception {
createMvByNereids("create materialized view mv1 BUILD DEFERRED REFRESH
COMPLETE ON MANUAL\n"
diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTaskTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTaskTest.java
index ac44fa26ca0..e995fc3b27f 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTaskTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTaskTest.java
@@ -558,6 +558,7 @@ public class MTMVTaskTest {
Mockito.when(mtmv.getExcludedTriggerTables()).thenReturn(excludedTriggerTables);
Mockito.when(mtmv.isIvm()).thenReturn(true);
Mockito.when(mtmv.getName()).thenReturn("test_mv");
+ Mockito.when(mtmv.getQuerySql()).thenReturn("select k1 from
test_db.base_table");
Mockito.when(mtmv.getDatabase()).thenReturn(null);
Mockito.when(mtmvPartitionInfo.getPartitionType()).thenReturn(MTMVPartitionType.FOLLOW_BASE_TABLE);
@@ -583,6 +584,8 @@ public class MTMVTaskTest {
public UpdateMvByPartitionCommand
answer(InvocationOnMock invocation) {
StatementContext statementContext =
invocation.getArgument(3);
Assertions.assertEquals(excludedTriggerTables,
statementContext.getExcludedTriggerTables());
+ Assertions.assertEquals("select k1 from
test_db.base_table",
+
statementContext.getOriginStatement().originStmt);
return command;
}
});
@@ -593,6 +596,8 @@ public class MTMVTaskTest {
public Void answer(InvocationOnMock invocation) {
StatementContext statementContext =
invocation.getArgument(2);
Assertions.assertEquals(excludedTriggerTables,
statementContext.getExcludedTriggerTables());
+ Assertions.assertEquals("select k1 from
test_db.base_table",
+
statementContext.getOriginStatement().originStmt);
return null;
}
});
diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTest.java
index 98aeb274eb9..5e7ea127b46 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVTest.java
@@ -33,6 +33,7 @@ import org.apache.doris.catalog.ScalarType;
import org.apache.doris.catalog.SinglePartitionInfo;
import org.apache.doris.catalog.info.TableNameInfo;
import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.Config;
import org.apache.doris.common.jmockit.Deencapsulation;
import org.apache.doris.common.util.PropertyAnalyzer;
import org.apache.doris.job.common.IntervalUnit;
@@ -43,11 +44,14 @@ import
org.apache.doris.mtmv.MTMVRefreshEnum.MTMVRefreshState;
import org.apache.doris.mtmv.MTMVRefreshEnum.MTMVState;
import org.apache.doris.mtmv.MTMVRefreshEnum.RefreshMethod;
import org.apache.doris.mtmv.MTMVRefreshEnum.RefreshTrigger;
+import org.apache.doris.nereids.StatementContext;
import org.apache.doris.persist.AlterMTMV;
import org.apache.doris.persist.EditLog;
import org.apache.doris.persist.EditLog.EditLogItem;
import org.apache.doris.persist.OperationType;
import org.apache.doris.persist.gson.GsonUtils;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.SessionVariable;
import org.apache.doris.thrift.TStorageType;
import com.google.common.collect.Lists;
@@ -433,6 +437,230 @@ public class MTMVTest {
Mockito.verify(editLogItem).await();
}
+ @Test
+ public void testRefreshPublishAdvancesCacheGeneration() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ HookedMTMV mtmv = buildHookedMTMV();
+ MTMVCache refreshedGuarded = Mockito.mock(MTMVCache.class);
+ MTMVCache refreshedUnguarded = Mockito.mock(MTMVCache.class);
+ mtmv.refreshGuardedCache = refreshedGuarded;
+ mtmv.refreshUnguardedCache = refreshedUnguarded;
+ long generationBefore = Deencapsulation.getField(mtmv,
"rewriteCacheGeneration");
+
+ Env env = mockEnv(manager);
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+
Assertions.assertTrue(mtmv.addTaskResult(buildSuccessTaskResult(mtmv), false));
+ }
+
+ long generationAfter = Deencapsulation.getField(mtmv,
"rewriteCacheGeneration");
+ Assertions.assertEquals(generationBefore + 1, generationAfter);
+ Assertions.assertSame(refreshedGuarded,
manager.getIfPresent(mtmv.getId(), true));
+ Assertions.assertSame(refreshedUnguarded,
manager.getIfPresent(mtmv.getId(), false));
+ }
+
+ @Test
+ public void testRefreshSkipsPlanBuildWhenCacheDisabled() {
+ int originalMaxSize = Config.mtmv_cache_manage_num;
+ try {
+ Config.mtmv_cache_manage_num = 0;
+ MTMVCacheManager manager = new MTMVCacheManager();
+ HookedMTMV mtmv = buildHookedMTMV();
+ mtmv.refreshGuardedCache = Mockito.mock(MTMVCache.class);
+ mtmv.refreshUnguardedCache = Mockito.mock(MTMVCache.class);
+ long generationBefore = Deencapsulation.getField(mtmv,
"rewriteCacheGeneration");
+
+ Env env = mockEnv(manager);
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+
Assertions.assertTrue(mtmv.addTaskResult(buildSuccessTaskResult(mtmv), false));
+ }
+
+ // The generation/invalidation transition still happens, but
neither plan was built.
+ long generationAfter = Deencapsulation.getField(mtmv,
"rewriteCacheGeneration");
+ Assertions.assertEquals(generationBefore + 1, generationAfter);
+ Assertions.assertEquals(0, mtmv.refreshBuildCount);
+ Assertions.assertNull(manager.getIfPresent(mtmv.getId(), true));
+ Assertions.assertNull(manager.getIfPresent(mtmv.getId(), false));
+ } finally {
+ Config.mtmv_cache_manage_num = originalMaxSize;
+ }
+ }
+
+ @Test
+ public void testDisabledCacheReusesPlanWithinSameStatement() throws
Exception {
+ int originalMaxSize = Config.mtmv_cache_manage_num;
+ try {
+ Config.mtmv_cache_manage_num = 0;
+ MTMVCacheManager manager = new MTMVCacheManager();
+ Assertions.assertFalse(manager.isEnabled());
+
+ HookedMTMV mtmv = buildHookedMTMV();
+ MTMVCache plan = Mockito.mock(MTMVCache.class);
+ mtmv.lazyCaches.add(plan);
+
+ ConnectContext context = mockConnectContext();
+ StatementContext statementContext = new StatementContext();
+
Mockito.when(context.getStatementContext()).thenReturn(statementContext);
+
+ Env env = mockEnv(manager);
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+ MTMVCache first = mtmv.getOrGenerateCache(context);
+ MTMVCache second = mtmv.getOrGenerateCache(context);
+ Assertions.assertSame(plan, first);
+ Assertions.assertSame(first, second);
+ }
+
+ Assertions.assertEquals(1, mtmv.lazyBuildCount);
+ Assertions.assertNull(manager.getIfPresent(mtmv.getId(), false));
+ Assertions.assertSame(plan,
statementContext.getQueryLocalMtmvCache(mtmv.getId(), false));
+ } finally {
+ Config.mtmv_cache_manage_num = originalMaxSize;
+ }
+ }
+
+ @Test
+ public void testPausedBuilderCannotRepublishPreRefreshPlan() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ HookedMTMV mtmv = buildHookedMTMV();
+ MTMVCache prePublishPlan = Mockito.mock(MTMVCache.class);
+ MTMVCache rebuiltPlan = Mockito.mock(MTMVCache.class);
+ MTMVCache refreshedUnguarded = Mockito.mock(MTMVCache.class);
+ mtmv.lazyCaches.add(prePublishPlan);
+ mtmv.lazyCaches.add(rebuiltPlan);
+ mtmv.refreshGuardedCache = Mockito.mock(MTMVCache.class);
+ mtmv.refreshUnguardedCache = refreshedUnguarded;
+
+ Env env = mockEnv(manager);
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+ // The builder snapshotted the generation and is now paused
outside the MV lock: the refresh
+ // publishes its pair and the fresh entry is then evicted before
the builder resumes.
+ mtmv.duringLazyBuild = () -> {
+
Assertions.assertTrue(mtmv.addTaskResult(buildSuccessTaskResult(mtmv), false));
+ Assertions.assertSame(refreshedUnguarded,
manager.getIfPresent(mtmv.getId(), false));
+ manager.invalidate(mtmv.getId());
+ };
+ MTMVCache published =
mtmv.getOrGenerateCache(mockConnectContext());
+
+ Assertions.assertSame(rebuiltPlan, published);
+ Assertions.assertSame(rebuiltPlan,
manager.getIfPresent(mtmv.getId(), false));
+ Assertions.assertNotSame(prePublishPlan,
manager.getIfPresent(mtmv.getId(), false));
+ }
+ }
+
+ @Test
+ public void testTaskCompletionDoesNotPublishForDroppedMv() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ HookedMTMV mtmv = buildHookedMTMV();
+ mtmv.refreshGuardedCache = Mockito.mock(MTMVCache.class);
+ mtmv.refreshUnguardedCache = Mockito.mock(MTMVCache.class);
+
+ Env env = mockEnv(manager);
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+ manager.put(mtmv.getId(), true, Mockito.mock(MTMVCache.class));
+ // The task builds its caches outside the MV lock; the drop lands
in that window.
+ mtmv.duringRefreshBuild = mtmv::markDropped;
+
Assertions.assertTrue(mtmv.addTaskResult(buildSuccessTaskResult(mtmv), false));
+
+ Assertions.assertTrue(mtmv.isDropped);
+ Assertions.assertNull(manager.getIfPresent(mtmv.getId(), true));
+ Assertions.assertNull(manager.getIfPresent(mtmv.getId(), false));
+ }
+ }
+
+ @Test
+ public void testDropStopsPausedBuilderFromPublishing() {
+ MTMVCacheManager manager = new MTMVCacheManager();
+ HookedMTMV mtmv = buildHookedMTMV();
+ MTMVCache builtPlan = Mockito.mock(MTMVCache.class);
+ mtmv.lazyCaches.add(builtPlan);
+
+ Env env = mockEnv(manager);
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+ manager.put(mtmv.getId(), true, Mockito.mock(MTMVCache.class));
+ mtmv.duringLazyBuild = mtmv::markDropped;
+ MTMVCache generated =
mtmv.getOrGenerateCache(mockConnectContext());
+
+ Assertions.assertSame(builtPlan, generated);
+ Assertions.assertNull(manager.getIfPresent(mtmv.getId(), true));
+ Assertions.assertNull(manager.getIfPresent(mtmv.getId(), false));
+ }
+ }
+
+ private HookedMTMV buildHookedMTMV() {
+ HookedMTMV mtmv = configureMTMV(new HookedMTMV());
+ mtmv.getIvmInfo();
+ return mtmv;
+ }
+
+ private Env mockEnv(MTMVCacheManager manager) {
+ Env env = Mockito.mock(Env.class);
+ EditLog editLog = Mockito.mock(EditLog.class);
+ Mockito.when(env.getEditLog()).thenReturn(editLog);
+
Mockito.when(env.getMtmvService()).thenReturn(Mockito.mock(MTMVService.class));
+ Mockito.when(env.getMtmvCacheManager()).thenReturn(manager);
+ Mockito.when(editLog.submitEdit(Mockito.anyShort(), Mockito.any()))
+ .thenReturn(Mockito.mock(EditLogItem.class));
+ return env;
+ }
+
+ private ConnectContext mockConnectContext() {
+ ConnectContext context = Mockito.mock(ConnectContext.class);
+ SessionVariable sessionVariable = Mockito.mock(SessionVariable.class);
+ Mockito.when(context.getSessionVariable()).thenReturn(sessionVariable);
+
Mockito.when(sessionVariable.getAffectQueryResultInPlanVariables()).thenReturn(Map.of());
+ return context;
+ }
+
+ private AlterMTMV buildSuccessTaskResult(MTMV mtmv) {
+ MTMVRelation relation = mtmv.getRelation();
+ MTMVTask task = new MTMVTask(mtmv, relation, null);
+ task.setStatus(TaskStatus.SUCCESS);
+ AlterMTMV alterMTMV = new AlterMTMV(new TableNameInfo("db1", "mv1"),
MTMVAlterOpType.ADD_TASK);
+ alterMTMV.setTask(task);
+ alterMTMV.setRelation(relation);
+ alterMTMV.setPartitionSnapshots(Map.of());
+ return alterMTMV;
+ }
+
+ /**
+ * Runs a hook inside the lock-free cache build so a refresh or a drop can
be interleaved with an
+ * in-flight build deterministically, without threads.
+ */
+ private static class HookedMTMV extends MTMV {
+ private final List<MTMVCache> lazyCaches = Lists.newArrayList();
+ private Runnable duringRefreshBuild;
+ private Runnable duringLazyBuild;
+ private MTMVCache refreshGuardedCache;
+ private MTMVCache refreshUnguardedCache;
+ private int lazyBuildCount;
+ private int refreshBuildCount;
+
+ @Override
+ protected MTMVCache createRewriteCache(ConnectContext currentContext,
boolean needLock,
+ boolean addSessionVarGuard) {
+ // needLock is true only on the refresh path, false on the lazy
query path.
+ Runnable hook = needLock ? duringRefreshBuild : duringLazyBuild;
+ if (needLock) {
+ duringRefreshBuild = null;
+ } else {
+ duringLazyBuild = null;
+ }
+ if (hook != null) {
+ hook.run();
+ }
+ if (needLock) {
+ refreshBuildCount++;
+ return addSessionVarGuard ? refreshGuardedCache :
refreshUnguardedCache;
+ }
+ return lazyCaches.get(Math.min(lazyBuildCount++, lazyCaches.size()
- 1));
+ }
+ }
+
private void replayAlterMvProperties(MTMV mtmv, Map<String, String>
properties) {
AlterMTMV alterMTMV = new AlterMTMV(
new TableNameInfo("db", "mv"), MTMVAlterOpType.ALTER_PROPERTY);
@@ -454,7 +682,10 @@ public class MTMVTest {
}
private MTMV buildSerializableMTMV() {
- MTMV mtmv = new MTMV();
+ return configureMTMV(new MTMV());
+ }
+
+ private <T extends MTMV> T configureMTMV(T mtmv) {
mtmv.setId(1L);
mtmv.setQualifiedDbName("db1");
mtmv.setRefreshInfo(buildMTMVRefreshInfo(mtmv));
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmIncrRefreshManagerTest.java
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmIncrRefreshManagerTest.java
index 3a4a4519e12..59a57d660a0 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmIncrRefreshManagerTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/mtmv/ivm/IvmIncrRefreshManagerTest.java
@@ -90,6 +90,33 @@ public class IvmIncrRefreshManagerTest {
Assertions.assertEquals(excluded,
captured.get().getExcludedTriggerTables());
}
+ @Test
+ public void testIncrementalRefreshParsesSetVarWithItsStatementContext()
throws Exception {
+ MTMV mtmv = mockMtmv();
+ Mockito.when(mtmv.getQuerySql()).thenReturn("SELECT /*+
SET_VAR(query_timeout=10) */ 1 AS k1");
+ Mockito.when(mtmv.getInsertedColumnNames()).thenReturn(List.of("k1"));
+ ConnectContext connectContext = new ConnectContext();
+ IvmIncrRefreshContext context = new IvmIncrRefreshContext(mtmv,
connectContext, "audit",
+ queryId -> { }, null);
+ connectContext.setThreadLocalInfo();
+ try (MockedStatic<MTMVPlanUtil> mockedUtil =
Mockito.mockStatic(MTMVPlanUtil.class)) {
+ mockedUtil.when(() -> MTMVPlanUtil.executeCommand(
+ Mockito.<ConnectContext>any(), Mockito.any(),
Mockito.any(), Mockito.any(), Mockito.any()))
+ .thenAnswer(inv -> {
+ StatementContext stmtCtx = inv.getArgument(2);
+ Assertions.assertSame(stmtCtx,
connectContext.getStatementContext());
+ Assertions.assertEquals(mtmv.getQuerySql(),
stmtCtx.getOriginStatement().originStmt);
+ Assertions.assertEquals(10,
connectContext.getSessionVariable().getQueryTimeoutS());
+ return null;
+ });
+ new IvmIncrRefreshManager().executeInternalRefresh(context);
+ mockedUtil.verify(() -> MTMVPlanUtil.executeCommand(
+ Mockito.eq(connectContext), Mockito.any(), Mockito.any(),
Mockito.any(), Mockito.any()));
+ } finally {
+ ConnectContext.remove();
+ }
+ }
+
@Test
public void testManagerReturnsSuccessForEmptyBundles() throws Exception {
MTMV mtmv = mockMtmv();
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/RefreshMTMVCommandTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/RefreshMTMVCommandTest.java
index 689d7fb0dad..81becc8b5a9 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/RefreshMTMVCommandTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/RefreshMTMVCommandTest.java
@@ -172,6 +172,24 @@ public class RefreshMTMVCommandTest {
Assertions.assertEquals(excluded, stmtCtx.getExcludedTriggerTables());
}
+ @Test
+ public void testIncrementalDryRunParsesSetVarWithItsStatementContext()
throws Exception {
+ RefreshMTMVInfo info = extractRefreshInfo("REFRESH MATERIALIZED VIEW
db1.mv1 INCREMENTAL");
+ TestRefreshMTMVCommand command = new TestRefreshMTMVCommand(info,
true);
+ MTMV mtmv = Mockito.mock(MTMV.class);
+ Mockito.when(mtmv.getQuerySql()).thenReturn("SELECT /*+
SET_VAR(query_timeout=10) */ 1 AS k1");
+ ConnectContext internalCtx = new ConnectContext();
+ internalCtx.setThreadLocalInfo();
+ try {
+ StatementContext stmtCtx =
command.createDryRunStatementContext(mtmv, internalCtx);
+ new IvmIncrRefreshManager().buildQueryPlan(mtmv);
+ Assertions.assertSame(stmtCtx, internalCtx.getStatementContext());
+ Assertions.assertEquals(10,
internalCtx.getSessionVariable().getQueryTimeoutS());
+ } finally {
+ ConnectContext.remove();
+ }
+ }
+
@Test
public void testIncrementalExplainCarriesExcludedTriggerTables() throws
Exception {
RefreshMTMVInfo info = extractRefreshInfo("REFRESH MATERIALIZED VIEW
db1.mv1 INCREMENTAL");
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommandTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommandTest.java
index 0eb933b45a4..4228865cfa4 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommandTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/UpdateMvByPartitionCommandTest.java
@@ -46,6 +46,7 @@ import org.apache.doris.nereids.util.PlanChecker;
import org.apache.doris.planner.ExchangeNode;
import org.apache.doris.planner.OlapTableSink;
import org.apache.doris.planner.PlanFragment;
+import org.apache.doris.qe.OriginStatement;
import org.apache.doris.qe.StmtExecutor;
import org.apache.doris.thrift.TPartitionType;
import org.apache.doris.utframe.TestWithFeService;
@@ -241,6 +242,7 @@ class UpdateMvByPartitionCommandTest extends
TestWithFeService {
void testRunRefreshCommandExecutesIncrementalMtmv() throws Exception {
MTMV mtmv = getMtmv("ivm_mv");
StatementContext statementContext = createStatementCtx("refresh
materialized view test.ivm_mv");
+ OriginStatement originStatement =
statementContext.getOriginStatement();
statementContext.setIvmRewriteContext(Optional.of(IvmRewriteContext.full(mtmv)));
UpdateMvByPartitionCommand command = newRefreshCommand(mtmv);
AtomicReference<StmtExecutor> executorRef = new AtomicReference<>();
@@ -260,12 +262,15 @@ class UpdateMvByPartitionCommandTest extends
TestWithFeService {
executor.getContext().getStatementContext().getIvmRewriteContext().orElseThrow().getMode());
Assertions.assertSame(executor.getContext(),
statementContext.getConnectContext());
Assertions.assertSame(statementContext,
executor.getContext().getStatementContext());
+ Assertions.assertSame(originStatement,
statementContext.getOriginStatement());
+ Assertions.assertSame(originStatement,
executor.getParsedStmt().getOrigStmt());
}
@Test
void testExecuteCommandRebindsTaskStatementContextToExecutionContext()
throws Exception {
MTMV mtmv = getMtmv("ivm_mv");
- StatementContext statementContext = new StatementContext();
+ StatementContext statementContext = createStatementCtx("refresh
materialized view test.ivm_mv");
+ OriginStatement originStatement =
statementContext.getOriginStatement();
statementContext.setIvmRewriteContext(Optional.of(IvmRewriteContext.full(mtmv)));
UpdateMvByPartitionCommand command = UpdateMvByPartitionCommand.from(
mtmv, Sets.newHashSet(), ImmutableMap.of(), statementContext);
@@ -281,6 +286,8 @@ class UpdateMvByPartitionCommandTest extends
TestWithFeService {
Assertions.assertSame(executor.getContext(),
statementContext.getConnectContext());
Assertions.assertSame(statementContext,
executor.getContext().getStatementContext());
+ Assertions.assertSame(originStatement,
statementContext.getOriginStatement());
+ Assertions.assertSame(originStatement,
executor.getParsedStmt().getOrigStmt());
}
@Test
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/info/ColumnDefinitionTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/info/ColumnDefinitionTest.java
index 6cbb4073df2..bcd54023faa 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/info/ColumnDefinitionTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/info/ColumnDefinitionTest.java
@@ -17,13 +17,38 @@
package org.apache.doris.nereids.trees.plans.commands.info;
+import org.apache.doris.catalog.AggregateType;
+import org.apache.doris.catalog.KeysType;
+import org.apache.doris.common.Config;
+import org.apache.doris.nereids.exceptions.AnalysisException;
+import org.apache.doris.nereids.types.AggStateType;
+import org.apache.doris.nereids.types.DataType;
+import org.apache.doris.nereids.types.HllType;
+import org.apache.doris.nereids.types.IntegerType;
+import org.apache.doris.nereids.types.QuantileStateType;
import org.apache.doris.nereids.types.StringType;
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableSet;
+import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import java.util.Optional;
+
public class ColumnDefinitionTest {
+ @BeforeEach
+ public void setUp() {
+ Config.enable_non_aggregate_table_state_types = false;
+ }
+
+ @AfterEach
+ public void tearDown() {
+ Config.enable_non_aggregate_table_state_types = false;
+ }
+
@Test
public void testNameEquals() {
ColumnDefinition columnDefinition = new ColumnDefinition("col1", null,
false, null, false, null, null);
@@ -43,4 +68,72 @@ public class ColumnDefinitionTest {
String sql = columnDefinition.toSql();
Assertions.assertTrue(sql.endsWith("COMMENT \"\""));
}
+
+ @Test
+ public void testStateTypesRequireAggregateKeyTableByDefault() {
+ for (KeysType keysType : ImmutableList.of(KeysType.DUP_KEYS,
KeysType.UNIQUE_KEYS)) {
+ for (DataType type : aggregateTableOnlyTypes()) {
+ ColumnDefinition column = new ColumnDefinition(
+ "v", type, false, null, false, Optional.empty(), "");
+
+ AnalysisException exception =
Assertions.assertThrows(AnalysisException.class,
+ () -> validateColumn(column, keysType));
+ Assertions.assertTrue(exception.getMessage().contains(
+ type.toSql() + " type is only supported in aggregate
key tables"));
+ }
+ }
+ }
+
+ @Test
+ public void testTemporaryConfigAllowsStateTypesInNonAggregateTable() {
+ Config.enable_non_aggregate_table_state_types = true;
+
+ for (KeysType keysType : ImmutableList.of(KeysType.DUP_KEYS,
KeysType.UNIQUE_KEYS)) {
+ for (DataType type : aggregateTableOnlyTypes()) {
+ ColumnDefinition column = new ColumnDefinition(
+ "v", type, false, null, false, Optional.empty(), "");
+ Assertions.assertDoesNotThrow(() -> validateColumn(column,
keysType));
+ }
+ }
+ }
+
+ @Test
+ public void testStateTypesRemainSupportedInAggregateKeyTable() {
+ Assertions.assertDoesNotThrow(() -> validateColumn(new
ColumnDefinition(
+ "v", HllType.INSTANCE, false, AggregateType.HLL_UNION, false,
Optional.empty(), ""),
+ KeysType.AGG_KEYS));
+ Assertions.assertDoesNotThrow(() -> validateColumn(new
ColumnDefinition(
+ "v", QuantileStateType.INSTANCE, false,
AggregateType.QUANTILE_UNION, false, Optional.empty(), ""),
+ KeysType.AGG_KEYS));
+ Assertions.assertDoesNotThrow(() -> validateColumn(new
ColumnDefinition(
+ "v", aggStateType(), false, AggregateType.GENERIC, false,
Optional.empty(), ""),
+ KeysType.AGG_KEYS));
+ }
+
+ @Test
+ public void testSystemGeneratedTableAllowsStateTypesInNonAggregateTable() {
+ for (KeysType keysType : ImmutableList.of(KeysType.DUP_KEYS,
KeysType.UNIQUE_KEYS)) {
+ for (DataType type : aggregateTableOnlyTypes()) {
+ ColumnDefinition column = new ColumnDefinition(
+ "v", type, false, null, false, Optional.empty(), "");
+ Assertions.assertDoesNotThrow(() ->
validateSystemGeneratedColumn(column, keysType));
+ }
+ }
+ }
+
+ private static ImmutableList<DataType> aggregateTableOnlyTypes() {
+ return ImmutableList.of(HllType.INSTANCE, QuantileStateType.INSTANCE,
aggStateType());
+ }
+
+ private static AggStateType aggStateType() {
+ return new AggStateType("sum", ImmutableList.of(IntegerType.INSTANCE),
ImmutableList.of(false), false);
+ }
+
+ private static void validateColumn(ColumnDefinition column, KeysType
keysType) {
+ column.validate(true, ImmutableSet.of("k"), ImmutableSet.of(), true,
keysType);
+ }
+
+ private static void validateSystemGeneratedColumn(ColumnDefinition column,
KeysType keysType) {
+ column.validate(true, ImmutableSet.of("k"), ImmutableSet.of(), true,
keysType, true);
+ }
}
diff --git a/regression-test/data/mtmv_p0/test_mtmv_cache_proc.out
b/regression-test/data/mtmv_p0/test_mtmv_cache_proc.out
new file mode 100644
index 00000000000..ef0bfb7ea65
--- /dev/null
+++ b/regression-test/data/mtmv_p0/test_mtmv_cache_proc.out
@@ -0,0 +1,5 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !mtmv_cache_dir --
+hot Top hot mtmv cache entries
+stat Global cache stats
+
diff --git a/regression-test/suites/correctness_p0/test_default_hll.groovy
b/regression-test/suites/correctness_p0/test_default_hll.groovy
index b21869e30e3..dc7c612bdf2 100644
--- a/regression-test/suites/correctness_p0/test_default_hll.groovy
+++ b/regression-test/suites/correctness_p0/test_default_hll.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("test_default_hll") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
def tableName = "test_default_hll"
sql """ DROP TABLE IF EXISTS ${tableName} """
@@ -96,4 +98,6 @@ suite("test_default_hll") {
qt_stream_load_csv1 """ select HLL_CARDINALITY(h1) from ${tableName} order
by k; """
-}
\ No newline at end of file
+ }
+ }
+}
diff --git
a/regression-test/suites/data_model_p0/duplicate/storage/test_duplicate_hll.groovy
b/regression-test/suites/data_model_p0/duplicate/storage/test_duplicate_hll.groovy
index 0c88f276a06..3c61f7b0fa9 100644
---
a/regression-test/suites/data_model_p0/duplicate/storage/test_duplicate_hll.groovy
+++
b/regression-test/suites/data_model_p0/duplicate/storage/test_duplicate_hll.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("test_duplicate_table_hll") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
sql "sync;"
@@ -68,4 +70,6 @@ suite("test_duplicate_table_hll") {
DISTRIBUTED BY HASH(k) BUCKETS 1 properties("replication_num"
= "1"); """
exception "Key column can not set complex type:k"
}
+ }
+ }
}
diff --git
a/regression-test/suites/data_model_p0/duplicate/storage/test_duplicate_quantile_state.groovy
b/regression-test/suites/data_model_p0/duplicate/storage/test_duplicate_quantile_state.groovy
index 9c4e07094b6..1715bbacbdd 100644
---
a/regression-test/suites/data_model_p0/duplicate/storage/test_duplicate_quantile_state.groovy
+++
b/regression-test/suites/data_model_p0/duplicate/storage/test_duplicate_quantile_state.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("test_duplicate_table_quantile_state") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
sql "sync;"
@@ -64,4 +66,6 @@ suite("test_duplicate_table_quantile_state") {
DISTRIBUTED BY HASH(k) BUCKETS 1 properties("replication_num"
= "1"); """
exception "Key column can not set complex type:k"
}
+ }
+ }
}
diff --git
a/regression-test/suites/data_model_p0/test_state_types_only_in_aggregate_table.groovy
b/regression-test/suites/data_model_p0/test_state_types_only_in_aggregate_table.groovy
new file mode 100644
index 00000000000..f74f026e7c3
--- /dev/null
+++
b/regression-test/suites/data_model_p0/test_state_types_only_in_aggregate_table.groovy
@@ -0,0 +1,113 @@
+// 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_state_types_only_in_aggregate_table") {
+ context.reconnectToMasterFe()
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: false]) {
+ sql "set enable_agg_state=true"
+ // SHOW executes an internal query, but must not exempt subsequent
user DDL.
+ sql "show table status"
+
+ sql "drop table if exists state_type_dup_hll"
+ test {
+ sql """
+ create table state_type_dup_hll (
+ k int,
+ v hll not null
+ ) duplicate key(k)
+ distributed by hash(k) buckets 1
+ properties("replication_num" = "1")
+ """
+ exception "type is only supported in aggregate key tables"
+ }
+
+ sql "drop table if exists state_type_unique_quantile"
+ test {
+ sql """
+ create table state_type_unique_quantile (
+ k int,
+ v quantile_state not null
+ ) unique key(k)
+ distributed by hash(k) buckets 1
+ properties("replication_num" = "1")
+ """
+ exception "type is only supported in aggregate key tables"
+ }
+
+ sql "drop table if exists state_type_dup_agg_state"
+ test {
+ sql """
+ create table state_type_dup_agg_state (
+ k int,
+ v agg_state<sum(int not null)> generic
+ ) duplicate key(k)
+ distributed by hash(k) buckets 1
+ properties("replication_num" = "1")
+ """
+ exception "DUP_KEYS table should not specify aggregate type"
+ }
+
+ sql "drop table if exists state_type_alter_dup"
+ sql """
+ create table state_type_alter_dup (
+ k int,
+ v int
+ ) duplicate key(k)
+ distributed by hash(k) buckets 1
+ properties("replication_num" = "1")
+ """
+ test {
+ sql "alter table state_type_alter_dup add column h hll not
null"
+ exception "type is only supported in aggregate key tables"
+ }
+ test {
+ sql "alter table state_type_alter_dup add column q
quantile_state not null"
+ exception "type is only supported in aggregate key tables"
+ }
+ test {
+ sql "alter table state_type_alter_dup add column a
agg_state<sum(int not null)> generic"
+ exception "type is only supported in aggregate key tables"
+ }
+
+ sql "drop table if exists state_type_aggregate"
+ sql """
+ create table state_type_aggregate (
+ k int,
+ h hll hll_union not null,
+ q quantile_state quantile_union not null,
+ a agg_state<sum(int not null)> generic
+ ) aggregate key(k)
+ distributed by hash(k) buckets 1
+ properties("replication_num" = "1")
+ """
+
+ setFeConfigTemporary([enable_non_aggregate_table_state_types:
true]) {
+ sql "drop table if exists state_type_compatibility_dup"
+ sql """
+ create table state_type_compatibility_dup (
+ k int,
+ h hll not null,
+ q quantile_state not null
+ ) duplicate key(k)
+ distributed by hash(k) buckets 1
+ properties("replication_num" = "1")
+ """
+ }
+ }
+ }
+}
diff --git a/regression-test/suites/data_model_p0/unique/test_unique_hll.groovy
b/regression-test/suites/data_model_p0/unique/test_unique_hll.groovy
index 035f6b1cb37..a26e266abd2 100644
--- a/regression-test/suites/data_model_p0/unique/test_unique_hll.groovy
+++ b/regression-test/suites/data_model_p0/unique/test_unique_hll.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("test_unique_table_hll") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
for (def enable_mow : [true, false]) {
sql "sync;"
@@ -70,4 +72,6 @@ suite("test_unique_table_hll") {
exception "Key column can not set complex type:k"
}
}
+ }
+ }
}
diff --git
a/regression-test/suites/data_model_p0/unique/test_unique_quantile_state.groovy
b/regression-test/suites/data_model_p0/unique/test_unique_quantile_state.groovy
index 9f2b2a5475a..70d23a34b17 100644
---
a/regression-test/suites/data_model_p0/unique/test_unique_quantile_state.groovy
+++
b/regression-test/suites/data_model_p0/unique/test_unique_quantile_state.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("test_unique_table_quantile_state") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
for (def enable_mow : [true, false]) {
sql "sync;"
@@ -66,4 +68,6 @@ suite("test_unique_table_quantile_state") {
exception "Key column can not set complex type:k"
}
}
+ }
+ }
}
diff --git
a/regression-test/suites/external_table_p0/remote_doris/test_remote_doris_unique_table_select.groovy
b/regression-test/suites/external_table_p0/remote_doris/test_remote_doris_unique_table_select.groovy
index 768deb9c81b..0d6a00026f6 100644
---
a/regression-test/suites/external_table_p0/remote_doris/test_remote_doris_unique_table_select.groovy
+++
b/regression-test/suites/external_table_p0/remote_doris/test_remote_doris_unique_table_select.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("test_remote_doris_unique_table_select", "p0,external") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
String remote_doris_host =
context.config.otherConfigs.get("extArrowFlightSqlHost")
String remote_doris_arrow_port =
context.config.otherConfigs.get("extArrowFlightSqlPort")
String remote_doris_http_port =
context.config.otherConfigs.get("extArrowFlightHttpPort")
@@ -235,4 +237,6 @@ suite("test_remote_doris_unique_table_select",
"p0,external") {
sql """ DROP DATABASE IF EXISTS `${db_name}` """
sql """ DROP CATALOG IF EXISTS `${catalog_name}` """
+ }
+ }
}
diff --git a/regression-test/suites/mtmv_p0/ivm/test_ivm_refresh_dry_run.groovy
b/regression-test/suites/mtmv_p0/ivm/test_ivm_refresh_dry_run.groovy
index f2583e14e8d..073d400a1b5 100644
--- a/regression-test/suites/mtmv_p0/ivm/test_ivm_refresh_dry_run.groovy
+++ b/regression-test/suites/mtmv_p0/ivm/test_ivm_refresh_dry_run.groovy
@@ -48,7 +48,7 @@ suite("test_ivm_refresh_dry_run") {
BUILD DEFERRED REFRESH INCREMENTAL ON MANUAL
DISTRIBUTED BY RANDOM BUCKETS 2
PROPERTIES ('replication_num' = '1')
- AS SELECT k1, COUNT(*) AS cnt, SUM(v1) AS sum_v1
+ AS SELECT /*+ SET_VAR(query_timeout=180) */ k1, COUNT(*) AS cnt,
SUM(v1) AS sum_v1
FROM test_ivm_refresh_dry_run_base
GROUP BY k1
"""
diff --git a/regression-test/suites/mtmv_p0/test_mtmv_cache_proc.groovy
b/regression-test/suites/mtmv_p0/test_mtmv_cache_proc.groovy
new file mode 100644
index 00000000000..ea64ea65e57
--- /dev/null
+++ b/regression-test/suites/mtmv_p0/test_mtmv_cache_proc.groovy
@@ -0,0 +1,83 @@
+// 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_mtmv_cache_proc", "mtmv,nonConcurrent") {
+ def dbName = "regression_test_mtmv_p0"
+ def mvName = "mtmv_cache_proc_mv"
+
+ sql """drop materialized view if exists ${mvName}"""
+ sql """drop table if exists t_test_mtmv_cache_proc_user"""
+
+ sql """
+ CREATE TABLE IF NOT EXISTS t_test_mtmv_cache_proc_user (
+ event_day DATE,
+ id BIGINT,
+ username VARCHAR(20)
+ )
+ DISTRIBUTED BY HASH(id) BUCKETS 2
+ PROPERTIES ('replication_num' = '1');
+ """
+
+ // The '/mtmv_cache' directory listing is a fixed pair of children.
+ order_qt_mtmv_cache_dir """SHOW PROC '/mtmv_cache'"""
+
+ // SHOW PROC '/mtmv_cache/stat' returns the same KV rows regardless of
cache contents; the
+ // values themselves are runtime-dependent, so only the key set can be
checked here.
+ def statRows = sql """SHOW PROC '/mtmv_cache/stat'"""
+ def statKeys = statRows.collect { it[0] }
+ ["size", "hitCount", "missCount", "evictionCount", "hitRate"].each {
+ assertTrue(statKeys.contains(it), "stat missing key: ${it}")
+ }
+
+ // Create an MV and trigger cache fill via a rewrite-eligible query.
+ sql """
+ CREATE MATERIALIZED VIEW ${mvName}
+ BUILD DEFERRED REFRESH COMPLETE ON MANUAL
+ DISTRIBUTED BY RANDOM BUCKETS 2
+ PROPERTIES ('replication_num' = '1')
+ AS
+ SELECT event_day, id, username FROM t_test_mtmv_cache_proc_user;
+ """
+ def jobName = getJobName(dbName, mvName)
+ sql """REFRESH MATERIALIZED VIEW ${mvName} AUTO"""
+ waitingMTMVTaskFinished(jobName)
+ // Query the base table so nereids checks the MV — fills the cache.
+ sql """SELECT event_day, id, username FROM t_test_mtmv_cache_proc_user"""
+
+ // hot proc: 5 columns; our MV MUST appear with its real DbName/MvName.
+ def hotRows = sql """SHOW PROC '/mtmv_cache/hot'"""
+ assertTrue(!hotRows.isEmpty(),
+ "hot cache should contain at least one entry after the MV was
queried")
+ assertEquals(5, hotRows[0].size())
+ def mvRow = hotRows.find { it[2] == mvName }
+ assertNotNull(mvRow, "MV ${mvName} should be visible in /mtmv_cache/hot
after query")
+ assertEquals(dbName, mvRow[1])
+ assertTrue(mvRow[3] == "Yes" || mvRow[3] == "No")
+ assertTrue((mvRow[4] as Long) >= 0L, "IdleMs must be non-negative")
+
+ // mtmv_cache_hot_show_num caps the row count.
+ def originalCap = sql """ADMIN SHOW FRONTEND CONFIG LIKE
'mtmv_cache_hot_show_num'"""
+ def originalCapVal = originalCap.isEmpty() ? "500" : originalCap[0][1]
+ try {
+ sql """ADMIN SET FRONTEND CONFIG ('mtmv_cache_hot_show_num' = '1')"""
+ def capped = sql """SHOW PROC '/mtmv_cache/hot'"""
+ assertEquals(1, capped.size(),
+ "hot row count should be exactly 1 after capping to 1, got
${capped.size()}")
+ } finally {
+ sql """ADMIN SET FRONTEND CONFIG ('mtmv_cache_hot_show_num' =
'${originalCapVal}')"""
+ }
+}
diff --git a/regression-test/suites/mv_p0/mv_negative/dup_negative_test.groovy
b/regression-test/suites/mv_p0/mv_negative/dup_negative_test.groovy
index 446954d6fdb..149cee3568f 100644
--- a/regression-test/suites/mv_p0/mv_negative/dup_negative_test.groovy
+++ b/regression-test/suites/mv_p0/mv_negative/dup_negative_test.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("dup_negative_mv_test", "mv_negative") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
// this mv rewrite would not be rewritten in RBO phase, so set TRY_IN_RBO
explicitly to make case stable
sql "set pre_materialized_view_rewrite_strategy = TRY_IN_RBO"
@@ -153,4 +155,6 @@ suite("dup_negative_mv_test", "mv_negative") {
}
+ }
+ }
}
diff --git a/regression-test/suites/mv_p0/mv_negative/mor_negative_test.groovy
b/regression-test/suites/mv_p0/mv_negative/mor_negative_test.groovy
index 5cd3264d6a7..e806507e039 100644
--- a/regression-test/suites/mv_p0/mv_negative/mor_negative_test.groovy
+++ b/regression-test/suites/mv_p0/mv_negative/mor_negative_test.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("mor_negative_mv_test", "mv_negative") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
String db = context.config.getDbNameByFile(context.file)
def prefix_str = "mv_mor_negative"
@@ -157,4 +159,6 @@ suite("mor_negative_mv_test", "mv_negative") {
}
+ }
+ }
}
diff --git a/regression-test/suites/mv_p0/mv_negative/mow_negative_test.groovy
b/regression-test/suites/mv_p0/mv_negative/mow_negative_test.groovy
index 760614a2038..7d598e3aabf 100644
--- a/regression-test/suites/mv_p0/mv_negative/mow_negative_test.groovy
+++ b/regression-test/suites/mv_p0/mv_negative/mow_negative_test.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("mow_negative_mv_test", "mv_negative") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
String db = context.config.getDbNameByFile(context.file)
def prefix_str = "mv_mow_negative"
@@ -158,4 +160,6 @@ suite("mow_negative_mv_test", "mv_negative") {
}
+ }
+ }
}
diff --git
a/regression-test/suites/nereids_p0/local_shuffle/test_local_shuffle_rqg_bugs.groovy
b/regression-test/suites/nereids_p0/local_shuffle/test_local_shuffle_rqg_bugs.groovy
index 4401f0362d3..f53ab6204f5 100644
---
a/regression-test/suites/nereids_p0/local_shuffle/test_local_shuffle_rqg_bugs.groovy
+++
b/regression-test/suites/nereids_p0/local_shuffle/test_local_shuffle_rqg_bugs.groovy
@@ -1663,5 +1663,40 @@ suite("test_local_shuffle_rqg_bugs") {
assertTrue(false, "Bug 26: ${t.message}")
}
+ // Bug 27: do not insert a local exchange in a pipeline whose parent
pipeline is serial.
+ // Otherwise the local exchange raises the lower AggSink pipeline to N
tasks while its
+ // paired AggSource remains at one task, leaving task 1+ without a source
dependency.
+ def bug27Query = { planner -> """
+ SELECT /*+SET_VAR(enable_local_shuffle_planner=${planner},
+ enable_local_exchange_before_agg=false,
+ enable_local_exchange_before_streaming_agg=true,
+ enable_broadcast_join_force_passthrough=true,
+ enable_share_hash_table_for_broadcast_join=false,
+ parallel_pipeline_task_num=3,
+ enable_sql_cache=false)*/
+ COUNT(*),
+ COUNT(DISTINCT CAST(f.pk AS STRING)),
+ MIN(CAST(f.pk AS STRING)),
+ MAX(CAST(f.pk AS STRING)),
+ COUNT(DISTINCT CAST(d.pk AS STRING)),
+ MIN(CAST(d.pk AS STRING)),
+ MAX(CAST(d.pk AS STRING)),
+ COUNT(DISTINCT CAST(f.col_int_undef_signed AS STRING)),
+ MIN(CAST(f.col_int_undef_signed AS STRING)),
+ MAX(CAST(f.col_int_undef_signed AS STRING))
+ FROM rqg_t1 f
+ INNER JOIN (
+ SELECT * FROM rqg_t2 d
+ WHERE d.col_int_undef_signed = 1 AND COALESCE(d.pk, -1) >= 10
+ ) d ON d.col_int_undef_signed = f.col_int_undef_signed
+ AND f.col_int_undef_signed = f.col_int_undef_signed2
+ AND f.pk = d.pk
+ """ }
+
+ def bug27BeResult = sql bug27Query(false)
+ for (int i = 0; i < 20; i++) {
+ assertEquals(bug27BeResult, sql(bug27Query(true)), "Bug 27 run ${i}")
+ }
+
logger.info("=== All RQG bug reproduction tests completed ===")
}
diff --git
a/regression-test/suites/query_p0/aggregate/support_type/any_value/any_value.groovy
b/regression-test/suites/query_p0/aggregate/support_type/any_value/any_value.groovy
index 68df1a9b57d..b43ae350e66 100644
---
a/regression-test/suites/query_p0/aggregate/support_type/any_value/any_value.groovy
+++
b/regression-test/suites/query_p0/aggregate/support_type/any_value/any_value.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("any_value") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
sql "set enable_decimal256 = true;"
sql """
drop table if exists d_table;
@@ -102,4 +104,6 @@ suite("any_value") {
qt_sql_bitmap """select bitmap_to_string(any_value(col_bitmap)) from
d_table;"""
qt_sql_hll """select hll_cardinality(any_value(col_hll)) from d_table;"""
qt_sql_quantile_state """select
QUANTILE_PERCENT(any_value(col_quantile_state), 0.5) from d_table;"""
-}
\ No newline at end of file
+ }
+ }
+}
diff --git a/regression-test/suites/query_p0/join/test_join_on.groovy
b/regression-test/suites/query_p0/join/test_join_on.groovy
index 042e16b1b2c..577ff414a25 100644
--- a/regression-test/suites/query_p0/join/test_join_on.groovy
+++ b/regression-test/suites/query_p0/join/test_join_on.groovy
@@ -16,6 +16,8 @@
// under the License.
suite("test_join_on", "query_p0") {
+ withGlobalLock("enable_non_aggregate_table_state_types") {
+ setFeConfigTemporary([enable_non_aggregate_table_state_types: true]) {
sql "DROP TABLE IF EXISTS join_on"
sql """
@@ -49,4 +51,6 @@ suite("test_join_on", "query_p0") {
sql """select * from (select cast('' as variant) as a) t1 join (select
cast('' as variant) as a) t2 on t1.a = t2.a"""
exception "could not used in ComparisonPredicate (a = a)"
}
+ }
+ }
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]