This is an automated email from the ASF dual-hosted git repository.
luwei16 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 607c8a702f0 [fix](binlog) Reject invalid Table Stream inputs (#68792)
607c8a702f0 is described below
commit 607c8a702f063249b0df7c5fb6fbdc8283afe9d0
Author: Luwei <[email protected]>
AuthorDate: Fri Oct 9 14:27:50 2026 +0800
[fix](binlog) Reject invalid Table Stream inputs (#68792)
### What problem does this PR solve?
Issue Number: N/A
Related PR: N/A
Problem Summary: CREATE STREAM silently converts invalid
show_initial_rows strings such as "garbage" to false, which can omit the
expected initial rows. The snapshot and reset scan modes also accept and
ignore unsupported arguments. Use strict boolean parsing for the Stream
property and reject non-empty map and positional arguments during Stream
binding. Invalid inputs now fail before Stream publication or query
execution, while valid booleans and parameterless calls retain their
existing behavior.
### Release note
CREATE STREAM rejects invalid show_initial_rows values. Table Stream
snapshot and reset reads reject unsupported arguments with explicit
errors.
### Check List (For Author)
- Test: Unit Test
- 39 focused FE tests passed; the three new checks failed on the
original implementation.
- FE Checkstyle and git diff --check passed.
- Added SQL regression coverage; not run because validation used the FE
unit-test harness without a deployed cluster.
- Behavior changed: Yes, reject invalid boolean values and unsupported
scan arguments.
- Does this need documentation: No, enforce the existing property and
zero-argument contracts.
---
.../org/apache/doris/analysis/TableScanParams.java | 5 ++
.../doris/catalog/stream/BaseTableStream.java | 14 +++-
.../apache/doris/analysis/TableScanParamsTest.java | 16 +++++
.../doris/catalog/CreateTableStreamTest.java | 32 +++++++++
.../trees/plans/ExplainTableStreamPlanTest.java | 27 +++++++
.../test_table_stream_input_validation.groovy | 84 ++++++++++++++++++++++
6 files changed, 175 insertions(+), 3 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/analysis/TableScanParams.java
b/fe/fe-core/src/main/java/org/apache/doris/analysis/TableScanParams.java
index b2189aa13ef..74e38484d91 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/analysis/TableScanParams.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/analysis/TableScanParams.java
@@ -17,6 +17,8 @@
package org.apache.doris.analysis;
+import org.apache.doris.nereids.exceptions.AnalysisException;
+
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.ImmutableSet;
@@ -75,6 +77,9 @@ public class TableScanParams {
if (!VALID_OLAP_TABLE_STREAM_PARAM_TYPES.contains(paramType)) {
throw new IllegalArgumentException("Invalid param type for olap
table stream : " + paramType);
}
+ if (!mapParams.isEmpty() || !listParams.isEmpty()) {
+ throw new AnalysisException(paramType + " does not accept
parameters");
+ }
}
public TableScanParams(String paramType, Map<String, String> mapParams,
List<String> listParams) {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/BaseTableStream.java
b/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/BaseTableStream.java
index d241d02a318..1cfa6986a61 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/BaseTableStream.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/BaseTableStream.java
@@ -23,6 +23,7 @@ import org.apache.doris.catalog.TableIf;
import org.apache.doris.common.UserException;
import org.apache.doris.common.io.Text;
import org.apache.doris.common.util.PropertyAnalyzer;
+import org.apache.doris.common.util.Util;
import org.apache.doris.nereids.exceptions.AnalysisException;
import org.apache.doris.persist.gson.GsonUtils;
import org.apache.doris.thrift.TBinlogScanType;
@@ -152,9 +153,16 @@ public abstract class BaseTableStream extends Table {
}
public void setProperties(Map<String, String> properties) throws
org.apache.doris.common.AnalysisException {
- showInitialRows = PropertyAnalyzer.analyzeBooleanProp(properties,
- PropertyAnalyzer.PROPERTIES_STREAM_SHOW_INITIAL_ROWS,
- false);
+ showInitialRows = false;
+ String key = PropertyAnalyzer.PROPERTIES_STREAM_SHOW_INITIAL_ROWS;
+ if (properties != null && properties.containsKey(key)) {
+ try {
+ showInitialRows =
Util.parseBooleanProperty(properties.get(key), key);
+ } catch (org.apache.doris.common.AnalysisException e) {
+ throw new org.apache.doris.common.AnalysisException(key + "
must be `true` or `false`", e);
+ }
+ properties.remove(key);
+ }
streamScanType = PropertyAnalyzer.analyzeStreamType(properties);
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/analysis/TableScanParamsTest.java
b/fe/fe-core/src/test/java/org/apache/doris/analysis/TableScanParamsTest.java
index 5adfa5f0ea2..b45d7292728 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/analysis/TableScanParamsTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/analysis/TableScanParamsTest.java
@@ -17,6 +17,8 @@
package org.apache.doris.analysis;
+import org.apache.doris.nereids.exceptions.AnalysisException;
+
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import org.junit.jupiter.api.Assertions;
@@ -85,6 +87,20 @@ public class TableScanParamsTest {
new TableScanParams(TableScanParams.RESET, EMPTY_MAP,
EMPTY_LIST).validateOlapTableStream();
}
+ @Test
+ public void testValidateOlapTableStreamRejectsArguments() {
+ for (String mode : ImmutableList.of(TableScanParams.SNAPSHOT,
TableScanParams.RESET)) {
+ AnalysisException mapError =
Assertions.assertThrows(AnalysisException.class,
+ () -> new TableScanParams(mode, ImmutableMap.of("unknown",
"x"), EMPTY_LIST)
+ .validateOlapTableStream());
+ Assertions.assertEquals(mode + " does not accept parameters",
mapError.getMessage());
+ AnalysisException listError =
Assertions.assertThrows(AnalysisException.class,
+ () -> new TableScanParams(mode, EMPTY_MAP,
ImmutableList.of("x"))
+ .validateOlapTableStream());
+ Assertions.assertEquals(mode + " does not accept parameters",
listError.getMessage());
+ }
+ }
+
@Test
public void testValidateOlapTableStreamRejectsOthers() {
IllegalArgumentException e =
Assertions.assertThrows(IllegalArgumentException.class,
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableStreamTest.java
b/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableStreamTest.java
index f311ca426a7..73f09ab67c4 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableStreamTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/catalog/CreateTableStreamTest.java
@@ -120,6 +120,38 @@ public class CreateTableStreamTest extends
TestWithFeService {
dropDatabase("test_stream_type_validation");
}
+ @Test
+ public void testCreateStreamShowInitialRowsValidation() throws Exception {
+ createDatabase("test_stream_boolean_validation");
+ createTable("create table test_stream_boolean_validation.base_table
(k1 int, k2 int) "
+ + "unique key(k1) distributed by hash(k1) buckets 1 "
+ + "properties('replication_num' = '1', 'binlog.enable' =
'true', "
+ + "'binlog.format' = 'ROW', 'binlog.need_historical_value' =
'true')");
+ Database db =
Env.getCurrentInternalCatalog().getDbOrDdlException("test_stream_boolean_validation");
+
+ for (String value : new String[] {"garbage", "yes", "1", "", " true
"}) {
+ ExceptionChecker.expectThrowsWithMsg(DdlException.class,
+ "show_initial_rows must be `true` or `false`",
+ () -> createTable("create stream
test_stream_boolean_validation.invalid_stream "
+ + "on table
test_stream_boolean_validation.base_table "
+ + "properties('show_initial_rows' = '" + value +
"')"));
+ Assertions.assertFalse(db.getTable("invalid_stream").isPresent());
+ }
+
+ String[] validValues = {"true", "false", "TrUe", "FaLsE"};
+ for (int i = 0; i < validValues.length; i++) {
+ String value = validValues[i];
+ createTable("create stream test_stream_boolean_validation.stream_"
+ i
+ + " on table test_stream_boolean_validation.base_table "
+ + "properties('show_initial_rows' = '" + value + "')");
+ BaseTableStream stream = (BaseTableStream)
db.getTableOrDdlException("stream_" + i);
+ Assertions.assertEquals("true".equalsIgnoreCase(value),
stream.isShowInitialRows());
+ }
+ createTable("create stream
test_stream_boolean_validation.default_stream "
+ + "on table test_stream_boolean_validation.base_table");
+ Assertions.assertFalse(((BaseTableStream)
db.getTableOrDdlException("default_stream")).isShowInitialRows());
+ }
+
@Test
public void testCreateStreamAbnormalOLAP() throws Exception {
createDatabase("test_stream");
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java
index 06698b6b966..4c7d4b9bfdc 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java
@@ -38,6 +38,7 @@ import org.apache.doris.common.Pair;
import org.apache.doris.common.util.TimeUtils;
import org.apache.doris.nereids.NereidsPlanner;
import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.nereids.exceptions.AnalysisException;
import org.apache.doris.nereids.glue.translator.PhysicalPlanTranslator;
import org.apache.doris.nereids.glue.translator.PlanTranslatorContext;
import org.apache.doris.nereids.parser.NereidsParser;
@@ -460,6 +461,32 @@ public class ExplainTableStreamPlanTest extends
TestWithFeService {
Assertions.assertTrue(assertedAtLeastOne);
}
+ @Test
+ public void testStreamReadModesRejectArgumentsWithoutChangingOffsets()
throws Exception {
+ Database db =
Env.getCurrentInternalCatalog().getDbOrMetaException("test_stream");
+ for (String streamName : new String[] {"s2", "s_dup"}) {
+ OlapTableStream stream = (OlapTableStream)
db.getTableOrMetaException(streamName);
+ OlapTable baseTable = stream.getBaseTableNullable();
+ Map<Long, Pair<Long, Long>> offsets = new java.util.HashMap<>();
+ for (Partition partition : baseTable.getPartitions()) {
+ offsets.put(partition.getId(),
stream.getStreamUpdate(partition.getId()));
+ }
+ for (String mode : new String[] {"snapshot", "reset"}) {
+ for (String arguments : new String[] {"'unknown'='x'", "x",
"`x`", "x,y"}) {
+ String sql = "select * from test_stream." + streamName +
"@" + mode + "(" + arguments + ")";
+ AnalysisException error =
Assertions.assertThrows(AnalysisException.class,
+ () ->
PlanChecker.from(connectContext).analyze(sql));
+ Assertions.assertTrue(error.getMessage().contains(mode + "
does not accept parameters"),
+ sql + ": " + error.getMessage());
+ for (Partition partition : baseTable.getPartitions()) {
+ Assertions.assertEquals(offsets.get(partition.getId()),
+ stream.getStreamUpdate(partition.getId()));
+ }
+ }
+ }
+ }
+ }
+
@Test
public void testAppendOnlyStreamSnapshotCanBePlanned() throws Exception {
ConnectContext ctx = createDefaultCtx();
diff --git
a/regression-test/suites/table_stream_p0/test_table_stream_input_validation.groovy
b/regression-test/suites/table_stream_p0/test_table_stream_input_validation.groovy
new file mode 100644
index 00000000000..5cdd0e661b7
--- /dev/null
+++
b/regression-test/suites/table_stream_p0/test_table_stream_input_validation.groovy
@@ -0,0 +1,84 @@
+// 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_table_stream_input_validation") {
+ sql "DROP STREAM IF EXISTS stream_validation_initial FORCE"
+ sql "DROP STREAM IF EXISTS stream_validation_incremental FORCE"
+ sql "DROP STREAM IF EXISTS stream_validation_default FORCE"
+ sql "DROP TABLE IF EXISTS stream_validation_sink"
+ sql "DROP TABLE IF EXISTS stream_validation_source"
+
+ sql """
+ CREATE TABLE stream_validation_source (id BIGINT, value INT)
+ UNIQUE KEY(id)
+ DISTRIBUTED BY HASH(id) BUCKETS 1
+ PROPERTIES (
+ "replication_num" = "1",
+ "enable_unique_key_merge_on_write" = "true",
+ "binlog.enable" = "true",
+ "binlog.format" = "ROW",
+ "binlog.need_historical_value" = "true"
+ )
+ """
+ sql """
+ CREATE TABLE stream_validation_sink (id BIGINT, value INT)
+ DUPLICATE KEY(id)
+ DISTRIBUTED BY HASH(id) BUCKETS 1
+ PROPERTIES ("replication_num" = "1")
+ """
+ sql "INSERT INTO stream_validation_source VALUES (1, 10)"
+ sql "sync"
+
+ ["garbage", "yes", "1", "", " true "].each { value ->
+ test {
+ sql """
+ CREATE STREAM stream_validation_initial ON TABLE
stream_validation_source
+ PROPERTIES ("show_initial_rows" = "${value}")
+ """
+ exception "show_initial_rows must be `true` or `false`"
+ }
+ }
+ // Reusing the failed name proves invalid properties did not publish a
Stream.
+ sql """
+ CREATE STREAM stream_validation_initial ON TABLE
stream_validation_source
+ PROPERTIES ("show_initial_rows" = "TrUe")
+ """
+ sql """
+ CREATE STREAM stream_validation_incremental ON TABLE
stream_validation_source
+ PROPERTIES ("show_initial_rows" = "FaLsE")
+ """
+ sql "CREATE STREAM stream_validation_default ON TABLE
stream_validation_source"
+
+ ["snapshot", "reset"].each { mode ->
+ ["'unknown'='x'", "x", "`x`", "x,y"].each { arguments ->
+ test {
+ sql "SELECT * FROM
stream_validation_incremental@${mode}(${arguments}) ORDER BY id"
+ exception "${mode} does not accept parameters"
+ }
+ test {
+ sql """
+ INSERT INTO stream_validation_sink
+ SELECT id, value FROM
stream_validation_incremental@${mode}(${arguments})
+ """
+ exception "${mode} does not accept parameters"
+ }
+ }
+ }
+ sql "SELECT * FROM stream_validation_initial ORDER BY id"
+ sql "SELECT * FROM stream_validation_incremental@snapshot() ORDER BY id"
+ sql "SELECT * FROM stream_validation_incremental@reset() ORDER BY id"
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]