This is an automated email from the ASF dual-hosted git repository.
jacktengg pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 157e03eeb16 [fix](fe) Bound session execution and scan concurrency
(#68281)
157e03eeb16 is described below
commit 157e03eeb16deb8a84ae58d9bed721f56f8fedb3
Author: TengJianPing <[email protected]>
AuthorDate: Fri Oct 9 16:12:59 2026 +0800
[fix](fe) Bound session execution and scan concurrency (#68281)
Issue Number: None
Related PR: None
Problem Summary: Setting parallel_pipeline_task_num to 2147483647 passes
its lower-bound-only validation and can create excessive BE pipeline
tasks.
Reject larger values before changing the session variable.
Also check upper bounds of colocate parallelism, four scanner
concurrency variables, parallel scan scanner counts, send-batch
parallelism,
and load streams per node. These values control instance creation,
scanner
splitting/scheduling, sender concurrency, or direct load-stream
allocation.
Preserve existing non-positive fallback values, and require positive
colocate and load-stream counts.
Reuse the existing integer validation and VariableMgr setter/checker
paths
so SET, SET GLOBAL, and SET_VAR hints reject excessive values.
Colocate parallelism and load streams per node must
also be positive. Automatic pipeline and scanner fallback values remain
supported.
- Behavior changed: Yes; reject excessive concurrency settings and
non-positive colocate/load-stream counts during variable assignment.
- Does this need documentation: No; configuration descriptions are
included.
### What problem does this PR solve?
Issue Number: close #xxx
Related PR: #xxx
Problem Summary:
### Release note
None
### Check List (For Author)
- Test <!-- At least one of them must be included. -->
- [ ] Regression test
- [ ] Unit Test
- [ ] Manual test (add detailed scripts or steps below)
- [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change.
- [ ] No code files have been changed.
- [ ] Other reason <!-- Add your reason? -->
- Behavior changed:
- [ ] No.
- [ ] Yes. <!-- Explain the behavior change -->
- Does this need documentation?
- [ ] No.
- [ ] Yes. <!-- Add document PR link here. eg:
https://github.com/apache/doris-website/pull/1214 -->
### Check List (For Reviewer who merge this PR)
- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->
---
.../apache/doris/mysql/privilege/UserProperty.java | 5 +
.../doris/nereids/load/NereidsStreamLoadTask.java | 6 +
.../java/org/apache/doris/qe/SessionVariable.java | 100 ++++++++--
.../java/org/apache/doris/qe/StmtExecutor.java | 12 +-
.../main/java/org/apache/doris/qe/VariableMgr.java | 6 +-
.../org/apache/doris/catalog/UserPropertyTest.java | 26 ++-
.../nereids/load/NereidsStreamLoadTaskTest.java | 68 +++++++
.../doris/qe/SessionVariableParallelismTest.java | 217 +++++++++++++++++++++
.../org/apache/doris/qe/SessionVariablesTest.java | 27 +++
.../test_stream_load_send_batch_parallelism.out | 52 +++++
.../test_session_concurrency_limits.out | 112 +++++++++++
.../test_user_parallelism_limit.out | 22 +++
.../test_stream_load_send_batch_parallelism.groovy | 52 +++++
.../test_session_concurrency_limits.groovy | 83 ++++++++
.../test_user_parallelism_limit.groovy | 39 ++++
15 files changed, 800 insertions(+), 27 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/mysql/privilege/UserProperty.java
b/fe/fe-core/src/main/java/org/apache/doris/mysql/privilege/UserProperty.java
index 21b3f9c49f0..019f87103b8 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/mysql/privilege/UserProperty.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/mysql/privilege/UserProperty.java
@@ -253,6 +253,11 @@ public class UserProperty {
} catch (NumberFormatException e) {
throw new
DdlException(PROP_PARALLEL_FRAGMENT_EXEC_INSTANCE_NUM + " is not number");
}
+ // Replay must accept historical values; SessionVariable caps
their effective parallelism.
+ if (!isReplay && newParallelFragmentExecInstanceNum > 256) {
+ throw new
DdlException(PROP_PARALLEL_FRAGMENT_EXEC_INSTANCE_NUM
+ + " must be less than or equal to 256, got " +
value);
+ }
} else if (keyArr[0].equalsIgnoreCase(PROP_SQL_BLOCK_RULES)) {
// set property "sql_block_rules" = "test_rule1,test_rule2"
if (keyArr.length != 1) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/nereids/load/NereidsStreamLoadTask.java
b/fe/fe-core/src/main/java/org/apache/doris/nereids/load/NereidsStreamLoadTask.java
index f5ddca41f19..be92b7cde06 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/nereids/load/NereidsStreamLoadTask.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/nereids/load/NereidsStreamLoadTask.java
@@ -30,6 +30,7 @@ import org.apache.doris.nereids.analyzer.UnboundSlot;
import org.apache.doris.nereids.trees.expressions.BinaryOperator;
import org.apache.doris.nereids.trees.expressions.Expression;
import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.SessionVariable;
import org.apache.doris.task.LoadTaskInfo;
import org.apache.doris.thrift.TFileCompressType;
import org.apache.doris.thrift.TFileFormatType;
@@ -474,6 +475,11 @@ public class NereidsStreamLoadTask implements
NereidsLoadTaskInfo {
sequenceCol = request.getSequenceCol();
}
if (request.isSetSendBatchParallelism()) {
+ try {
+
SessionVariable.checkSendBatchParallelism(Integer.toString(request.getSendBatchParallelism()));
+ } catch (Exception e) {
+ throw new UserException(e.getMessage(), e);
+ }
sendBatchParallelism = request.getSendBatchParallelism();
}
if (request.isSetMaxFilterRatio()) {
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
index 6fb6656ce76..b7654c9051a 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
@@ -1157,11 +1157,13 @@ public class SessionVariable implements Serializable,
Writable {
// 100MB
public long maxScanQueueMemByte = 2147483648L / 20;
- @VarAttrDef.VarAttr(name = MAX_SCANNERS_CONCURRENCY, needForward = true,
description = "The max threads to read "
+ @VarAttrDef.VarAttr(name = MAX_SCANNERS_CONCURRENCY, needForward = true,
+ checker = "checkMaxScannersConcurrency", description = "The max
threads to read "
+ "data of ScanNode, default 4")
public int maxScannersConcurrency = 4;
- @VarAttrDef.VarAttr(name = MAX_FILE_SCANNERS_CONCURRENCY, needForward =
true, description = "The max threads to "
+ @VarAttrDef.VarAttr(name = MAX_FILE_SCANNERS_CONCURRENCY, needForward =
true,
+ checker = "checkMaxFileScannersConcurrency", description = "The
max threads to "
+ "read data of FileScanNode, default 16")
public int maxFileScannersConcurrency = 16;
@@ -1172,11 +1174,13 @@ public class SessionVariable implements Serializable,
Writable {
@VarAttrDef.VarAttr(name = LOCAL_EXCHANGE_FREE_BLOCKS_LIMIT)
public int localExchangeFreeBlocksLimit = 4;
- @VarAttrDef.VarAttr(name = MIN_SCANNERS_CONCURRENCY, needForward = true,
description = "The min concurrency of "
+ @VarAttrDef.VarAttr(name = MIN_SCANNERS_CONCURRENCY, needForward = true,
+ checker = "checkMinScannersConcurrency", description = "The min
concurrency of "
+ "Scanner, default 1")
public int minScannersConcurrency = 1;
- @VarAttrDef.VarAttr(name = MIN_FILE_SCANNERS_CONCURRENCY, needForward =
true, description = "The min concurrency "
+ @VarAttrDef.VarAttr(name = MIN_FILE_SCANNERS_CONCURRENCY, needForward =
true,
+ checker = "checkMinFileScannersConcurrency", description = "The
min concurrency "
+ "of Remote Scanner, default 1")
public int minFileScannersConcurrency = 1;
@@ -1459,7 +1463,8 @@ public class SessionVariable implements Serializable,
Writable {
setter = "setFragmentInstanceNum", varType =
VariableAnnotation.DEPRECATED)
public int parallelExecInstanceNum = 8;
- @VarAttrDef.VarAttr(name = COLOCATE_MAX_PARALLEL_NUM, needForward = true,
fuzzy = false)
+ @VarAttrDef.VarAttr(name = COLOCATE_MAX_PARALLEL_NUM, needForward = true,
fuzzy = false,
+ checker = "checkColocateMaxParallelNum")
public int colocateMaxParallelNum = 128;
@VarAttrDef.VarAttr(name = PARALLEL_PIPELINE_TASK_NUM, fuzzy = true,
needForward = true,
@@ -1646,7 +1651,7 @@ public class SessionVariable implements Serializable,
Writable {
@VarAttrDef.VarAttr(name = DELETE_WITHOUT_PARTITION, needForward = true)
public boolean deleteWithoutPartition = false;
- @VarAttrDef.VarAttr(name = SEND_BATCH_PARALLELISM, needForward = true)
+ @VarAttrDef.VarAttr(name = SEND_BATCH_PARALLELISM, needForward = true,
checker = "checkSendBatchParallelism")
public int sendBatchParallelism = 1;
@VarAttrDef.VarAttr(name = ENABLE_NEREIDS_DML, varType =
VariableAnnotation.REMOVED)
@@ -1685,7 +1690,8 @@ public class SessionVariable implements Serializable,
Writable {
private boolean optimizeIndexScanParallelism = true;
@VarAttrDef.VarAttr(name = PARALLEL_SCAN_MAX_SCANNERS_COUNT, fuzzy = true,
- varType = VariableAnnotation.EXPERIMENTAL, needForward = true)
+ varType = VariableAnnotation.EXPERIMENTAL, needForward = true,
+ checker = "checkParallelScanMaxScannersCount")
private int parallelScanMaxScannersCount = 0;
@VarAttrDef.VarAttr(name = PARALLEL_SCAN_MIN_ROWS_PER_SCANNER, fuzzy =
true,
@@ -2699,7 +2705,7 @@ public class SessionVariable implements Serializable,
Writable {
@VarAttrDef.VarAttr(name = ENABLE_MEMTABLE_ON_SINK_NODE, needForward =
true)
public boolean enableMemtableOnSinkNode = true;
- @VarAttrDef.VarAttr(name = LOAD_STREAM_PER_NODE)
+ @VarAttrDef.VarAttr(name = LOAD_STREAM_PER_NODE, checker =
"checkLoadStreamPerNode")
public int loadStreamPerNode = 2;
@VarAttrDef.VarAttr(name = GROUP_COMMIT, needForward = true)
@@ -4172,10 +4178,45 @@ public class SessionVariable implements Serializable,
Writable {
}
public void setPipelineTaskNum(String value) throws Exception {
- int val = checkFieldValue(PARALLEL_PIPELINE_TASK_NUM, 0, value);
+ int val = checkFieldValue(PARALLEL_PIPELINE_TASK_NUM, 0, 256, value);
this.parallelPipelineTaskNum = val;
}
+ public void checkColocateMaxParallelNum(String value) throws Exception {
+ checkFieldValue(COLOCATE_MAX_PARALLEL_NUM, 1, 256, value);
+ }
+
+ public void checkMaxScannersConcurrency(String value) throws Exception {
+ // Non-positive scanner concurrency values select the BE defaults.
+ checkFieldValue(MAX_SCANNERS_CONCURRENCY, Integer.MIN_VALUE, 256,
value);
+ }
+
+ public void checkMaxFileScannersConcurrency(String value) throws Exception
{
+ checkFieldValue(MAX_FILE_SCANNERS_CONCURRENCY, Integer.MIN_VALUE, 256,
value);
+ }
+
+ public void checkMinScannersConcurrency(String value) throws Exception {
+ checkFieldValue(MIN_SCANNERS_CONCURRENCY, Integer.MIN_VALUE, 256,
value);
+ }
+
+ public void checkMinFileScannersConcurrency(String value) throws Exception
{
+ checkFieldValue(MIN_FILE_SCANNERS_CONCURRENCY, Integer.MIN_VALUE, 256,
value);
+ }
+
+ public void checkParallelScanMaxScannersCount(String value) throws
Exception {
+ // Non-positive values select the number of CPU cores on the BE.
+ checkFieldValue(PARALLEL_SCAN_MAX_SCANNERS_COUNT, Integer.MIN_VALUE,
256, value);
+ }
+
+ public static void checkSendBatchParallelism(String value) throws
Exception {
+ // The tablet writer uses one sender for values less than or equal to
one.
+ checkFieldValue(SEND_BATCH_PARALLELISM, Integer.MIN_VALUE, 256, value);
+ }
+
+ public void checkLoadStreamPerNode(String value) throws Exception {
+ checkFieldValue(LOAD_STREAM_PER_NODE, 1, 256, value);
+ }
+
public void setEnablePipelineEngine(String value) throws Exception {
if (value.equalsIgnoreCase("ON")
|| value.equalsIgnoreCase("TRUE")
@@ -4242,7 +4283,7 @@ public class SessionVariable implements Serializable,
Writable {
return val;
}
- private int checkFieldValue(String variableName, int minValue, String
value) throws Exception {
+ private static int checkFieldValue(String variableName, int minValue,
String value) throws Exception {
int val = Integer.valueOf(value);
if (val < minValue) {
throw new Exception(
@@ -4252,6 +4293,15 @@ public class SessionVariable implements Serializable,
Writable {
return val;
}
+ private static int checkFieldValue(String variableName, int minValue, int
maxValue, String value) throws Exception {
+ int val = checkFieldValue(variableName, minValue, value);
+ if (val > maxValue) {
+ throw new Exception(variableName + " value should less than or
equal " + maxValue
+ + ", you set value is: " + value);
+ }
+ return val;
+ }
+
public String getWorkloadGroup() {
return workloadGroup;
}
@@ -4392,6 +4442,11 @@ public class SessionVariable implements Serializable,
Writable {
}
public int getParallelExecInstanceNum(String clusterName) {
+ // Apply the limit after resolving user properties, automatic sizing,
and historical session values.
+ return Math.min(resolveParallelExecInstanceNum(clusterName), 256);
+ }
+
+ private int resolveParallelExecInstanceNum(String clusterName) {
ConnectContext connectContext = ConnectContext.get();
if (connectContext != null && connectContext.getEnv() != null &&
connectContext.getEnv().getAuth() != null) {
int userParallelExecInstanceNum = connectContext.getEnv().getAuth()
@@ -5846,12 +5901,29 @@ public class SessionVariable implements Serializable,
Writable {
}
}
- private static int normalizeIntValue(String name, String value) {
+ static int normalizeIntValue(String name, String value) {
int intValue = Integer.valueOf(value);
- if (RUNTIME_FILTER_TYPE.equalsIgnoreCase(name)) {
- return (int)
RuntimeFilterTypeHelper.normalizeDeprecatedRuntimeFilterTypes(intValue);
+ switch (name.toLowerCase(Locale.ROOT)) {
+ case RUNTIME_FILTER_TYPE:
+ return (int)
RuntimeFilterTypeHelper.normalizeDeprecatedRuntimeFilterTypes(intValue);
+ // Historical and forwarded values bypass statement-time
validation. Keep them valid
+ // before SET_VAR saves the original value for restoration through
the strict setters.
+ case PARALLEL_PIPELINE_TASK_NUM:
+ return Math.max(0, Math.min(intValue, 256));
+ case COLOCATE_MAX_PARALLEL_NUM:
+ case LOAD_STREAM_PER_NODE:
+ return Math.max(1, Math.min(intValue, 256));
+ case MAX_SCANNERS_CONCURRENCY:
+ case MAX_FILE_SCANNERS_CONCURRENCY:
+ case MIN_SCANNERS_CONCURRENCY:
+ case MIN_FILE_SCANNERS_CONCURRENCY:
+ case PARALLEL_SCAN_MAX_SCANNERS_COUNT:
+ case SEND_BATCH_PARALLELISM:
+ // Non-positive values retain their existing BE
default-selection semantics.
+ return Math.min(intValue, 256);
+ default:
+ return intValue;
}
- return intValue;
}
/**
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
index 5d5d4e6b53f..3e100238f03 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
@@ -794,12 +794,12 @@ public class StmtExecutor {
// revert Session Value
try {
VariableMgr.revertSessionValue(sessionVariable);
- // origin value init
- sessionVariable.setIsSingleSetVar(false);
- sessionVariable.clearSessionOriginValue();
} catch (DdlException e) {
LOG.warn("failed to revert Session value. {}",
context.getQueryIdentifier(), e);
context.getState().setError(e.getMysqlErrorCode(),
e.getMessage());
+ } finally {
+ sessionVariable.setIsSingleSetVar(false);
+ sessionVariable.clearSessionOriginValue();
}
}
}
@@ -2323,12 +2323,12 @@ public class StmtExecutor {
// revert Session Value
try {
VariableMgr.revertSessionValue(sessionVariable);
- // origin value init
- sessionVariable.setIsSingleSetVar(false);
- sessionVariable.clearSessionOriginValue();
} catch (DdlException e) {
LOG.warn("failed to revert Session value. {}",
context.getQueryIdentifier(), e);
context.getState().setError(e.getMysqlErrorCode(),
e.getMessage());
+ } finally {
+ sessionVariable.setIsSingleSetVar(false);
+ sessionVariable.clearSessionOriginValue();
}
}
}
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/VariableMgr.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/VariableMgr.java
index 144ab256416..3531f1565cb 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/VariableMgr.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/VariableMgr.java
@@ -249,11 +249,7 @@ public class VariableMgr {
field.setShort(obj, Short.parseShort(value));
break;
case "int":
- int intValue = Integer.parseInt(value);
- if
(SessionVariable.RUNTIME_FILTER_TYPE.equalsIgnoreCase(name)) {
- intValue = (int)
RuntimeFilterTypeHelper.normalizeDeprecatedRuntimeFilterTypes(intValue);
- }
- field.setInt(obj, intValue);
+ field.setInt(obj, SessionVariable.normalizeIntValue(name,
value));
break;
case "long":
field.setLong(obj, Long.parseLong(value));
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/catalog/UserPropertyTest.java
b/fe/fe-core/src/test/java/org/apache/doris/catalog/UserPropertyTest.java
index 46c5bea29ec..773a340c425 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/catalog/UserPropertyTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/catalog/UserPropertyTest.java
@@ -20,6 +20,7 @@ package org.apache.doris.catalog;
import org.apache.doris.analysis.UserIdentity;
import org.apache.doris.authentication.BasicPrincipal;
import org.apache.doris.blockrule.SqlBlockRuleMgr;
+import org.apache.doris.common.DdlException;
import org.apache.doris.common.Pair;
import org.apache.doris.common.UserException;
import org.apache.doris.datasource.CatalogIf;
@@ -74,7 +75,7 @@ public class UserPropertyTest {
List<Pair<String, String>> properties = Lists.newArrayList();
properties.add(Pair.of("MAX_USER_CONNECTIONS", "100"));
properties.add(Pair.of("max_qUERY_instances", "3000"));
- properties.add(Pair.of("parallel_fragment_exec_instance_num", "2000"));
+ properties.add(Pair.of("parallel_fragment_exec_instance_num", "256"));
properties.add(Pair.of("sql_block_rules", "rule1,rule2"));
properties.add(Pair.of("cpu_resource_limit", "2"));
properties.add(Pair.of("query_timeout", "500"));
@@ -85,7 +86,7 @@ public class UserPropertyTest {
userProperty.update(properties);
Assertions.assertEquals(100, userProperty.getMaxConn());
Assertions.assertEquals(3000, userProperty.getMaxQueryInstances());
- Assertions.assertEquals(2000,
userProperty.getParallelFragmentExecInstanceNum());
+ Assertions.assertEquals(256,
userProperty.getParallelFragmentExecInstanceNum());
Assertions.assertArrayEquals(new String[]{"rule1", "rule2"},
userProperty.getSqlBlockRules());
Assertions.assertEquals(2, userProperty.getCpuResourceLimit());
Assertions.assertEquals(500, userProperty.getQueryTimeout());
@@ -123,6 +124,27 @@ public class UserPropertyTest {
Assertions.assertEquals(3, userProperty.getSqlBlockRules().length);
}
+ @Test
+ public void testParallelFragmentExecInstanceNumLimit() throws
UserException {
+ UserProperty property = new UserProperty();
+ for (int value : new int[] {Integer.MIN_VALUE, -1, 0, 1, 256}) {
+ property.update(Lists.newArrayList(
+ Pair.of("parallel_fragment_exec_instance_num",
Integer.toString(value))));
+ Assertions.assertEquals(value,
property.getParallelFragmentExecInstanceNum());
+ }
+
+ for (int value : new int[] {257, 2000, Integer.MAX_VALUE}) {
+ DdlException exception =
Assertions.assertThrows(DdlException.class,
+ () -> property.update(Lists.newArrayList(
+ Pair.of("max_user_connections", "200"),
+ Pair.of("PARALLEL_FRAGMENT_EXEC_INSTANCE_NUM",
Integer.toString(value)))));
+ Assertions.assertTrue(exception.getMessage().contains(
+ "parallel_fragment_exec_instance_num must be less than or
equal to 256, got " + value));
+ Assertions.assertEquals(256,
property.getParallelFragmentExecInstanceNum());
+ Assertions.assertEquals(100, property.getMaxConn());
+ }
+ }
+
@Test
public void testValidation() throws UserException {
List<Pair<String, String>> properties = Lists.newArrayList();
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/load/NereidsStreamLoadTaskTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/load/NereidsStreamLoadTaskTest.java
new file mode 100644
index 00000000000..7e690bec0fa
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/load/NereidsStreamLoadTaskTest.java
@@ -0,0 +1,68 @@
+// 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.nereids.load;
+
+import org.apache.doris.common.UserException;
+import org.apache.doris.thrift.TFileCompressType;
+import org.apache.doris.thrift.TFileFormatType;
+import org.apache.doris.thrift.TFileType;
+import org.apache.doris.thrift.TStreamLoadPutRequest;
+import org.apache.doris.thrift.TUniqueId;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+public class NereidsStreamLoadTaskTest {
+ @Test
+ public void testDefaultSendBatchParallelism() throws UserException {
+ Assertions.assertEquals(1,
+
NereidsStreamLoadTask.fromTStreamLoadPutRequest(newRequest()).getSendBatchParallelism());
+ }
+
+ @Test
+ public void testSendBatchParallelismBoundary() throws UserException {
+ for (int parallelism : new int[] {Integer.MIN_VALUE, -1, 0, 1, 256}) {
+ TStreamLoadPutRequest request = newRequest();
+ request.setSendBatchParallelism(parallelism);
+ Assertions.assertEquals(parallelism,
+
NereidsStreamLoadTask.fromTStreamLoadPutRequest(request).getSendBatchParallelism());
+ }
+ }
+
+ @Test
+ public void testRejectExcessiveSendBatchParallelism() {
+ for (int parallelism : new int[] {257, Integer.MAX_VALUE}) {
+ TStreamLoadPutRequest request = newRequest();
+ request.setSendBatchParallelism(parallelism);
+ UserException exception =
Assertions.assertThrows(UserException.class,
+ () ->
NereidsStreamLoadTask.fromTStreamLoadPutRequest(request));
+ Assertions.assertTrue(exception.getMessage().contains(
+ "send_batch_parallelism value should less than or equal
256, you set value is: " + parallelism));
+ }
+ }
+
+ private TStreamLoadPutRequest newRequest() {
+ TStreamLoadPutRequest request = new TStreamLoadPutRequest();
+ request.setLoadId(new TUniqueId(1, 2));
+ request.setTxnId(3);
+ request.setFileType(TFileType.FILE_STREAM);
+ request.setFormatType(TFileFormatType.FORMAT_CSV_PLAIN);
+ request.setCompressType(TFileCompressType.UNKNOWN);
+ return request;
+ }
+}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/qe/SessionVariableParallelismTest.java
b/fe/fe-core/src/test/java/org/apache/doris/qe/SessionVariableParallelismTest.java
new file mode 100644
index 00000000000..ff0b6ef49c9
--- /dev/null
+++
b/fe/fe-core/src/test/java/org/apache/doris/qe/SessionVariableParallelismTest.java
@@ -0,0 +1,217 @@
+// 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.qe;
+
+import org.apache.doris.analysis.UserIdentity;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.Pair;
+import org.apache.doris.common.io.Text;
+import org.apache.doris.datasource.InternalCatalog;
+import org.apache.doris.mysql.privilege.Auth;
+import org.apache.doris.mysql.privilege.UserPropertyInfo;
+import org.apache.doris.mysql.privilege.UserPropertyMgr;
+import org.apache.doris.planner.DataPartition;
+import org.apache.doris.planner.PlanFragment;
+import org.apache.doris.planner.PlanFragmentId;
+import org.apache.doris.system.Backend;
+import org.apache.doris.system.SystemInfoService;
+
+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 org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.DataInputStream;
+import java.io.DataOutputStream;
+import java.util.Collections;
+
+public class SessionVariableParallelismTest {
+ private static final String USER = "parallelism_test";
+ private UserPropertyMgr propertyMgr;
+ private SessionVariable sessionVariable;
+
+ @BeforeEach
+ public void setUp() throws Exception {
+ propertyMgr = new UserPropertyMgr();
+ propertyMgr.addUserResource(USER);
+ Auth auth = Mockito.mock(Auth.class);
+ Mockito.when(auth.getParallelFragmentExecInstanceNum(USER))
+ .thenAnswer(invocation ->
propertyMgr.getParallelFragmentExecInstanceNum(USER));
+ Env env = Mockito.mock(Env.class);
+ Mockito.when(env.getAuth()).thenReturn(auth);
+ Mockito.when(env.getInternalCatalog()).thenReturn(new
InternalCatalog());
+ ConnectContext context = new ConnectContext();
+ context.setEnv(env);
+
context.setCurrentUserIdentity(UserIdentity.createAnalyzedUserIdentWithIp(USER,
"%"));
+ context.setThreadLocalInfo();
+ sessionVariable = context.getSessionVariable();
+ sessionVariable.setPipelineTaskNum("8");
+ }
+
+ @AfterEach
+ public void tearDown() {
+ ConnectContext.remove();
+ }
+
+ @Test
+ public void testAutomaticParallelismLimit() throws Exception {
+ sessionVariable.setPipelineTaskNum("0");
+ Backend backend = new Backend(1, "127.0.0.1", 9050);
+ SystemInfoService systemInfo = new SystemInfoService();
+ systemInfo.addBackend(backend);
+ try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+ mockedEnv.when(Env::getCurrentSystemInfo).thenReturn(systemInfo);
+ // Executor report, max_instance_num, effective parallelism.
+ for (int[] values : new int[][] {{1024, 1024, 256}, {1024, 64,
64}, {1024, 256, 256},
+ {510, 1024, 255}, {511, 1024, 256}, {512, 1024, 256},
{513, 1024, 256},
+ {8, 1024, 4}, {9, 1024, 5}}) {
+ backend.setPipelineExecutorSize(values[0]);
+ sessionVariable.maxInstanceNum = values[1];
+ assertEffectiveParallelism(values[2]);
+ }
+
+ // A positive user property still takes precedence over automatic
sizing.
+ propertyMgr.updateUserProperty(USER, Collections.singletonList(
+ Pair.of("parallel_fragment_exec_instance_num", "7")),
false);
+ assertEffectiveParallelism(7);
+ }
+ }
+
+ @Test
+ public void testHistoricalSessionParallelismLimit() throws Exception {
+ for (int value : new int[] {257, 2000, Integer.MAX_VALUE}) {
+ sessionVariable.readFromJson("{\"parallel_pipeline_task_num\":" +
value + "}");
+ assertEffectiveParallelism(256);
+ }
+ sessionVariable.setPipelineTaskNum("8");
+ assertEffectiveParallelism(8);
+ }
+
+ @Test
+ public void testHistoricalConcurrencyHintRestoration() throws Exception {
+ for (String variable : new String[]
{SessionVariable.PARALLEL_PIPELINE_TASK_NUM,
+ SessionVariable.COLOCATE_MAX_PARALLEL_NUM,
SessionVariable.MAX_SCANNERS_CONCURRENCY,
+ SessionVariable.MAX_FILE_SCANNERS_CONCURRENCY,
SessionVariable.MIN_SCANNERS_CONCURRENCY,
+ SessionVariable.MIN_FILE_SCANNERS_CONCURRENCY,
SessionVariable.PARALLEL_SCAN_MAX_SCANNERS_COUNT,
+ SessionVariable.SEND_BATCH_PARALLELISM,
SessionVariable.LOAD_STREAM_PER_NODE}) {
+ for (int value : new int[] {257, 2000, Integer.MAX_VALUE}) {
+ assertHistoricalHintRestoration(variable, value, 256);
+ }
+ assertHistoricalHintRestoration(variable, 256, 256);
+ assertHistoricalHintRestoration(variable, 8, 8);
+ }
+ }
+
+ @Test
+ public void testHistoricalConcurrencyDefaultValues() throws Exception {
+
assertHistoricalHintRestoration(SessionVariable.PARALLEL_PIPELINE_TASK_NUM, 0,
0);
+ for (String variable : new String[]
{SessionVariable.COLOCATE_MAX_PARALLEL_NUM,
+ SessionVariable.LOAD_STREAM_PER_NODE}) {
+ assertHistoricalHintRestoration(variable, 0, 1);
+ assertHistoricalHintRestoration(variable, -1, 1);
+ }
+ for (String variable : new String[]
{SessionVariable.MAX_SCANNERS_CONCURRENCY,
+ SessionVariable.MAX_FILE_SCANNERS_CONCURRENCY,
SessionVariable.MIN_SCANNERS_CONCURRENCY,
+ SessionVariable.MIN_FILE_SCANNERS_CONCURRENCY,
SessionVariable.PARALLEL_SCAN_MAX_SCANNERS_COUNT,
+ SessionVariable.SEND_BATCH_PARALLELISM}) {
+ for (int value : new int[] {0, -1, Integer.MIN_VALUE}) {
+ assertHistoricalHintRestoration(variable, value, value);
+ }
+ }
+ }
+
+ private void assertHistoricalHintRestoration(String variable, int value,
int expected) throws Exception {
+ SessionVariable fromJson = new SessionVariable();
+ fromJson.readFromJson("{\"" + variable + "\":" + value + "}");
+ assertHintRestoration(fromJson, variable, expected);
+
+ SessionVariable fromMap = new SessionVariable();
+ fromMap.readFromMap(Collections.singletonMap(variable,
Integer.toString(value)));
+ assertHintRestoration(fromMap, variable, expected);
+
+ SessionVariable forwarded = new SessionVariable();
+ if (forwarded.getForwardVariables().containsKey(variable)) {
+
forwarded.setForwardedSessionVariables(Collections.singletonMap(variable,
Integer.toString(value)));
+ assertHintRestoration(forwarded, variable, expected);
+ }
+ }
+
+ private void assertHintRestoration(SessionVariable restored, String
variable, int expected) throws Exception {
+ Assertions.assertTrue(restored.setVarOnce(variable, "8"));
+ VariableMgr.revertSessionValue(restored);
+ Assertions.assertEquals(expected,
VariableMgr.getVarContext(variable).getField().getInt(restored), variable);
+ }
+
+ @Test
+ public void testUserPropertyPrecedenceAndReset() throws Exception {
+ assertEffectiveParallelism(8);
+ for (int value : new int[] {1, 256, 0, -1, Integer.MIN_VALUE}) {
+ propertyMgr.updateUserProperty(USER, Collections.singletonList(
+ Pair.of("parallel_fragment_exec_instance_num",
Integer.toString(value))), false);
+ assertEffectiveParallelism(value > 0 ? value : 8);
+ }
+ }
+
+ @Test
+ public void testHistoricalJournalParallelism() throws Exception {
+ for (int value : new int[] {257, 2000, Integer.MAX_VALUE}) {
+ UserPropertyInfo info = new UserPropertyInfo(USER,
Collections.singletonList(
+ Pair.of("parallel_fragment_exec_instance_num",
Integer.toString(value))));
+ ByteArrayOutputStream bytes = new ByteArrayOutputStream();
+ info.write(new DataOutputStream(bytes));
+ UserPropertyInfo restored = UserPropertyInfo.read(
+ new DataInputStream(new
ByteArrayInputStream(bytes.toByteArray())));
+ propertyMgr.updateUserProperty(restored.getUser(),
restored.getProperties(), true);
+ Assertions.assertEquals(value,
propertyMgr.getParallelFragmentExecInstanceNum(USER));
+ assertEffectiveParallelism(256);
+
+ bytes.reset();
+ propertyMgr.write(new DataOutputStream(bytes));
+ propertyMgr = UserPropertyMgr.read(new DataInputStream(new
ByteArrayInputStream(bytes.toByteArray())));
+ assertEffectiveParallelism(256);
+ }
+ }
+
+ @Test
+ public void testHistoricalImageParallelism() throws Exception {
+ // Cover both the current image field names and their legacy aliases.
+ for (String[] fields : new String[][] {{"cp", "pfei"},
+ {"commonProperties", "parallelFragmentExecInstanceNum"}}) {
+ for (int value : new int[] {257, 2000, Integer.MAX_VALUE}) {
+ String json =
String.format("{\"propertyMap\":{\"%s\":{\"qu\":\"%s\",\"%s\":{\"%s\":%d}}}}",
+ USER, USER, fields[0], fields[1], value);
+ ByteArrayOutputStream bytes = new ByteArrayOutputStream();
+ Text.writeString(new DataOutputStream(bytes), json);
+ propertyMgr = UserPropertyMgr.read(new DataInputStream(new
ByteArrayInputStream(bytes.toByteArray())));
+ Assertions.assertEquals(value,
propertyMgr.getParallelFragmentExecInstanceNum(USER));
+ assertEffectiveParallelism(256);
+ }
+ }
+ }
+
+ private void assertEffectiveParallelism(int expected) {
+ Assertions.assertEquals(expected,
sessionVariable.getParallelExecInstanceNum(""));
+ Assertions.assertEquals(expected,
sessionVariable.toThrift().getParallelInstance());
+ PlanFragment fragment = new PlanFragment(new PlanFragmentId(0), null,
DataPartition.RANDOM);
+ Assertions.assertEquals(expected, fragment.getParallelExecNum());
+ }
+}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/qe/SessionVariablesTest.java
b/fe/fe-core/src/test/java/org/apache/doris/qe/SessionVariablesTest.java
index 315ac6554dc..306c6570f02 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/qe/SessionVariablesTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/qe/SessionVariablesTest.java
@@ -269,6 +269,33 @@ public class SessionVariablesTest extends
TestWithFeService {
Assertions.assertEquals(false,
connectContext.getSessionVariable().enableNereidsDmlWithPipeline);
}
+ @Test
+ public void testHistoricalConcurrencyHintCleanup() throws Exception {
+ SessionVariable original = connectContext.getSessionVariable();
+ try {
+ for (String variable : new String[]
{SessionVariable.PARALLEL_PIPELINE_TASK_NUM,
+ SessionVariable.COLOCATE_MAX_PARALLEL_NUM,
SessionVariable.MAX_SCANNERS_CONCURRENCY,
+ SessionVariable.MAX_FILE_SCANNERS_CONCURRENCY,
SessionVariable.MIN_SCANNERS_CONCURRENCY,
+ SessionVariable.MIN_FILE_SCANNERS_CONCURRENCY,
SessionVariable.PARALLEL_SCAN_MAX_SCANNERS_COUNT,
+ SessionVariable.SEND_BATCH_PARALLELISM,
SessionVariable.LOAD_STREAM_PER_NODE}) {
+ SessionVariable restored = new SessionVariable();
+ restored.readFromJson("{\"" + variable + "\":2000}");
+ connectContext.setSessionVariable(restored);
+
+ executeNereidsSql("SELECT /*+ SET_VAR(" + variable + "=8) */
1");
+ Assertions.assertEquals(256,
VariableMgr.getVarContext(variable).getField().getInt(restored), variable);
+ Assertions.assertFalse(restored.getIsSingleSetVar());
+
Assertions.assertTrue(restored.getSessionOriginValue().isEmpty());
+
+ executeNereidsSql("SELECT 1");
+ Assertions.assertEquals(256,
VariableMgr.getVarContext(variable).getField().getInt(restored), variable);
+
Assertions.assertTrue(restored.getSessionOriginValue().isEmpty());
+ }
+ } finally {
+ connectContext.setSessionVariable(original);
+ }
+ }
+
@Test
public void testAiSessionVariableChecker() throws Exception {
SessionVariable sv = new SessionVariable();
diff --git
a/regression-test/data/load_p0/stream_load/test_stream_load_send_batch_parallelism.out
b/regression-test/data/load_p0/stream_load/test_stream_load_send_batch_parallelism.out
new file mode 100644
index 00000000000..354b718d19f
--- /dev/null
+++
b/regression-test/data/load_p0/stream_load/test_stream_load_send_batch_parallelism.out
@@ -0,0 +1,52 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !false_default --
+Success 1 false
+
+-- !false_-2147483648 --
+Success 1 false
+
+-- !false_-1 --
+Success 1 false
+
+-- !false_0 --
+Success 1 false
+
+-- !false_1 --
+Success 1 false
+
+-- !false_256 --
+Success 1 false
+
+-- !false_257 --
+Fail 0 true
+
+-- !false_2147483647 --
+Fail 0 true
+
+-- !true_default --
+Success 1 false
+
+-- !true_-2147483648 --
+Success 1 false
+
+-- !true_-1 --
+Success 1 false
+
+-- !true_0 --
+Success 1 false
+
+-- !true_1 --
+Success 1 false
+
+-- !true_256 --
+Success 1 false
+
+-- !true_257 --
+Fail 0 true
+
+-- !true_2147483647 --
+Fail 0 true
+
+-- !loaded_rows --
+12
+
diff --git
a/regression-test/data/query_p0/session_variable/test_session_concurrency_limits.out
b/regression-test/data/query_p0/session_variable/test_session_concurrency_limits.out
new file mode 100644
index 00000000000..7b276c14f44
--- /dev/null
+++
b/regression-test/data/query_p0/session_variable/test_session_concurrency_limits.out
@@ -0,0 +1,112 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !parallel_pipeline_task_num_upper_bound --
+true
+
+-- !parallel_pipeline_task_num_valid_hint --
+true
+
+-- !parallel_pipeline_task_num_hint_restored --
+true
+
+-- !parallel_pipeline_task_num_unchanged --
+true true
+
+-- !colocate_max_parallel_num_upper_bound --
+true
+
+-- !colocate_max_parallel_num_valid_hint --
+true
+
+-- !colocate_max_parallel_num_hint_restored --
+true
+
+-- !colocate_max_parallel_num_unchanged --
+true true
+
+-- !max_scanners_concurrency_upper_bound --
+true
+
+-- !max_scanners_concurrency_valid_hint --
+true
+
+-- !max_scanners_concurrency_hint_restored --
+true
+
+-- !max_scanners_concurrency_unchanged --
+true true
+
+-- !max_file_scanners_concurrency_upper_bound --
+true
+
+-- !max_file_scanners_concurrency_valid_hint --
+true
+
+-- !max_file_scanners_concurrency_hint_restored --
+true
+
+-- !max_file_scanners_concurrency_unchanged --
+true true
+
+-- !min_scanners_concurrency_upper_bound --
+true
+
+-- !min_scanners_concurrency_valid_hint --
+true
+
+-- !min_scanners_concurrency_hint_restored --
+true
+
+-- !min_scanners_concurrency_unchanged --
+true true
+
+-- !min_file_scanners_concurrency_upper_bound --
+true
+
+-- !min_file_scanners_concurrency_valid_hint --
+true
+
+-- !min_file_scanners_concurrency_hint_restored --
+true
+
+-- !min_file_scanners_concurrency_unchanged --
+true true
+
+-- !parallel_scan_max_scanners_count_upper_bound --
+true
+
+-- !parallel_scan_max_scanners_count_valid_hint --
+true
+
+-- !parallel_scan_max_scanners_count_hint_restored --
+true
+
+-- !parallel_scan_max_scanners_count_unchanged --
+true true
+
+-- !send_batch_parallelism_upper_bound --
+true
+
+-- !send_batch_parallelism_valid_hint --
+true
+
+-- !send_batch_parallelism_hint_restored --
+true
+
+-- !send_batch_parallelism_unchanged --
+true true
+
+-- !load_stream_per_node_upper_bound --
+true
+
+-- !load_stream_per_node_valid_hint --
+true
+
+-- !load_stream_per_node_hint_restored --
+true
+
+-- !load_stream_per_node_unchanged --
+true true
+
+-- !pipeline_auto --
+0
+
diff --git
a/regression-test/data/query_p0/session_variable/test_user_parallelism_limit.out
b/regression-test/data/query_p0/session_variable/test_user_parallelism_limit.out
new file mode 100644
index 00000000000..0024f68687e
--- /dev/null
+++
b/regression-test/data/query_p0/session_variable/test_user_parallelism_limit.out
@@ -0,0 +1,22 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !upper_bound --
+parallel_fragment_exec_instance_num 256
+
+-- !unchanged_parallelism --
+parallel_fragment_exec_instance_num 256
+
+-- !unchanged_connections --
+max_user_connections 100
+
+-- !value_1 --
+parallel_fragment_exec_instance_num 1
+
+-- !value_0 --
+parallel_fragment_exec_instance_num 0
+
+-- !value_-1 --
+parallel_fragment_exec_instance_num -1
+
+-- !value_-2147483648 --
+parallel_fragment_exec_instance_num -2147483648
+
diff --git
a/regression-test/suites/load_p0/stream_load/test_stream_load_send_batch_parallelism.groovy
b/regression-test/suites/load_p0/stream_load/test_stream_load_send_batch_parallelism.groovy
new file mode 100644
index 00000000000..a8057537523
--- /dev/null
+++
b/regression-test/suites/load_p0/stream_load/test_stream_load_send_batch_parallelism.groovy
@@ -0,0 +1,52 @@
+// 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_stream_load_send_batch_parallelism", "p0") {
+ sql "DROP TABLE IF EXISTS test_stream_load_send_batch_parallelism"
+ sql """
+ CREATE TABLE test_stream_load_send_batch_parallelism (id INT)
+ DUPLICATE KEY(id)
+ DISTRIBUTED BY HASH(id) BUCKETS 1
+ PROPERTIES ("replication_num" = "1")
+ """
+
+ // Both writer versions must validate the HTTP override at the FE planning
boundary.
+ for (def memtableOnSinkNode : ["false", "true"]) {
+ for (def parallelism : [null, "-2147483648", "-1", "0", "1", "256",
"257", "2147483647"]) {
+ streamLoad {
+ table "test_stream_load_send_batch_parallelism"
+ set "memtable_on_sink_node", memtableOnSinkNode
+ if (parallelism != null) {
+ set "send_batch_parallelism", parallelism
+ }
+ inputText "1\n"
+ check { result, exception, startTime, endTime ->
+ if (exception != null) {
+ throw exception
+ }
+ def json = parseJson(result)
+ def limitError = "send_batch_parallelism value should less
than or equal 256, " +
+ "you set value is: ${parallelism}"
+ "qt_${memtableOnSinkNode}_${parallelism ?: 'default'}"(
+ "SELECT '${json.Status}',
${json.NumberLoadedRows}, ${json.Message.contains(limitError)}")
+ }
+ }
+ }
+ }
+ sql "sync"
+ qt_loaded_rows "SELECT count(*) FROM
test_stream_load_send_batch_parallelism"
+}
diff --git
a/regression-test/suites/query_p0/session_variable/test_session_concurrency_limits.groovy
b/regression-test/suites/query_p0/session_variable/test_session_concurrency_limits.groovy
new file mode 100644
index 00000000000..9edf3fe4504
--- /dev/null
+++
b/regression-test/suites/query_p0/session_variable/test_session_concurrency_limits.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_session_concurrency_limits") {
+ def variables = [
+ "parallel_pipeline_task_num",
+ "colocate_max_parallel_num",
+ "max_scanners_concurrency",
+ "max_file_scanners_concurrency",
+ "min_scanners_concurrency",
+ "min_file_scanners_concurrency",
+ "parallel_scan_max_scanners_count",
+ "send_batch_parallelism",
+ "load_stream_per_node"
+ ]
+ def upperBound = 256
+
+ variables.each { variable ->
+ def original = (sql "SELECT @@${variable}")[0][0]
+ def originalGlobal = (sql "SELECT @@global.${variable}")[0][0]
+ try {
+ sql "SET ${variable} = ${upperBound}"
+ "order_qt_${variable}_upper_bound"("SELECT @@${variable} =
${upperBound}")
+
+ "order_qt_${variable}_valid_hint"("SELECT /*+
SET_VAR(${variable}=8) */ @@${variable} = 8")
+ "order_qt_${variable}_hint_restored"("SELECT @@${variable} =
${upperBound}")
+
+ [upperBound + 1, 2147483647L].each { invalid ->
+ test {
+ sql "SET ${variable} = ${invalid}"
+ exception "${variable} value should less than or equal
${upperBound}"
+ }
+ }
+ test {
+ sql "SET GLOBAL ${variable} = ${upperBound + 1}"
+ exception "${variable} value should less than or equal
${upperBound}"
+ }
+ test {
+ sql "SELECT /*+ SET_VAR(${variable}=${upperBound + 1}) */ 1"
+ exception "Can not set session variable '${variable}'"
+ }
+ "order_qt_${variable}_unchanged"(
+ "SELECT @@${variable} = ${upperBound},
@@global.${variable} = ${originalGlobal}")
+ } finally {
+ sql "SET ${variable} = ${original}"
+ }
+ }
+
+ def originalPipeline = (sql "SELECT @@parallel_pipeline_task_num")[0][0]
+ try {
+ sql "SET parallel_pipeline_task_num = 0"
+ order_qt_pipeline_auto "SELECT @@parallel_pipeline_task_num"
+ test {
+ sql "SET parallel_pipeline_task_num = -1"
+ exception "parallel_pipeline_task_num value should greater than or
equal 0"
+ }
+ } finally {
+ sql "SET parallel_pipeline_task_num = ${originalPipeline}"
+ }
+
+ ["colocate_max_parallel_num", "load_stream_per_node"].each { variable ->
+ [0, -1].each { invalid ->
+ test {
+ sql "SET ${variable} = ${invalid}"
+ exception "${variable} value should greater than or equal 1"
+ }
+ }
+ }
+}
diff --git
a/regression-test/suites/query_p0/session_variable/test_user_parallelism_limit.groovy
b/regression-test/suites/query_p0/session_variable/test_user_parallelism_limit.groovy
new file mode 100644
index 00000000000..02fa5d23320
--- /dev/null
+++
b/regression-test/suites/query_p0/session_variable/test_user_parallelism_limit.groovy
@@ -0,0 +1,39 @@
+// 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_user_parallelism_limit") {
+ sql "DROP USER IF EXISTS 'test_user_parallelism_limit'"
+ sql "CREATE USER 'test_user_parallelism_limit'"
+
+ sql "SET PROPERTY FOR 'test_user_parallelism_limit'
'parallel_fragment_exec_instance_num' = '256'"
+ qt_upper_bound "SHOW PROPERTY FOR 'test_user_parallelism_limit' LIKE
'parallel_fragment_exec_instance_num'"
+
+ for (def value : [257, 2000, 2147483647]) {
+ test {
+ sql """SET PROPERTY FOR 'test_user_parallelism_limit'
+ 'max_user_connections' = '200',
'PARALLEL_FRAGMENT_EXEC_INSTANCE_NUM' = '${value}'"""
+ exception "parallel_fragment_exec_instance_num must be less than
or equal to 256, got ${value}"
+ }
+ }
+ qt_unchanged_parallelism "SHOW PROPERTY FOR 'test_user_parallelism_limit'
LIKE 'parallel_fragment_exec_instance_num'"
+ qt_unchanged_connections "SHOW PROPERTY FOR 'test_user_parallelism_limit'
LIKE 'max_user_connections'"
+
+ for (def value : [1, 0, -1, -2147483648]) {
+ sql "SET PROPERTY FOR 'test_user_parallelism_limit'
'parallel_fragment_exec_instance_num' = '${value}'"
+ "qt_value_${value}"("SHOW PROPERTY FOR 'test_user_parallelism_limit'
LIKE 'parallel_fragment_exec_instance_num'")
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]