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]

Reply via email to