This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new bde96dafa5e branch-4.1: [fix](workload policy)Remove set session
variable workload policy action #64856 (#67009)
bde96dafa5e is described below
commit bde96dafa5e0f87df9a31d29eeb8876f2ee8c17e
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Mon Aug 24 09:51:31 2026 +0800
branch-4.1: [fix](workload policy)Remove set session variable workload
policy action #64856 (#67009)
Cherry-picked from #64856
Co-authored-by: feiniaofeiafei <[email protected]>
---
.../antlr4/org/apache/doris/nereids/DorisLexer.g4 | 1 -
.../antlr4/org/apache/doris/nereids/DorisParser.g4 | 4 +-
.../main/java/org/apache/doris/catalog/Env.java | 1 -
.../doris/nereids/parser/LogicalPlanBuilder.java | 15 +-
.../workloadschedpolicy/WorkloadAction.java | 2 -
.../WorkloadActionSetSessionVar.java | 67 --------
.../workloadschedpolicy/WorkloadActionType.java | 1 +
.../workloadschedpolicy/WorkloadSchedPolicy.java | 21 +--
.../WorkloadSchedPolicyMgr.java | 179 +--------------------
.../WorkloadSchedPolicyMgrTest.java | 63 +-------
.../test_workload_sched_policy.out | 12 +-
.../test_workload_sched_policy.groovy | 89 ++--------
12 files changed, 40 insertions(+), 415 deletions(-)
diff --git a/fe/fe-core/src/main/antlr4/org/apache/doris/nereids/DorisLexer.g4
b/fe/fe-core/src/main/antlr4/org/apache/doris/nereids/DorisLexer.g4
index e61155ed7ff..1b6b76a2059 100644
--- a/fe/fe-core/src/main/antlr4/org/apache/doris/nereids/DorisLexer.g4
+++ b/fe/fe-core/src/main/antlr4/org/apache/doris/nereids/DorisLexer.g4
@@ -524,7 +524,6 @@ SESSION: 'SESSION';
SESSION_USER: 'SESSION_USER';
SET: 'SET';
SETS: 'SETS';
-SET_SESSION_VARIABLE: 'SET_SESSION_VARIABLE';
SHAPE: 'SHAPE';
SHOW: 'SHOW';
SIGNED: 'SIGNED';
diff --git a/fe/fe-core/src/main/antlr4/org/apache/doris/nereids/DorisParser.g4
b/fe/fe-core/src/main/antlr4/org/apache/doris/nereids/DorisParser.g4
index 0edad842cf7..f3c07685299 100644
--- a/fe/fe-core/src/main/antlr4/org/apache/doris/nereids/DorisParser.g4
+++ b/fe/fe-core/src/main/antlr4/org/apache/doris/nereids/DorisParser.g4
@@ -925,8 +925,7 @@ workloadPolicyActions
;
workloadPolicyAction
- : SET_SESSION_VARIABLE STRING_LITERAL
- | identifier (STRING_LITERAL)?
+ : identifier (STRING_LITERAL)?
;
workloadPolicyConditions
@@ -2339,7 +2338,6 @@ nonReserved
| MICROSECOND
| SEPARATOR
| SERIALIZABLE
- | SET_SESSION_VARIABLE
| SESSION
| SESSION_USER
| SHAPE
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 1bcd2af4d1f..43235d9ac51 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
@@ -2035,7 +2035,6 @@ public class Env {
dnsCache.start();
- workloadSchedPolicyMgr.start();
workloadRuntimeStatusMgr.start();
admissionControl.start();
splitSourceManager.start();
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/parser/LogicalPlanBuilder.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/parser/LogicalPlanBuilder.java
index c9e90765ce0..a33f7ba4f64 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/parser/LogicalPlanBuilder.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/parser/LogicalPlanBuilder.java
@@ -8389,16 +8389,11 @@ public class LogicalPlanBuilder extends
DorisParserBaseVisitor<Object> {
for (DorisParser.WorkloadPolicyActionContext actionCtx :
ctx.workloadPolicyActions().workloadPolicyAction()) {
try {
- if (actionCtx.SET_SESSION_VARIABLE() != null) {
- actions.add(new
WorkloadActionMeta("SET_SESSION_VARIABLE",
-
stripQuotes(actionCtx.STRING_LITERAL().getText())));
- } else {
- String identifier = actionCtx.identifier().getText();
- String value = actionCtx.STRING_LITERAL() != null
- ?
stripQuotes(actionCtx.STRING_LITERAL().getText())
- : null;
- actions.add(new WorkloadActionMeta(identifier, value));
- }
+ String identifier = actionCtx.identifier().getText();
+ String value = actionCtx.STRING_LITERAL() != null
+ ? stripQuotes(actionCtx.STRING_LITERAL().getText())
+ : null;
+ actions.add(new WorkloadActionMeta(identifier, value));
} catch (UserException e) {
throw new AnalysisException(e.getMessage(), e);
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadAction.java
b/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadAction.java
index 661ea6a45fa..0b8d2980530 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadAction.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadAction.java
@@ -30,8 +30,6 @@ public interface WorkloadAction {
throws UserException {
if (WorkloadActionType.CANCEL_QUERY.equals(workloadActionMeta.action))
{
return WorkloadActionCancelQuery.createWorkloadAction();
- } else if
(WorkloadActionType.SET_SESSION_VARIABLE.equals(workloadActionMeta.action)) {
- return
WorkloadActionSetSessionVar.createWorkloadAction(workloadActionMeta.actionArgs);
}
throw new UserException("invalid action type " +
workloadActionMeta.action);
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadActionSetSessionVar.java
b/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadActionSetSessionVar.java
deleted file mode 100644
index 775fa91b995..00000000000
---
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadActionSetSessionVar.java
+++ /dev/null
@@ -1,67 +0,0 @@
-// 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.resource.workloadschedpolicy;
-
-import org.apache.doris.analysis.SetVar;
-import org.apache.doris.analysis.StringLiteral;
-import org.apache.doris.common.UserException;
-import org.apache.doris.qe.VariableMgr;
-
-import org.apache.commons.lang3.StringUtils;
-import org.apache.logging.log4j.LogManager;
-import org.apache.logging.log4j.Logger;
-
-public class WorkloadActionSetSessionVar implements WorkloadAction {
-
- private static final Logger LOG =
LogManager.getLogger(WorkloadActionSetSessionVar.class);
-
- private String varName;
- private String varValue;
-
- public WorkloadActionSetSessionVar(String varName, String varValue) {
- this.varName = varName;
- this.varValue = varValue;
- }
-
- @Override
- public void exec(WorkloadQueryInfo queryInfo) {
- try {
- SetVar var = new SetVar(varName, new StringLiteral(varValue));
- VariableMgr.setVar(queryInfo.context.getSessionVariable(), var);
- } catch (Throwable t) {
- LOG.error("error happens when exec {}",
WorkloadActionType.SET_SESSION_VARIABLE, t);
- }
- }
-
- @Override
- public WorkloadActionType getWorkloadActionType() {
- return WorkloadActionType.SET_SESSION_VARIABLE;
- }
-
- public String getVarName() {
- return varName;
- }
-
- public static WorkloadAction createWorkloadAction(String actionCmdArgs)
throws UserException {
- String[] strs = actionCmdArgs.split("=");
- if (strs.length != 2 || StringUtils.isEmpty(strs[0].trim()) ||
StringUtils.isEmpty(strs[1].trim())) {
- throw new UserException("illegal arguments, it should be like
set_session_variable \"xxx=xxx\"");
- }
- return new WorkloadActionSetSessionVar(strs[0].trim(), strs[1].trim());
- }
-}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadActionType.java
b/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadActionType.java
index e80c202a763..b3cb2cfb72c 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadActionType.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadActionType.java
@@ -20,5 +20,6 @@ package org.apache.doris.resource.workloadschedpolicy;
public enum WorkloadActionType {
CANCEL_QUERY, // cancel query
MOVE_QUERY_TO_GROUP, // move query from one wg group to another
+ // Deprecated: only kept for deserializing old workload policy metadata.
SET_SESSION_VARIABLE
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicy.java
b/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicy.java
index ff27a08706b..b93e7a56900 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicy.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicy.java
@@ -73,8 +73,6 @@ public class WorkloadSchedPolicy implements Writable,
GsonPostProcessable {
private List<WorkloadCondition> workloadConditionList;
private List<WorkloadAction> workloadActionList;
- private Boolean isFePolicy = null;
-
// for ut
public WorkloadSchedPolicy() {
}
@@ -207,21 +205,6 @@ public class WorkloadSchedPolicy implements Writable,
GsonPostProcessable {
return actionMetaList;
}
- // true, current policy can only run in FE;
- // false, current policy can only run in BE
- public boolean isFePolicy() {
- if (isFePolicy == null) {
- isFePolicy = false;
- for (WorkloadAction action : workloadActionList) {
- if
(WorkloadSchedPolicyMgr.FE_ACTION_SET.contains(action.getWorkloadActionType()))
{
- isFePolicy = true;
- break;
- }
- }
- }
- return isFePolicy;
- }
-
public TopicInfo toTopicInfo() {
TWorkloadSchedPolicy tPolicy = new TWorkloadSchedPolicy();
tPolicy.setId(id);
@@ -292,6 +275,10 @@ public class WorkloadSchedPolicy implements Writable,
GsonPostProcessable {
List<WorkloadAction> actionList = new ArrayList<>();
for (WorkloadActionMeta actionMeta : actionMetaList) {
try {
+ if
(WorkloadActionType.SET_SESSION_VARIABLE.equals(actionMeta.action)) {
+ Log.warn("skip deprecated workload policy action " +
actionMeta.action);
+ continue;
+ }
WorkloadAction ret =
WorkloadAction.createWorkloadAction(actionMeta);
actionList.add(ret);
} catch (UserException ue) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicyMgr.java
b/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicyMgr.java
index 07d4288b643..3704d0ad455 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicyMgr.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicyMgr.java
@@ -25,14 +25,11 @@ import org.apache.doris.common.io.Text;
import org.apache.doris.common.io.Writable;
import org.apache.doris.common.proc.BaseProcResult;
import org.apache.doris.common.proc.ProcResult;
-import org.apache.doris.common.util.DebugUtil;
-import org.apache.doris.common.util.MasterDaemon;
import org.apache.doris.mysql.privilege.PrivPredicate;
import org.apache.doris.persist.gson.GsonPostProcessable;
import org.apache.doris.persist.gson.GsonUtils;
import org.apache.doris.qe.ConnectContext;
import org.apache.doris.resource.Tag;
-import org.apache.doris.service.ExecuteEnv;
import org.apache.doris.thrift.TCompareOperator;
import org.apache.doris.thrift.TUserIdentity;
import org.apache.doris.thrift.TWorkloadActionType;
@@ -58,13 +55,12 @@ import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
-import java.util.PriorityQueue;
import java.util.Queue;
import java.util.Set;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.locks.ReentrantReadWriteLock;
-public class WorkloadSchedPolicyMgr extends MasterDaemon implements Writable,
GsonPostProcessable {
+public class WorkloadSchedPolicyMgr implements Writable, GsonPostProcessable {
private static final Logger LOG =
LogManager.getLogger(WorkloadSchedPolicyMgr.class);
@@ -74,10 +70,6 @@ public class WorkloadSchedPolicyMgr extends MasterDaemon
implements Writable, Gs
private PolicyProcNode policyProcNode = new PolicyProcNode();
- public WorkloadSchedPolicyMgr() {
- super("workload-sched-thread",
Config.workload_sched_policy_interval_ms);
- }
-
public static final ImmutableList<String>
WORKLOAD_SCHED_POLICY_NODE_TITLE_NAMES
= new ImmutableList.Builder<String>()
.add("Id").add("Name").add("Condition").add("Action").add("Priority").add("Enabled").add("Version")
@@ -92,13 +84,6 @@ public class WorkloadSchedPolicyMgr extends MasterDaemon
implements Writable, Gs
.put(WorkloadConditionOperator.LESS, TCompareOperator.LESS)
.put(WorkloadConditionOperator.LESS_EQUAl,
TCompareOperator.LESS_EQUAL).build();
- public static final ImmutableSet<WorkloadActionType> FE_ACTION_SET
- = new
ImmutableSet.Builder<WorkloadActionType>().add(WorkloadActionType.SET_SESSION_VARIABLE).build();
-
- public static final ImmutableSet<WorkloadMetricType> FE_METRIC_SET
- = new
ImmutableSet.Builder<WorkloadMetricType>().add(WorkloadMetricType.USERNAME)
- .build();
-
public static final ImmutableSet<WorkloadActionType> BE_ACTION_SET
= new
ImmutableSet.Builder<WorkloadActionType>().add(WorkloadActionType.MOVE_QUERY_TO_GROUP)
.add(WorkloadActionType.CANCEL_QUERY).build();
@@ -130,16 +115,9 @@ public class WorkloadSchedPolicyMgr extends MasterDaemon
implements Writable, Gs
public static final Map<String, WorkloadActionType> STRING_ACTION_MAP =
new HashMap<>();
static {
- for (WorkloadMetricType metricType : FE_METRIC_SET) {
- STRING_METRIC_MAP.put(metricType.toString(), metricType);
- }
for (WorkloadMetricType metricType : BE_METRIC_SET) {
STRING_METRIC_MAP.put(metricType.toString(), metricType);
}
-
- for (WorkloadActionType actionType : FE_ACTION_SET) {
- STRING_ACTION_MAP.put(actionType.toString(), actionType);
- }
for (WorkloadActionType actionType : BE_ACTION_SET) {
STRING_ACTION_MAP.put(actionType.toString(), actionType);
}
@@ -154,45 +132,6 @@ public class WorkloadSchedPolicyMgr extends MasterDaemon
implements Writable, Gs
}
};
- @Override
- protected void runAfterCatalogReady() {
- try {
- // todo(wb) add more query info source, not only comes from
connectionmap
- // 1 get query info map
- Map<Integer, ConnectContext> connectMap =
ExecuteEnv.getInstance().getScheduler()
- .getConnectionMap();
- List<WorkloadQueryInfo> queryInfoList = new ArrayList<>();
-
- // a snapshot for connect context
- Set<Integer> keySet = new HashSet<>();
- keySet.addAll(connectMap.keySet());
-
- for (Integer connectId : keySet) {
- ConnectContext cctx = connectMap.get(connectId);
- if (cctx == null || cctx.isKilled()) {
- continue;
- }
-
- String username = cctx.getQualifiedUser();
- WorkloadQueryInfo policyQueryInfo = new WorkloadQueryInfo();
- policyQueryInfo.queryId = cctx.queryId() == null ? null :
DebugUtil.printId(cctx.queryId());
- policyQueryInfo.tUniqueId = cctx.queryId();
- policyQueryInfo.context = cctx;
- policyQueryInfo.metricMap = new HashMap<>();
- policyQueryInfo.metricMap.put(WorkloadMetricType.USERNAME,
username);
-
- queryInfoList.add(policyQueryInfo);
- }
-
- // 2 exec policy
- if (queryInfoList.size() > 0) {
- execPolicy(queryInfoList);
- }
- } catch (Throwable t) {
- LOG.error("[policy thread]error happens when exec policy");
- }
- }
-
public void createWorkloadSchedPolicy(String policyName, boolean
isIfNotExists,
List<WorkloadConditionMeta> originConditions,
List<WorkloadActionMeta> originActions,
Map<String, String> propMap) throws UserException {
@@ -202,7 +141,7 @@ public class WorkloadSchedPolicyMgr extends MasterDaemon
implements Writable, Gs
WorkloadCondition cond =
WorkloadCondition.createWorkloadCondition(cm);
policyConditionList.add(cond);
}
- Boolean feCondition = checkPolicyCondition(policyConditionList);
+ checkPolicyCondition(policyConditionList);
// 2 create action
List<WorkloadAction> policyActionList = new ArrayList<>();
@@ -211,11 +150,7 @@ public class WorkloadSchedPolicyMgr extends MasterDaemon
implements Writable, Gs
WorkloadAction ret =
WorkloadAction.createWorkloadAction(workloadActionMeta);
policyActionList.add(ret);
}
-
- boolean feAction = checkPolicyAction(policyActionList);
- if (feCondition != null && feAction != feCondition) {
- throw new UserException("action and metric must run in FE together
or run in BE together");
- }
+ checkPolicyAction(policyActionList);
// 3 create policy
if (propMap == null) {
@@ -251,68 +186,22 @@ public class WorkloadSchedPolicyMgr extends MasterDaemon
implements Writable, Gs
}
}
- private Boolean checkPolicyCondition(List<WorkloadCondition>
conditionList) throws UserException {
+ private void checkPolicyCondition(List<WorkloadCondition> conditionList)
throws UserException {
if (conditionList.size() >
Config.workload_max_condition_num_in_policy) {
throw new UserException(
"condition num in a policy can not exceed " +
Config.workload_max_condition_num_in_policy);
}
- boolean hasFeOnlyMetric = false;
- boolean hasBeOnlyMetric = false;
- for (WorkloadCondition cond : conditionList) {
- boolean isFe = FE_METRIC_SET.contains(cond.getMetricType());
- boolean isBe = BE_METRIC_SET.contains(cond.getMetricType());
-
- if (isFe && !isBe) {
- hasFeOnlyMetric = true;
- } else if (isBe && !isFe) {
- hasBeOnlyMetric = true;
- }
-
- if (hasFeOnlyMetric && hasBeOnlyMetric) {
- throw new UserException(
- "one policy can not contains fe only and be only
metric, FE metric list is " + FE_METRIC_SET
- + ", BE metric list is " + BE_METRIC_SET);
- }
- }
- if (hasFeOnlyMetric) {
- return true;
- } else if (hasBeOnlyMetric) {
- return false;
- } else {
- return null;
- }
}
- private boolean checkPolicyAction(List<WorkloadAction> actionList) throws
UserException {
+ private void checkPolicyAction(List<WorkloadAction> actionList) throws
UserException {
if (actionList.size() > Config.workload_max_action_num_in_policy) {
throw new UserException(
"action num in one policy can not exceed " +
Config.workload_max_action_num_in_policy);
}
Set<WorkloadActionType> actionTypeSet = new HashSet<>();
- Set<String> setSessionVarSet = new HashSet<>();
- boolean containsFeAction = false;
- boolean containsBeAction = false;
for (WorkloadAction action : actionList) {
- if (FE_ACTION_SET.contains(action.getWorkloadActionType())) {
- containsFeAction = true;
- }
- if (BE_ACTION_SET.contains(action.getWorkloadActionType())) {
- containsBeAction = true;
- }
- if (containsFeAction && containsBeAction) {
- throw new UserException(
- "one policy can not contains fe and be action, FE
action list is " + FE_ACTION_SET
- + ", BE action list is " + BE_ACTION_SET);
- }
- // set session var cmd can be duplicate, but args can not be
duplicate
- if
(action.getWorkloadActionType().equals(WorkloadActionType.SET_SESSION_VARIABLE))
{
- WorkloadActionSetSessionVar setAction =
(WorkloadActionSetSessionVar) action;
- if (!setSessionVarSet.add(setAction.getVarName())) {
- throw new UserException(
- "duplicate set_session_variable action args one
policy, " + setAction.getVarName());
- }
- } else if (!actionTypeSet.add(action.getWorkloadActionType())) {
+ if (!actionTypeSet.add(action.getWorkloadActionType())) {
throw new UserException("duplicate action in one policy");
}
}
@@ -322,59 +211,6 @@ public class WorkloadSchedPolicyMgr extends MasterDaemon
implements Writable, Gs
throw new UserException(String.format("%s and %s can not exist in
one policy at same time",
WorkloadActionType.CANCEL_QUERY,
WorkloadActionType.MOVE_QUERY_TO_GROUP));
}
- return containsFeAction;
- }
-
- public void execPolicy(List<WorkloadQueryInfo> queryInfoList) {
- // 1 get a snapshot of policy
- Set<Long> policyIdSet = new HashSet<>();
- readLock();
- try {
- for (Map.Entry<Long, WorkloadSchedPolicy> entry :
idToPolicy.entrySet()) {
- if (entry.getValue().isFePolicy()) {
- policyIdSet.add(entry.getKey());
- }
- }
- } finally {
- readUnlock();
- }
-
- for (WorkloadQueryInfo queryInfo : queryInfoList) {
- try {
- // 1 check policy is match
- Map<WorkloadActionType, Queue<WorkloadSchedPolicy>>
matchedPolicyMap = Maps.newHashMap();
- for (Long policyId : policyIdSet) {
- WorkloadSchedPolicy policy = idToPolicy.get(policyId);
- if (policy == null) {
- continue;
- }
- if (policy.isEnabled() && policy.isMatch(queryInfo)) {
- WorkloadActionType actionType =
policy.getFirstActionType();
- // add to priority queue
- Queue<WorkloadSchedPolicy> queue =
matchedPolicyMap.get(actionType);
- if (queue == null) {
- queue = new PriorityQueue<>(policyComparator);
- matchedPolicyMap.put(actionType, queue);
- }
- queue.offer(policy);
- }
- }
-
- if (matchedPolicyMap.size() == 0) {
- continue;
- }
-
- // 2 pick higher priority policy when action conflicts
- List<WorkloadSchedPolicy> pickedPolicyList =
pickPolicy(matchedPolicyMap);
-
- // 3 exec action
- for (WorkloadSchedPolicy policy : pickedPolicyList) {
- policy.execAction(queryInfo);
- }
- } catch (Throwable e) {
- LOG.warn("exec policy with query {} failed ",
queryInfo.queryId, e);
- }
- }
}
List<WorkloadSchedPolicy> pickPolicy(Map<WorkloadActionType,
Queue<WorkloadSchedPolicy>> policyMap) {
@@ -579,9 +415,6 @@ public class WorkloadSchedPolicyMgr extends MasterDaemon
implements Writable, Gs
readLock();
try {
for (Map.Entry<Long, WorkloadSchedPolicy> entry :
idToPolicy.entrySet()) {
- if (entry.getValue().isFePolicy()) {
- continue;
- }
TopicInfo tInfo = entry.getValue().toTopicInfo();
if (tInfo != null) {
topicInfoList.add(tInfo);
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicyMgrTest.java
b/fe/fe-core/src/test/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicyMgrTest.java
index 6d4b6e2e6c3..7d68d77f276 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicyMgrTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/resource/workloadschedpolicy/WorkloadSchedPolicyMgrTest.java
@@ -91,20 +91,7 @@ public class WorkloadSchedPolicyMgrTest {
Assert.fail("Should not throw exception for mixed USERNAME and BE
metrics: " + e.getMessage());
}
- // Case 2: USERNAME (Shared) + FE Action -> OK
- try {
- List<WorkloadConditionMeta> conditionMetas = new ArrayList<>();
- conditionMetas.add(new WorkloadConditionMeta("username", "=",
"user1"));
-
- List<WorkloadActionMeta> actionMetas = new ArrayList<>();
- actionMetas.add(new WorkloadActionMeta("set_session_variable",
"workload_group=normal"));
-
- mgr.createWorkloadSchedPolicy("policy_fe_only", false,
conditionMetas, actionMetas, null);
- } catch (UserException e) {
- Assert.fail("Should not throw exception for USERNAME + FE Action:
" + e.getMessage());
- }
-
- // Case 3: USERNAME (Shared) + BE Action -> OK
+ // Case 2: USERNAME (Shared) + BE Action -> OK
try {
List<WorkloadConditionMeta> conditionMetas = new ArrayList<>();
conditionMetas.add(new WorkloadConditionMeta("username", "=",
"user1"));
@@ -116,34 +103,15 @@ public class WorkloadSchedPolicyMgrTest {
} catch (UserException e) {
Assert.fail("Should not throw exception for USERNAME + BE Action:
" + e.getMessage());
}
+ }
- // Case 4: BE Metric + FE Action -> Error
- try {
- List<WorkloadConditionMeta> conditionMetas = new ArrayList<>();
- conditionMetas.add(new WorkloadConditionMeta("query_time", ">",
"1000"));
-
- List<WorkloadActionMeta> actionMetas = new ArrayList<>();
- actionMetas.add(new WorkloadActionMeta("set_session_variable",
"workload_group=normal"));
-
- mgr.createWorkloadSchedPolicy("policy_error_1", false,
conditionMetas, actionMetas, null);
- Assert.fail("Should throw exception for BE Metric + FE Action");
- } catch (UserException e) {
- Assert.assertTrue(e.getMessage().contains("action and metric must
run in FE together or run in BE together"));
- }
-
- // Case 5: USERNAME + BE Metric + FE Action -> Error
+ @Test
+ public void testSetSessionVariableActionIsRejected() {
try {
- List<WorkloadConditionMeta> conditionMetas = new ArrayList<>();
- conditionMetas.add(new WorkloadConditionMeta("username", "=",
"user1"));
- conditionMetas.add(new WorkloadConditionMeta("query_time", ">",
"1000"));
-
- List<WorkloadActionMeta> actionMetas = new ArrayList<>();
- actionMetas.add(new WorkloadActionMeta("set_session_variable",
"workload_group=normal"));
-
- mgr.createWorkloadSchedPolicy("policy_error_2", false,
conditionMetas, actionMetas, null);
- Assert.fail("Should throw exception for USERNAME + BE Metric + FE
Action");
+ new WorkloadActionMeta("set_session_variable",
"workload_group=normal");
+ Assert.fail("Should throw exception for removed
set_session_variable action");
} catch (UserException e) {
- Assert.assertTrue(e.getMessage().contains("action and metric must
run in FE together or run in BE together"));
+ Assert.assertTrue(e.getMessage().contains("invalid action type
set_session_variable"));
}
}
@@ -372,21 +340,4 @@ public class WorkloadSchedPolicyMgrTest {
Assert.assertTrue(e.getMessage().contains("remote scan bytes"));
}
}
-
- @Test
- public void testRemoteScanBytesMetricCanNotMixWithFeAction() throws
UserException {
- try {
- List<WorkloadConditionMeta> conditionMetas = new ArrayList<>();
- // Validate the new metric follows the existing BE-only action
compatibility rules.
- conditionMetas.add(new
WorkloadConditionMeta("be_scan_bytes_from_remote_storage", ">", "100"));
- List<WorkloadActionMeta> actionMetas = new ArrayList<>();
- actionMetas.add(new WorkloadActionMeta("set_session_variable",
"workload_group=normal"));
-
-
mgr.createWorkloadSchedPolicy("policy_remote_scan_bytes_with_fe_action", false,
conditionMetas,
- actionMetas, null);
- Assert.fail("Should throw exception for remote scan bytes metric
with FE action");
- } catch (UserException e) {
- Assert.assertTrue(e.getMessage().contains("action and metric must
run in FE together or run in BE together"));
- }
- }
}
diff --git
a/regression-test/data/workload_manager_p0/test_workload_sched_policy.out
b/regression-test/data/workload_manager_p0/test_workload_sched_policy.out
index 3152367e9a1..ed49b9675fc 100644
--- a/regression-test/data/workload_manager_p0/test_workload_sched_policy.out
+++ b/regression-test/data/workload_manager_p0/test_workload_sched_policy.out
@@ -1,24 +1,22 @@
-- This file is automatically generated. You should know what you did if you
want to edit this
-- !select_policy_tvf --
be_policy query_time > 10 cancel_query 10 false 0
-fe_policy username = root set_session_variable "workload_group=normal"
10 false 0
-set_action_policy username = root set_session_variable
"workload_group=normal" 0 false 0
test_cancel_policy query_time > 10 cancel_query 0 false 0
-- !select_policy_tvf_after_drop --
-- !select_alter_1 --
-test_alter_policy username = test_alter_policy_user
set_session_variable "parallel_pipeline_task_num=0" 0 true 0
normal
+test_alter_policy username = test_alter_policy_user cancel_query
0 true 0 normal
-- !select_alter_2 --
-test_alter_policy username = test_alter_policy_user
set_session_variable "parallel_pipeline_task_num=0" 0 true 1
+test_alter_policy username = test_alter_policy_user cancel_query
0 true 1
-- !select_alter_3 --
-test_alter_policy username = test_alter_policy_user
set_session_variable "parallel_pipeline_task_num=0" 0 false 2
+test_alter_policy username = test_alter_policy_user cancel_query
0 false 2
-- !select_alter_4 --
-test_alter_policy username = test_alter_policy_user
set_session_variable "parallel_pipeline_task_num=0" 9 false 3
+test_alter_policy username = test_alter_policy_user cancel_query
9 false 3
-- !select_alter_5 --
-test_alter_policy username = test_alter_policy_user
set_session_variable "parallel_pipeline_task_num=0" 9 false 4
normal
+test_alter_policy username = test_alter_policy_user cancel_query
9 false 4 normal
diff --git
a/regression-test/suites/workload_manager_p0/test_workload_sched_policy.groovy
b/regression-test/suites/workload_manager_p0/test_workload_sched_policy.groovy
index b205ce9dda0..55e198df2b6 100644
---
a/regression-test/suites/workload_manager_p0/test_workload_sched_policy.groovy
+++
b/regression-test/suites/workload_manager_p0/test_workload_sched_policy.groovy
@@ -30,8 +30,6 @@ suite("test_workload_sched_policy") {
}
sql "drop workload policy if exists test_cancel_policy;"
- sql "drop workload policy if exists set_action_policy;"
- sql "drop workload policy if exists fe_policy;"
sql "drop workload policy if exists be_policy;"
sql "drop workload policy if exists be_scan_row_policy;"
sql "drop workload policy if exists be_scan_bytes_policy;"
@@ -42,21 +40,7 @@ suite("test_workload_sched_policy") {
" conditions(query_time > 10) " +
" actions(cancel_query) properties('enabled'='false'); "
- // 2 create set policy
- sql "create workload policy set_action_policy " +
- "conditions(username='root') " +
- "actions(set_session_variable 'workload_group=normal')
properties('enabled'='false');"
-
- // 3 create policy run in fe
- sql "create workload policy fe_policy " +
- "conditions(username='root') " +
- "actions(set_session_variable 'workload_group=normal') " +
- "properties( " +
- "'enabled' = 'false', " +
- "'priority'='10' " +
- ");"
-
- // 4 create policy run in be
+ // 2 create policy run in be
sql "create workload policy be_policy " +
"conditions(query_time > 10) " +
"actions(cancel_query) " +
@@ -65,10 +49,10 @@ suite("test_workload_sched_policy") {
"'priority'='10' " +
");"
- qt_select_policy_tvf "select
name,condition,action,priority,enabled,version from
information_schema.workload_policy where name
in('be_policy','fe_policy','set_action_policy','test_cancel_policy') order by
name;"
+ qt_select_policy_tvf "select
name,condition,action,priority,enabled,version from
information_schema.workload_policy where name
in('be_policy','test_cancel_policy') order by name;"
// test_alter
- sql "alter workload policy fe_policy properties('priority'='2',
'enabled'='false');"
+ sql "alter workload policy be_policy properties('priority'='2',
'enabled'='false');"
// create failed check
test {
@@ -86,19 +70,19 @@ suite("test_workload_sched_policy") {
}
test {
- sql "alter workload policy fe_policy properties('priority'='abc');"
+ sql "alter workload policy be_policy properties('priority'='abc');"
exception "invalid priority property value"
}
test {
- sql "alter workload policy fe_policy properties('enabled'='abc');"
+ sql "alter workload policy be_policy properties('enabled'='abc');"
exception "invalid enabled property value"
}
test {
- sql "alter workload policy fe_policy properties('priority'='10000');"
+ sql "alter workload policy be_policy properties('priority'='10000');"
exception "priority can only between"
}
@@ -112,11 +96,11 @@ suite("test_workload_sched_policy") {
}
test {
- sql "create workload policy conflict_policy " +
+ sql "create workload policy invalid_set_session_policy " +
"conditions (username = 'root') " +
- "actions(set_session_variable 'workload_group=normal',
set_session_variable 'workload_group=abc');"
+ "actions(set_session_variable 'workload_group=normal');"
- exception "duplicate set_session_variable action args one policy"
+ exception "invalid action type set_session_variable"
}
test {
@@ -145,66 +129,16 @@ suite("test_workload_sched_policy") {
// drop
sql "drop workload policy test_cancel_policy;"
- sql "drop workload policy set_action_policy;"
- sql "drop workload policy fe_policy;"
sql "drop workload policy be_policy;"
sql "drop workload policy be_scan_row_policy;"
sql "drop workload policy be_scan_bytes_policy;"
sql "drop workload policy query_be_memory_used;"
- qt_select_policy_tvf_after_drop "select
name,condition,action,priority,enabled,version from
information_schema.workload_policy where name
in('be_policy','fe_policy','set_action_policy','test_cancel_policy') order by
name;"
-
- // test workload policy
- sql """drop user if exists test_workload_sched_user"""
- sql """create user test_workload_sched_user identified by '12345'"""
- sql """grant ADMIN_PRIV on *.*.* to test_workload_sched_user"""
- sql "drop workload group if exists test_set_session_wg
$forComputeGroupStr;"
- sql "drop workload group if exists test_set_session_wg2
$forComputeGroupStr;"
- sql "create workload group test_set_session_wg $forComputeGroupStr
properties('min_cpu_percent'='0');"
- sql "create workload group test_set_session_wg2 $forComputeGroupStr
properties('min_cpu_percent'='0');"
-
- sql "drop workload policy if exists test_set_var_policy;"
- sql "drop workload policy if exists test_set_var_policy2;"
-
- // 1 create test_set_var_policy
- sql "create workload policy test_set_var_policy
conditions(username='test_workload_sched_user')" +
- "actions(set_session_variable
'workload_group=test_set_session_wg');"
- def result1 = connect('test_workload_sched_user', '12345',
context.config.jdbcUrl) {
- logger.info("begin sleep 15s to wait")
- Thread.sleep(15000)
- sql "show variables like 'workload_group';"
- }
- assertEquals("workload_group", result1[0][0])
- assertEquals("test_set_session_wg", result1[0][1])
-
- // 2 create test_set_var_policy2 with higher priority
- sql "create workload policy test_set_var_policy2
conditions(username='test_workload_sched_user') " +
- "actions(set_session_variable
'workload_group=test_set_session_wg2') properties('priority'='10');"
- def result2 = connect('test_workload_sched_user', '12345',
context.config.jdbcUrl) {
- Thread.sleep(3000)
- sql "show variables like 'workload_group';"
- }
- assertEquals("workload_group", result2[0][0])
- assertEquals("test_set_session_wg2", result2[0][1])
-
- // 3 disable test_set_var_policy2
- sql "alter workload policy test_set_var_policy2
properties('enabled'='false');"
- def result3 = connect('test_workload_sched_user', '12345',
context.config.jdbcUrl) {
- Thread.sleep(3000)
- sql "show variables like 'workload_group';"
- }
- assertEquals("workload_group", result3[0][0])
- assertEquals("test_set_session_wg", result3[0][1])
- sql "drop workload group if exists test_set_session_wg
$forComputeGroupStr;"
- sql "drop workload group if exists test_set_session_wg2
$forComputeGroupStr;"
-
- sql "drop workload policy if exists test_set_var_policy;"
- sql "drop workload policy if exists test_set_var_policy2;"
+ qt_select_policy_tvf_after_drop "select
name,condition,action,priority,enabled,version from
information_schema.workload_policy where name
in('be_policy','test_cancel_policy') order by name;"
sql "drop user if exists test_policy_user"
sql "drop workload policy if exists test_cancel_query_policy"
sql "drop workload policy if exists test_cancel_query_policy2"
- sql "drop workload policy if exists test_set_session"
sql "drop workload group if exists policy_group $forComputeGroupStr;"
sql "CREATE USER 'test_policy_user'@'%' IDENTIFIED BY '12345';"
sql """grant SELECT_PRIV on *.*.* to test_policy_user;"""
@@ -214,7 +148,6 @@ suite("test_workload_sched_policy") {
sql "GRANT USAGE_PRIV ON WORKLOAD GROUP 'policy_group2' TO
'test_policy_user'@'%';"
sql "create workload policy test_cancel_query_policy conditions(query_time
> 1000) actions(cancel_query)
properties('workload_group'='${currentCgName}policy_group')"
sql "create workload policy test_cancel_query_policy2
conditions(query_time > 0, be_scan_rows>1) actions(cancel_query)
properties('workload_group'='${currentCgName}policy_group')"
- sql "create workload policy test_set_session
conditions(username='test_policy_user') actions(set_session_variable
'parallel_pipeline_task_num=1')"
test {
sql "drop workload group policy_group $forComputeGroupStr;"
@@ -230,7 +163,7 @@ suite("test_workload_sched_policy") {
sql "drop user if exists test_alter_policy_user"
sql "CREATE USER 'test_alter_policy_user'@'%' IDENTIFIED BY '12345';"
sql "drop workload policy if exists test_alter_policy;"
- sql "create workload policy test_alter_policy
conditions(username='test_alter_policy_user') actions(set_session_variable
'parallel_pipeline_task_num=0')
properties('workload_group'='${currentCgName}normal');"
+ sql "create workload policy test_alter_policy
conditions(username='test_alter_policy_user') actions(cancel_query)
properties('workload_group'='${currentCgName}normal');"
qt_select_alter_1 "select
name,condition,action,PRIORITY,ENABLED,VERSION,WORKLOAD_GROUP from
information_schema.workload_policy where name='test_alter_policy'"
sql "alter workload policy test_alter_policy
properties('workload_group'='');"
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]