This is an automated email from the ASF dual-hosted git repository.

morningman 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 0482779cea6 [Enhance](system table) Push down predicates for 
table_stream_consumption metadata scan (#68227)
0482779cea6 is described below

commit 0482779cea68558cb7d1bf88cb601254e239777b
Author: linrrarity <[email protected]>
AuthorDate: Wed Sep 23 11:21:00 2026 +0800

    [Enhance](system table) Push down predicates for table_stream_consumption 
metadata scan (#68227)
    
    ### What problem does this PR solve?
    
    Issue Number: close #xxx
    
    Related PR: #xxx
    
    Problem Summary:
    
    Queries on `information_schema.table_stream_consumption` always scanned
    all streams and partitions, even when predicates such as `DB_NAME`,
    `STREAM_NAME`, `STREAM_ID`, or `UNIT` could exclude most of them.
    
    In cloud mode, this also caused unnecessary partition requests to
    MetaService.
    
    ### What is changed?
    
    - Forward schema-scan predicates from BE to FE for
    `table_stream_consumption`.
    - Push down `DB_NAME`, `STREAM_NAME`, and `STREAM_ID` predicates before
    scanning stream metadata.
    - Push down `UNIT` predicates after partition enumeration:
    - In cloud mode, only matching partitions are requested from
    MetaService.
      - In non-cloud mode, rows for unmatched partitions are not generated.
---
 .../schema_table_stream_consumption_scanner.cpp    |  14 +-
 .../schema_table_stream_consumption_scanner.h      |   4 +-
 ...chema_table_stream_consumption_scanner_test.cpp |  47 +++++++
 .../doris/catalog/stream/BaseTableStream.java      |   7 +-
 .../doris/catalog/stream/OlapTableStream.java      | 144 +++++++++++++-------
 .../doris/catalog/stream/TableStreamManager.java   | 101 ++++++++++++--
 .../rewrite/PushDownFilterIntoSchemaScan.java      |  10 +-
 .../trees/copier/LogicalPlanDeepCopier.java        |  11 ++
 .../doris/nereids/util/FrontendConjunctsUtils.java |  17 ++-
 .../doris/tablefunction/MetadataGenerator.java     |  70 +++++++++-
 .../stream/CloudTableStreamConsumptionTest.java    | 151 ++++++++++++++++++---
 .../stream/OlapTableStreamConsumptionTest.java     |  69 ++++++++++
 .../rules/rewrite/RewriteRuleSuiteTest.java        |  43 ++++++
 .../trees/copier/LogicalPlanDeepCopierTest.java    |  27 ++++
 .../nereids/util/FrontendConjunctsUtilsTest.java   |  30 ++++
 .../StreamConsumptionMetadataGeneratorTest.java    | 115 ++++++++++++++++
 .../test_stream_consumption_schema.groovy          |  37 ++++-
 17 files changed, 805 insertions(+), 92 deletions(-)

diff --git 
a/be/src/information_schema/schema_table_stream_consumption_scanner.cpp 
b/be/src/information_schema/schema_table_stream_consumption_scanner.cpp
index c484533ae7a..81e68b6c4f1 100644
--- a/be/src/information_schema/schema_table_stream_consumption_scanner.cpp
+++ b/be/src/information_schema/schema_table_stream_consumption_scanner.cpp
@@ -56,13 +56,21 @@ Status 
SchemaTableStreamConsumptionScanner::start(RuntimeState* state) {
     return Status::OK();
 }
 
-Status 
SchemaTableStreamConsumptionScanner::_get_table_stream_consumption_block_from_fe()
 {
-    TNetworkAddress master_addr = 
ExecEnv::GetInstance()->cluster_info()->master_fe_addr;
-
+TFetchSchemaTableDataRequest 
SchemaTableStreamConsumptionScanner::_build_fetch_request() const {
     TSchemaTableRequestParams schema_table_request_params;
+    if (_param->common_param->frontend_conjuncts) {
+        schema_table_request_params.__set_frontend_conjuncts(
+                *_param->common_param->frontend_conjuncts);
+    }
     TFetchSchemaTableDataRequest request;
     
request.__set_schema_table_name(TSchemaTableName::TABLE_STREAM_CONSUMPTION);
     request.__set_schema_table_params(schema_table_request_params);
+    return request;
+}
+
+Status 
SchemaTableStreamConsumptionScanner::_get_table_stream_consumption_block_from_fe()
 {
+    TNetworkAddress master_addr = 
ExecEnv::GetInstance()->cluster_info()->master_fe_addr;
+    TFetchSchemaTableDataRequest request = _build_fetch_request();
 
     TFetchSchemaTableDataResult result;
 
diff --git 
a/be/src/information_schema/schema_table_stream_consumption_scanner.h 
b/be/src/information_schema/schema_table_stream_consumption_scanner.h
index 9d7639cbfd8..bec6ab17a14 100644
--- a/be/src/information_schema/schema_table_stream_consumption_scanner.h
+++ b/be/src/information_schema/schema_table_stream_consumption_scanner.h
@@ -27,6 +27,7 @@ namespace doris {
 
 class RuntimeState;
 class Block;
+class TFetchSchemaTableDataRequest;
 
 class SchemaTableStreamConsumptionScanner : public SchemaScanner {
     ENABLE_FACTORY_CREATOR(SchemaTableStreamConsumptionScanner);
@@ -39,6 +40,7 @@ public:
     Status get_next_block_internal(Block* block, bool* eos) override;
 
 private:
+    TFetchSchemaTableDataRequest _build_fetch_request() const;
     Status _get_table_stream_consumption_block_from_fe();
     int _block_rows_limit = 4096;
     int _row_idx = 0;
@@ -48,4 +50,4 @@ private:
     static std::vector<SchemaScanner::ColumnDesc> 
_s_table_stream_consumption_columns;
 };
 
-} // namespace doris
\ No newline at end of file
+} // namespace doris
diff --git 
a/be/test/exec/schema_scanner/schema_table_stream_consumption_scanner_test.cpp 
b/be/test/exec/schema_scanner/schema_table_stream_consumption_scanner_test.cpp
new file mode 100644
index 00000000000..5909303a084
--- /dev/null
+++ 
b/be/test/exec/schema_scanner/schema_table_stream_consumption_scanner_test.cpp
@@ -0,0 +1,47 @@
+// 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.
+
+#include "information_schema/schema_table_stream_consumption_scanner.h"
+
+#include <gen_cpp/FrontendService_types.h>
+#include <gtest/gtest.h>
+
+#include <string>
+
+#include "common/object_pool.h"
+#include "testutil/mock/mock_runtime_state.h"
+
+namespace doris {
+
+TEST(SchemaTableStreamConsumptionScannerTest, 
forwards_frontend_conjuncts_in_fetch_request) {
+    const std::string frontend_conjuncts = 
R"([{"column":"UNIT","value":"p1"}])";
+    MockRuntimeState state;
+    SchemaScannerParam param;
+    param.common_param->frontend_conjuncts = &frontend_conjuncts;
+    ObjectPool pool;
+
+    SchemaTableStreamConsumptionScanner scanner;
+    ASSERT_TRUE(scanner.init(&state, &param, &pool).ok());
+
+    TFetchSchemaTableDataRequest request = scanner._build_fetch_request();
+    EXPECT_EQ(TSchemaTableName::TABLE_STREAM_CONSUMPTION, 
request.schema_table_name);
+    ASSERT_TRUE(request.__isset.schema_table_params);
+    ASSERT_TRUE(request.schema_table_params.__isset.frontend_conjuncts);
+    EXPECT_EQ(frontend_conjuncts, 
request.schema_table_params.frontend_conjuncts);
+}
+
+} // namespace doris
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 0111c243fea..d241d02a318 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
@@ -35,6 +35,7 @@ import java.io.DataOutput;
 import java.io.IOException;
 import java.util.List;
 import java.util.Map;
+import java.util.function.Predicate;
 
 public abstract class BaseTableStream extends Table {
     public enum StreamScanType {
@@ -212,7 +213,11 @@ public abstract class BaseTableStream extends Table {
     // fill table_stream_consumption info
     // @param dataBatch the data batch to fill
     // DB_NAME, STREAM_NAME, STREAM_ID, UNIT, CONSUMPTION_STATUS, LAG, 
LAST_CONSUMPTION_TIME
-    abstract void fillTableStreamConsumptionInfo(List<TRow> dataBatch);
+    void fillTableStreamConsumptionInfo(List<TRow> dataBatch) {
+        fillTableStreamConsumptionInfo(dataBatch, unit -> true);
+    }
+
+    abstract void fillTableStreamConsumptionInfo(List<TRow> dataBatch, 
Predicate<String> unitSelector);
 
     public <E extends Exception> TableIf 
getBaseTableOrException(java.util.function.Function<String, E> e)
             throws E {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/OlapTableStream.java 
b/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/OlapTableStream.java
index ec6fde1a97f..3de5a5ee58d 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/OlapTableStream.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/OlapTableStream.java
@@ -37,12 +37,13 @@ import com.google.gson.annotations.SerializedName;
 import java.io.DataInput;
 import java.io.DataOutput;
 import java.io.IOException;
+import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
-import java.util.stream.Collectors;
+import java.util.function.Predicate;
 
 public class OlapTableStream extends BaseTableStream {
 
@@ -165,52 +166,14 @@ public class OlapTableStream extends BaseTableStream {
         }
         if (table.readLockIfExist()) {
             try {
-                Map<Long, Partition> id2name = 
table.getPartitions().stream().collect(Collectors.toMap(
-                        p -> p.getId(),
-                        p -> p,
-                        (oldValue, newValue) -> newValue,
-                        HashMap::new
-                ));
-                for (Map.Entry<Long, Partition> entry : id2name.entrySet()) {
-                    TRow trow = new TRow();
-                    // DB_NAME
-                    trow.addToColumnValue(new 
TCell().setStringVal(qualifiedDbName));
-                    // STREAM_NAME
-                    trow.addToColumnValue(new TCell().setStringVal(name));
-                    // STREAM_ID
-                    trow.addToColumnValue(new TCell().setLongVal(id));
-                    // UNIT
-                    trow.addToColumnValue(new 
TCell().setStringVal(entry.getValue().getName()));
-                    if (partitionOffset.containsKey(entry.getKey())) {
-                        // CONSUMPTION_STATUS
-                        trow.addToColumnValue(new TCell()
-                                
.setStringVal(String.valueOf(partitionOffset.get(entry.getKey()))));
-                        // LAG
-                        trow.addToColumnValue(new TCell()
-                                .setStringVal(String.valueOf(
-                                        entry.getValue().getTso()
-                                                - 
partitionOffset.get(entry.getKey()))));
-                        // LAST_CONSUMPTION_TIME
-                        if 
(partitionConsumptionTime.containsKey(entry.getKey())) {
-                            trow.addToColumnValue(new TCell()
-                                    
.setLongVal(partitionConsumptionTime.get(entry.getKey())));
-                        } else {
-                            trow.addToColumnValue(new TCell().setLongVal(-1));
-                        }
-                    } else {
-                        // CONSUMPTION_STATUS
-                        trow.addToColumnValue(new TCell().setStringVal("N/A"));
-                        // LAG
-                        if (entry.getValue().hasData()) {
-                            // for partition with data and no consumption yet, 
lag is N/A
-                            trow.addToColumnValue(new 
TCell().setStringVal("N/A"));
-                        } else {
-                            trow.addToColumnValue(new 
TCell().setStringVal("0"));
-                        }
-                        // LAST_CONSUMPTION_TIME
-                        trow.addToColumnValue(new TCell().setLongVal(-1));
-                    }
-                    dataBatch.add(trow);
+                for (Partition partition : table.getPartitions()) {
+                    long partitionId = partition.getId();
+                    boolean hasOffset = 
partitionOffset.containsKey(partitionId);
+                    appendConsumptionRow(dataBatch, qualifiedDbName, name, id, 
partition.getName(), hasOffset,
+                            hasOffset ? partitionOffset.get(partitionId) : 0,
+                            hasOffset ? partition.getTso() : 0,
+                            !hasOffset && partition.hasData(),
+                            partitionConsumptionTime.getOrDefault(partitionId, 
-1L));
                 }
             } finally {
                 table.readUnlock();
@@ -218,6 +181,93 @@ public class OlapTableStream extends BaseTableStream {
         }
     }
 
+    @Override
+    void fillTableStreamConsumptionInfo(List<TRow> dataBatch, 
Predicate<String> unitSelector) {
+        for (StreamConsumptionUnitSnapshot snapshot : 
snapshotTableStreamConsumptionInfo()) {
+            if (unitSelector.test(snapshot.unit)) {
+                snapshot.appendRow(dataBatch, qualifiedDbName, name, id);
+            }
+        }
+    }
+
+    List<StreamConsumptionUnitSnapshot> snapshotTableStreamConsumptionInfo() {
+        // Copy row inputs under lock so UNIT expression rewriting and folding 
can run after unlocking.
+        OlapTable table = getBaseTableNullable();
+        if (table == null) {
+            return ImmutableList.of();
+        }
+        List<StreamConsumptionUnitSnapshot> snapshots = new ArrayList<>();
+        if (table.readLockIfExist()) {
+            try {
+                for (Partition partition : table.getPartitions()) {
+                    snapshots.add(snapshotPartition(partition));
+                }
+            } finally {
+                table.readUnlock();
+            }
+        }
+        return snapshots;
+    }
+
+    private StreamConsumptionUnitSnapshot snapshotPartition(Partition 
partition) {
+        long partitionId = partition.getId();
+        boolean hasOffset = partitionOffset.containsKey(partitionId);
+        return new StreamConsumptionUnitSnapshot(
+                partition.getName(), hasOffset,
+                hasOffset ? partitionOffset.get(partitionId) : 0,
+                hasOffset ? partition.getTso() : 0,
+                !hasOffset && partition.hasData(),
+                partitionConsumptionTime.getOrDefault(partitionId, -1L));
+    }
+
+    private static void appendConsumptionRow(List<TRow> dataBatch, String 
dbName, String streamName,
+            long streamId, String unit, boolean hasOffset, long offset, long 
endTso, boolean hasData,
+            long lastConsumptionTime) {
+        TRow row = new TRow();
+        row.addToColumnValue(new TCell().setStringVal(dbName));
+        row.addToColumnValue(new TCell().setStringVal(streamName));
+        row.addToColumnValue(new TCell().setLongVal(streamId));
+        row.addToColumnValue(new TCell().setStringVal(unit));
+        if (hasOffset) {
+            row.addToColumnValue(new 
TCell().setStringVal(String.valueOf(offset)));
+            row.addToColumnValue(new 
TCell().setStringVal(String.valueOf(endTso - offset)));
+            row.addToColumnValue(new TCell().setLongVal(lastConsumptionTime));
+        } else {
+            row.addToColumnValue(new TCell().setStringVal("N/A"));
+            row.addToColumnValue(new TCell().setStringVal(hasData ? "N/A" : 
"0"));
+            row.addToColumnValue(new TCell().setLongVal(-1));
+        }
+        dataBatch.add(row);
+    }
+
+    static class StreamConsumptionUnitSnapshot {
+        private final String unit;
+        private final boolean hasOffset;
+        private final long offset;
+        private final long endTso;
+        private final boolean hasData;
+        private final long lastConsumptionTime;
+
+        private StreamConsumptionUnitSnapshot(String unit, boolean hasOffset, 
long offset, long endTso,
+                boolean hasData, long lastConsumptionTime) {
+            this.unit = unit;
+            this.hasOffset = hasOffset;
+            this.offset = offset;
+            this.endTso = endTso;
+            this.hasData = hasData;
+            this.lastConsumptionTime = lastConsumptionTime;
+        }
+
+        String getUnit() {
+            return unit;
+        }
+
+        void appendRow(List<TRow> dataBatch, String dbName, String streamName, 
long streamId) {
+            appendConsumptionRow(dataBatch, dbName, streamName, streamId, 
unit, hasOffset, offset, endTso,
+                    hasData, lastConsumptionTime);
+        }
+    }
+
     public boolean hasData(Partition partition) {
         // if all available visible data has been consumed, return false
         return  (!partitionOffset.containsKey(partition.getId())
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/TableStreamManager.java
 
b/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/TableStreamManager.java
index 35275748397..a3e1d4e9a7b 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/TableStreamManager.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/TableStreamManager.java
@@ -64,6 +64,18 @@ import java.util.concurrent.locks.LockSupport;
 public class TableStreamManager extends MasterDaemon implements Writable, 
GsonPostProcessable {
     private static final Logger LOG = 
LogManager.getLogger(TableStreamManager.class);
     private static final String BASE_TABLE_NOT_FOUND_STALE_REASON = "Base 
table does not exist";
+
+    @FunctionalInterface
+    public interface StreamConsumptionSelector {
+        // unit is null while selecting a stream, and is set after its 
base-table partitions are known.
+        boolean test(String dbName, String streamName, long streamId, String 
unit);
+
+        // Callers override this to skip partition snapshots when no UNIT 
predicate exists.
+        default boolean hasUnitFilter() {
+            return true;
+        }
+    }
+
     @SerializedName(value = "dbStreamMap")
     private Map<Long, Set<Long>> dbStreamMap;
     protected MonitoredReentrantReadWriteLock rwLock;
@@ -398,9 +410,24 @@ public class TableStreamManager extends MasterDaemon 
implements Writable, GsonPo
 
     public void fillStreamConsumptionValuesMetadataResult(List<TRow> 
dataBatch) throws UserException {
         if (Config.isCloudMode()) {
-            fillCloudStreamConsumptionValuesMetadataResult(dataBatch);
+            fillCloudStreamConsumptionValuesMetadataResult(dataBatch, null);
             return;
         }
+        fillLocalStreamConsumptionValuesMetadataResult(dataBatch, null);
+    }
+
+    public void fillStreamConsumptionValuesMetadataResult(List<TRow> dataBatch,
+            StreamConsumptionSelector selector) throws UserException {
+        if (Config.isCloudMode()) {
+            fillCloudStreamConsumptionValuesMetadataResult(dataBatch, 
selector);
+            return;
+        }
+        fillLocalStreamConsumptionValuesMetadataResult(dataBatch, selector);
+    }
+
+    private void fillLocalStreamConsumptionValuesMetadataResult(List<TRow> 
dataBatch,
+            StreamConsumptionSelector selector) {
+        // Resolve registered streams and apply stream-level predicates before 
taking metadata locks.
         for (Map.Entry<Long, Set<Long>> entry : copyDbStreamMap().entrySet()) {
             Optional<Database> db = 
Env.getCurrentInternalCatalog().getDb(entry.getKey());
             if (db.isPresent()) {
@@ -412,22 +439,43 @@ public class TableStreamManager extends MasterDaemon 
implements Writable, GsonPo
                         }
                         continue;
                     }
-                    Preconditions.checkArgument(table.get() instanceof 
BaseTableStream);
-                    BaseTableStream stream = (BaseTableStream) table.get();
+                    Preconditions.checkArgument(table.get() instanceof 
OlapTableStream);
+                    OlapTableStream stream = (OlapTableStream) table.get();
+                    String dbName = db.get().getFullName();
+                    String streamName = stream.getName();
+                    long streamId = stream.getId();
+                    if (selector != null && !selector.test(dbName, streamName, 
streamId, null)) {
+                        continue;
+                    }
+                    List<OlapTableStream.StreamConsumptionUnitSnapshot> 
snapshots = Collections.emptyList();
                     if (stream.readLockIfExist()) {
                         try {
-                            stream.fillTableStreamConsumptionInfo(dataBatch);
+                            // Build rows directly unless UNIT evaluation 
requires a lock-protected metadata snapshot.
+                            if (selector == null || !selector.hasUnitFilter()) 
{
+                                
stream.fillTableStreamConsumptionInfo(dataBatch);
+                            } else {
+                                // UNIT evaluation rewrites and folds 
expressions, so defer it until after unlocking.
+                                snapshots = 
stream.snapshotTableStreamConsumptionInfo();
+                            }
                         } finally {
                             stream.readUnlock();
                         }
                     }
+                    // Both metadata locks are released here; materialize only 
units accepted by the selector.
+                    for (OlapTableStream.StreamConsumptionUnitSnapshot 
snapshot : snapshots) {
+                        if (selector.test(dbName, streamName, streamId, 
snapshot.getUnit())) {
+                            snapshot.appendRow(dataBatch, dbName, streamName, 
streamId);
+                        }
+                    }
                 }
             }
         }
     }
 
-    private void fillCloudStreamConsumptionValuesMetadataResult(List<TRow> 
dataBatch) throws UserException {
+    private void fillCloudStreamConsumptionValuesMetadataResult(List<TRow> 
dataBatch,
+            StreamConsumptionSelector selector) throws UserException {
         Map<Cloud.TableStreamIdentityPB, CloudStreamConsumptionSnapshot> 
snapshots = new LinkedHashMap<>();
+        // Resolve registered streams and apply stream-level predicates before 
taking metadata locks.
         for (Map.Entry<Long, Set<Long>> entry : copyDbStreamMap().entrySet()) {
             Optional<Database> db = 
Env.getCurrentInternalCatalog().getDb(entry.getKey());
             if (!db.isPresent()) {
@@ -440,44 +488,59 @@ public class TableStreamManager extends MasterDaemon 
implements Writable, GsonPo
                 }
                 Preconditions.checkArgument(table.get() instanceof 
OlapTableStream);
                 OlapTableStream stream = (OlapTableStream) table.get();
+                if (selector != null
+                        && !selector.test(db.get().getFullName(), 
stream.getName(), stream.getId(), null)) {
+                    continue;
+                }
                 if (!stream.readLockIfExist()) {
                     continue;
                 }
+                Cloud.TableStreamIdentityPB identity;
+                CloudStreamConsumptionSnapshot snapshot;
                 try {
+                    // Snapshot partition identities under lock; authoritative 
consumption state lives in MetaService.
                     OlapTable baseTable = stream.getBaseTableNullable();
                     if (baseTable == null || !baseTable.readLockIfExist()) {
                         continue;
                     }
                     try {
                         Map<Long, String> partitionNames = new 
LinkedHashMap<>();
-                        baseTable.getPartitions().forEach(partition ->
-                                partitionNames.put(partition.getId(), 
partition.getName()));
+                        baseTable.getPartitions().stream()
+                                .forEach(partition -> 
partitionNames.put(partition.getId(), partition.getName()));
                         if (partitionNames.isEmpty()) {
                             continue;
                         }
-                        Cloud.TableStreamIdentityPB identity = 
Cloud.TableStreamIdentityPB.newBuilder()
+                        identity = Cloud.TableStreamIdentityPB.newBuilder()
                                 
.setBaseDbId(stream.getBaseTableInfo().getDbId())
                                 
.setBaseTableId(stream.getBaseTableInfo().getTableId())
                                 .setStreamDbId(entry.getKey())
                                 .setStreamId(stream.getId())
                                 .build();
-                        CloudStreamConsumptionSnapshot previous = 
snapshots.put(identity,
-                                new 
CloudStreamConsumptionSnapshot(db.get().getFullName(), stream.getName(),
-                                        stream.getId(), partitionNames));
-                        Preconditions.checkState(previous == null,
-                                "Duplicate Cloud Table Stream identity %s", 
identity);
+                        snapshot = new 
CloudStreamConsumptionSnapshot(db.get().getFullName(), stream.getName(),
+                                stream.getId(), partitionNames);
                     } finally {
                         baseTable.readUnlock();
                     }
                 } finally {
                     stream.readUnlock();
                 }
+                // Evaluate UNIT expressions after unlocking, then retain only 
selected units in the global map.
+                if (selector != null && selector.hasUnitFilter()) {
+                    snapshot = snapshot.selectUnits(selector);
+                    if (snapshot == null) {
+                        continue;
+                    }
+                }
+                CloudStreamConsumptionSnapshot previous = 
snapshots.put(identity, snapshot);
+                Preconditions.checkState(previous == null,
+                        "Duplicate Cloud Table Stream identity %s", identity);
             }
         }
         if (snapshots.isEmpty()) {
             return;
         }
 
+        // Fetch authoritative states for selected units in one batch, then 
materialize the result rows.
         Map<Cloud.TableStreamIdentityPB, Set<Long>> requestedPartitions = new 
LinkedHashMap<>();
         snapshots.forEach((identity, snapshot) ->
                 requestedPartitions.put(identity, 
snapshot.partitionNames.keySet()));
@@ -503,6 +566,18 @@ public class TableStreamManager extends MasterDaemon 
implements Writable, GsonPo
             this.partitionNames = Collections.unmodifiableMap(partitionNames);
         }
 
+        private CloudStreamConsumptionSnapshot 
selectUnits(StreamConsumptionSelector selector) {
+            Map<Long, String> selectedPartitions = new LinkedHashMap<>();
+            for (Map.Entry<Long, String> entry : partitionNames.entrySet()) {
+                if (selector.test(dbName, streamName, streamId, 
entry.getValue())) {
+                    selectedPartitions.put(entry.getKey(), entry.getValue());
+                }
+            }
+            return selectedPartitions.isEmpty()
+                    ? null
+                    : new CloudStreamConsumptionSnapshot(dbName, streamName, 
streamId, selectedPartitions);
+        }
+
         private void fillRows(Map<Long, Cloud.TableStreamPartitionReadStatePB> 
partitionStates,
                 List<TRow> dataBatch) throws UserException {
             for (Map.Entry<Long, String> entry : partitionNames.entrySet()) {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/PushDownFilterIntoSchemaScan.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/PushDownFilterIntoSchemaScan.java
index 0d4f2fbcee3..36af32b9ff3 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/PushDownFilterIntoSchemaScan.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/PushDownFilterIntoSchemaScan.java
@@ -31,6 +31,7 @@ import org.apache.doris.nereids.trees.expressions.Not;
 import org.apache.doris.nereids.trees.expressions.NullSafeEqual;
 import org.apache.doris.nereids.trees.expressions.Or;
 import org.apache.doris.nereids.trees.expressions.SlotReference;
+import 
org.apache.doris.nereids.trees.expressions.functions.NoneMovableFunction;
 import org.apache.doris.nereids.trees.expressions.literal.VarcharLiteral;
 import org.apache.doris.nereids.trees.plans.logical.LogicalFilter;
 import org.apache.doris.nereids.trees.plans.logical.LogicalSchemaScan;
@@ -48,13 +49,20 @@ import java.util.Optional;
 public class PushDownFilterIntoSchemaScan extends OneRewriteRuleFactory {
 
     public static ImmutableSet<String> SUPPOPRT_FRONTEND_CONJUNCTS_TABLES =
-            ImmutableSet.of("view_dependency", "sql_block_rule_status");
+            ImmutableSet.of("view_dependency", "sql_block_rule_status", 
"table_stream_consumption");
 
     @Override
     public Rule build() {
         return logicalFilter(logicalSchemaScan()).when(p -> 
!p.child().isFilterPushed()).thenApply(ctx -> {
             LogicalFilter<LogicalSchemaScan> filter = ctx.root;
             LogicalSchemaScan scan = filter.child();
+            if (filter.getConjuncts().stream()
+                    .anyMatch(expression -> 
expression.containsType(NoneMovableFunction.class))) {
+                // A pushed sibling could prune every row and skip required 
evaluation such as assert_true.
+                LogicalSchemaScan rewrittenScan = scan.withFrontendConjuncts(
+                        Optional.empty(), Optional.empty(), Optional.empty(), 
ImmutableList.of());
+                return filter.withChildren(ImmutableList.of(rewrittenScan));
+            }
             List<Optional<String>> fixedFilter = getFixedFilter(filter);
             List<Expression> commonFilter = getCommonFilter(filter);
             LogicalSchemaScan rewrittenScan = 
scan.withFrontendConjuncts(fixedFilter.get(0), fixedFilter.get(1),
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/copier/LogicalPlanDeepCopier.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/copier/LogicalPlanDeepCopier.java
index 01472f1c920..abb9e015a27 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/copier/LogicalPlanDeepCopier.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/copier/LogicalPlanDeepCopier.java
@@ -56,6 +56,7 @@ import 
org.apache.doris.nereids.trees.plans.logical.LogicalRecursiveUnionAnchor;
 import 
org.apache.doris.nereids.trees.plans.logical.LogicalRecursiveUnionProducer;
 import org.apache.doris.nereids.trees.plans.logical.LogicalRelation;
 import org.apache.doris.nereids.trees.plans.logical.LogicalRepeat;
+import org.apache.doris.nereids.trees.plans.logical.LogicalSchemaScan;
 import org.apache.doris.nereids.trees.plans.logical.LogicalSink;
 import org.apache.doris.nereids.trees.plans.logical.LogicalSort;
 import org.apache.doris.nereids.trees.plans.logical.LogicalTopN;
@@ -128,6 +129,16 @@ public class LogicalPlanDeepCopier extends 
DefaultPlanRewriter<DeepCopierContext
                 .map(o -> (NamedExpression) 
ExpressionDeepCopier.INSTANCE.deepCopy(o, context))
                 .collect(ImmutableList.toImmutableList());
         newRelation = newRelation.withVirtualColumns(virtualColumns);
+        if (catalogRelation instanceof LogicalSchemaScan
+                && ((LogicalSchemaScan) catalogRelation).isFilterPushed()) {
+            LogicalSchemaScan oldSchemaScan = (LogicalSchemaScan) 
catalogRelation;
+            List<Expression> frontendConjuncts = 
oldSchemaScan.getFrontendConjuncts().stream()
+                    .map(expression -> 
ExpressionDeepCopier.INSTANCE.deepCopy(expression, context))
+                    .collect(ImmutableList.toImmutableList());
+            newRelation = ((LogicalSchemaScan) 
newRelation).withFrontendConjuncts(
+                    oldSchemaScan.getSchemaCatalog(), 
oldSchemaScan.getSchemaDatabase(), oldSchemaScan.getSchemaTable(),
+                    frontendConjuncts);
+        }
         context.putRelation(catalogRelation.getRelationId(), newRelation);
         return updateOperativeSlots(catalogRelation, newRelation);
     }
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/util/FrontendConjunctsUtils.java
 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/util/FrontendConjunctsUtils.java
index 0c3aa7b56c4..08c2a7bd81c 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/nereids/util/FrontendConjunctsUtils.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/nereids/util/FrontendConjunctsUtils.java
@@ -19,6 +19,7 @@ package org.apache.doris.nereids.util;
 
 import org.apache.doris.analysis.Expr;
 import org.apache.doris.analysis.ExprToSqlVisitor;
+import org.apache.doris.analysis.IntLiteral;
 import org.apache.doris.analysis.ToSqlParams;
 import org.apache.doris.nereids.analyzer.UnboundSlot;
 import org.apache.doris.nereids.parser.NereidsParser;
@@ -45,6 +46,12 @@ import java.util.stream.Collectors;
  */
 public class FrontendConjunctsUtils {
     private static final Logger LOG = 
LogManager.getLogger(FrontendConjunctsUtils.class);
+    private static final ExprToSqlVisitor FRONTEND_CONJUNCTS_TO_SQL_VISITOR = 
new ExprToSqlVisitor() {
+        @Override
+        public String visitIntLiteral(IntLiteral expr, ToSqlParams context) {
+            return "CAST(" + expr.getStringValue() + " AS " + 
expr.getType().toSql() + ")";
+        }
+    };
     private static List<String> nameParts;
 
     public static List<Expression> convertToExpression(String conjuncts) {
@@ -57,7 +64,7 @@ public class FrontendConjunctsUtils {
 
     public static Expression exprToExpression(Expr expr) {
         NereidsParser nereidsParser = new NereidsParser();
-        return 
nereidsParser.parseExpression(expr.accept(ExprToSqlVisitor.INSTANCE, 
ToSqlParams.WITH_TABLE));
+        return 
nereidsParser.parseExpression(expr.accept(FRONTEND_CONJUNCTS_TO_SQL_VISITOR, 
ToSqlParams.WITH_TABLE));
     }
 
     /**
@@ -117,10 +124,11 @@ public class FrontendConjunctsUtils {
      *
      * @param expression expression
      * @param values case insensitive map
-     * @return isFiltered
+     * @return true if the candidate can be proven not to match the expression
      */
     public static boolean isFiltered(Expression expression, TreeMap<String, 
Object> values) {
         try {
+            // Bind referenced slots to this candidate's values, producing a 
temporary constant expression.
             AtomicBoolean containsAllColumn = new AtomicBoolean(true);
             Expression rewrittenExpr = expression.rewriteUp(expr -> {
                 if (expr instanceof UnboundSlot) {
@@ -137,12 +145,11 @@ public class FrontendConjunctsUtils {
                 }
                 return expr;
             });
-            // expression is: c1=v1 or c2=v2,
-            // if values is {c1=v3}
-            // we should not return true, because c2 may equals v2
+            // Missing values make the predicate undecidable here, so preserve 
the candidate for downstream filtering.
             if (!containsAllColumn.get()) {
                 return false;
             }
+            // SQL WHERE rejects FALSE and NULL; any other result keeps the 
candidate.
             Expression evaluate = FoldConstantRuleOnFE.evaluate(rewrittenExpr, 
null);
             if (evaluate instanceof BooleanLiteral && !((BooleanLiteral) 
evaluate).getValue()) {
                 return true;
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/MetadataGenerator.java
 
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/MetadataGenerator.java
index 773b39dea65..475b625406e 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/tablefunction/MetadataGenerator.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/tablefunction/MetadataGenerator.java
@@ -48,6 +48,7 @@ import org.apache.doris.catalog.Type;
 import org.apache.doris.catalog.View;
 import org.apache.doris.catalog.info.TableNameInfo;
 import org.apache.doris.catalog.stream.BaseTableStream;
+import org.apache.doris.catalog.stream.TableStreamManager;
 import org.apache.doris.common.AnalysisException;
 import org.apache.doris.common.ClientPool;
 import org.apache.doris.common.Config;
@@ -91,6 +92,7 @@ import org.apache.doris.mtmv.MTMVRelation;
 import org.apache.doris.mtmv.MTMVStatus;
 import org.apache.doris.mtmv.ivm.IvmUtil;
 import org.apache.doris.mysql.privilege.PrivPredicate;
+import org.apache.doris.nereids.analyzer.UnboundSlot;
 import org.apache.doris.nereids.trees.expressions.Expression;
 import org.apache.doris.nereids.util.FrontendConjunctsUtils;
 import org.apache.doris.nereids.util.PlanUtils;
@@ -142,17 +144,23 @@ import org.jetbrains.annotations.NotNull;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.Collections;
+import java.util.HashSet;
 import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Locale;
 import java.util.Map;
 import java.util.Optional;
 import java.util.Set;
+import java.util.TreeMap;
 import java.util.concurrent.TimeUnit;
 import java.util.stream.Stream;
 
 public class MetadataGenerator {
     private static final Logger LOG = 
LogManager.getLogger(MetadataGenerator.class);
+    private static final Set<String> STREAM_CONSUMPTION_STREAM_COLUMNS =
+            Set.of("DB_NAME", "STREAM_NAME", "STREAM_ID");
+    private static final Set<String> STREAM_CONSUMPTION_SELECTOR_COLUMNS =
+            Set.of("DB_NAME", "STREAM_NAME", "STREAM_ID", "UNIT");
 
     private static final ImmutableMap<String, Integer> 
ACTIVE_QUERIES_COLUMN_TO_INDEX;
 
@@ -2329,8 +2337,68 @@ public class MetadataGenerator {
     private static TFetchSchemaTableDataResult 
streamConsumptionMetadataResult(TSchemaTableRequestParams params) {
         TFetchSchemaTableDataResult result = new TFetchSchemaTableDataResult();
         List<TRow> dataBatch = Lists.newArrayList();
+        // Decode the planner predicates carried back by BE. Conversion 
failures only disable FE pruning.
+        List<Expression> parsedConjuncts = Collections.emptyList();
+        if (params.isSetFrontendConjuncts()) {
+            try {
+                parsedConjuncts = 
FrontendConjunctsUtils.convertToExpression(params.getFrontendConjuncts());
+            } catch (RuntimeException e) {
+                LOG.warn("Failed to convert frontend conjuncts for 
table_stream_consumption; skip FE pruning", e);
+            }
+        }
+        List<Expression> conjuncts = parsedConjuncts;
         try {
-            
Env.getCurrentEnv().getTableStreamManager().fillStreamConsumptionValuesMetadataResult(dataBatch);
+            // Keep unfiltered scans on the direct path without allocating a 
selector or partition snapshots.
+            if (conjuncts.isEmpty()) {
+                
Env.getCurrentEnv().getTableStreamManager().fillStreamConsumptionValuesMetadataResult(dataBatch);
+                result.setDataBatch(dataBatch);
+                result.setStatus(new TStatus(TStatusCode.OK));
+                return result;
+            }
+            // Split predicates by the earliest metadata level that has every 
referenced column.
+            // Stream-only predicates run once per stream; predicates using 
UNIT run once per partition.
+            List<Expression> streamConjuncts = Lists.newArrayList();
+            List<Expression> unitConjuncts = Lists.newArrayList();
+            for (Expression conjunct : conjuncts) {
+                Set<String> referencedColumns = new HashSet<>();
+                for (UnboundSlot slot : 
conjunct.<UnboundSlot>collectToList(UnboundSlot.class::isInstance)) {
+                    List<String> nameParts = slot.getNameParts();
+                    if (!nameParts.isEmpty()) {
+                        referencedColumns.add(nameParts.get(nameParts.size() - 
1).toUpperCase(Locale.ROOT));
+                    }
+                }
+                // containsAll allows any subset, but rejects predicates that 
need unavailable columns.
+                if 
(STREAM_CONSUMPTION_STREAM_COLUMNS.containsAll(referencedColumns)) {
+                    streamConjuncts.add(conjunct);
+                } else if (referencedColumns.contains("UNIT")
+                        && 
STREAM_CONSUMPTION_SELECTOR_COLUMNS.containsAll(referencedColumns)) {
+                    unitConjuncts.add(conjunct);
+                }
+            }
+            // Bind each candidate's metadata values in the selector; 
unsupported predicates remain for BE filtering.
+            
Env.getCurrentEnv().getTableStreamManager().fillStreamConsumptionValuesMetadataResult(
+                    dataBatch, new 
TableStreamManager.StreamConsumptionSelector() {
+                        @Override
+                        public boolean test(String dbName, String streamName, 
long streamId, String unit) {
+                            List<Expression> currentConjuncts = unit == null ? 
streamConjuncts : unitConjuncts;
+                            if (currentConjuncts.isEmpty()) {
+                                return true;
+                            }
+                            TreeMap<String, Object> values = new 
TreeMap<>(String.CASE_INSENSITIVE_ORDER);
+                            values.put("DB_NAME", dbName);
+                            values.put("STREAM_NAME", streamName);
+                            values.put("STREAM_ID", streamId);
+                            if (unit != null) {
+                                values.put("UNIT", unit);
+                            }
+                            return 
!FrontendConjunctsUtils.isFiltered(currentConjuncts, values);
+                        }
+
+                        @Override
+                        public boolean hasUnitFilter() {
+                            return !unitConjuncts.isEmpty();
+                        }
+                    });
         } catch (UserException e) {
             return errorResult(e.getMessage());
         }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/CloudTableStreamConsumptionTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/CloudTableStreamConsumptionTest.java
index 60110061711..8f752dc8faa 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/CloudTableStreamConsumptionTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/CloudTableStreamConsumptionTest.java
@@ -17,6 +17,10 @@
 
 package org.apache.doris.catalog.stream;
 
+import org.apache.doris.analysis.BinaryPredicate;
+import org.apache.doris.analysis.IntLiteral;
+import org.apache.doris.analysis.SlotRef;
+import org.apache.doris.analysis.StringLiteral;
 import org.apache.doris.catalog.Database;
 import org.apache.doris.catalog.Env;
 import org.apache.doris.catalog.OlapTable;
@@ -24,7 +28,16 @@ import org.apache.doris.cloud.proto.Cloud;
 import org.apache.doris.cloud.rpc.MetaServiceProxy;
 import org.apache.doris.common.Config;
 import org.apache.doris.common.FeConstants;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.common.lock.MonitoredReentrantReadWriteLock;
+import org.apache.doris.persist.gson.GsonUtils;
+import org.apache.doris.tablefunction.MetadataGenerator;
+import org.apache.doris.thrift.TFetchSchemaTableDataRequest;
+import org.apache.doris.thrift.TFetchSchemaTableDataResult;
 import org.apache.doris.thrift.TRow;
+import org.apache.doris.thrift.TSchemaTableName;
+import org.apache.doris.thrift.TSchemaTableRequestParams;
+import org.apache.doris.thrift.TStatusCode;
 import org.apache.doris.utframe.TestWithFeService;
 
 import org.junit.jupiter.api.Assertions;
@@ -37,6 +50,7 @@ import java.util.Comparator;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Set;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 public class CloudTableStreamConsumptionTest extends TestWithFeService {
 
@@ -56,6 +70,9 @@ public class CloudTableStreamConsumptionTest extends 
TestWithFeService {
         createTable("create stream test_cloud_stream_consumption.s1 "
                 + "on table test_cloud_stream_consumption.base_table "
                 + "properties('show_initial_rows'='true')");
+        createTable("create stream test_cloud_stream_consumption.s2 "
+                + "on table test_cloud_stream_consumption.base_table "
+                + "properties('show_initial_rows'='true')");
         createTable("create table 
test_cloud_stream_consumption.empty_base_table (k1 int, k2 int) "
                 + "unique key(k1) partition by range(k1) "
                 + "(partition p1 values less than (100)) "
@@ -81,6 +98,8 @@ public class CloudTableStreamConsumptionTest extends 
TestWithFeService {
         Database db = (Database) Env.getCurrentInternalCatalog()
                 .getDbOrMetaException("test_cloud_stream_consumption");
         OlapTable table = (OlapTable) db.getTableOrMetaException("base_table");
+        OlapTableStream stream = (OlapTableStream) 
db.getTableOrMetaException("s1");
+        OlapTableStream secondStream = (OlapTableStream) 
db.getTableOrMetaException("s2");
         long p1 = table.getPartition("p1").getId();
         long p2 = table.getPartition("p2").getId();
 
@@ -93,8 +112,52 @@ public class CloudTableStreamConsumptionTest extends 
TestWithFeService {
             mockedProxy.when(MetaServiceProxy::getInstance).thenReturn(proxy);
             
Mockito.when(proxy.getTableStreamOffset(Mockito.any())).thenAnswer(invocation 
-> {
                 Cloud.GetTableStreamOffsetRequest request = 
invocation.getArgument(0);
+                if (request.getBindingsCount() == 2) {
+                    Assertions.assertEquals(Set.of(stream.getId(), 
secondStream.getId()),
+                            request.getBindingsList().stream()
+                                    .map(binding -> 
binding.getIdentity().getStreamId())
+                                    
.collect(java.util.stream.Collectors.toSet()));
+                    Cloud.GetTableStreamOffsetResponse.Builder response = 
Cloud.GetTableStreamOffsetResponse
+                            
.newBuilder().setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
+                                    .setCode(Cloud.MetaServiceCode.OK));
+                    request.getBindingsList().forEach(binding -> {
+                        Assertions.assertEquals(Set.of(p1, p2), new 
HashSet<>(binding.getPartitionIdsList()));
+                        Cloud.TableStreamReadBindingResultPB.Builder 
bindingResult =
+                                
Cloud.TableStreamReadBindingResultPB.newBuilder()
+                                        .setIdentity(binding.getIdentity());
+                        binding.getPartitionIdsList().forEach(partitionId -> 
bindingResult.addPartitionStates(
+                                
Cloud.TableStreamPartitionReadStatePB.newBuilder()
+                                        .setPartitionId(partitionId)
+                                        
.setOffsetState(Cloud.TableStreamOffsetStatePB.TABLE_STREAM_OFFSET_UNKNOWN)
+                                        .setEndTso(200)
+                                        .setVisibleVersion(1)));
+                        response.addBindings(bindingResult);
+                    });
+                    return response.build();
+                }
                 Assertions.assertEquals(1, request.getBindingsCount());
-                Assertions.assertEquals(Set.of(p1, p2),
+                if (request.getBindings(0).getIdentity().getStreamId() == 
secondStream.getId()) {
+                    Assertions.assertEquals(Set.of(p1, p2),
+                            new 
HashSet<>(request.getBindings(0).getPartitionIdsList()));
+                    return Cloud.GetTableStreamOffsetResponse.newBuilder()
+                            
.setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
+                                    .setCode(Cloud.MetaServiceCode.OK))
+                            
.addBindings(Cloud.TableStreamReadBindingResultPB.newBuilder()
+                                    
.setIdentity(request.getBindings(0).getIdentity())
+                                    
.addPartitionStates(Cloud.TableStreamPartitionReadStatePB.newBuilder()
+                                            .setPartitionId(p1)
+                                            
.setOffsetState(Cloud.TableStreamOffsetStatePB.TABLE_STREAM_OFFSET_UNKNOWN)
+                                            .setEndTso(200)
+                                            .setVisibleVersion(1))
+                                    
.addPartitionStates(Cloud.TableStreamPartitionReadStatePB.newBuilder()
+                                            .setPartitionId(p2)
+                                            
.setOffsetState(Cloud.TableStreamOffsetStatePB.TABLE_STREAM_OFFSET_UNKNOWN)
+                                            .setEndTso(200)
+                                            .setVisibleVersion(1)))
+                            .build();
+                }
+                Assertions.assertEquals(stream.getId(), 
request.getBindings(0).getIdentity().getStreamId());
+                Assertions.assertEquals(Set.of(p1),
                         new 
HashSet<>(request.getBindings(0).getPartitionIdsList()));
                 return Cloud.GetTableStreamOffsetResponse.newBuilder()
                         .setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
@@ -107,27 +170,51 @@ public class CloudTableStreamConsumptionTest extends 
TestWithFeService {
                                         .setOffsetTso(100)
                                         .setEndTso(130)
                                         .setVisibleVersion(8)
-                                        .setLastConsumptionTimeMs(999))
-                                
.addPartitionStates(Cloud.TableStreamPartitionReadStatePB.newBuilder()
-                                        .setPartitionId(p2)
-                                        
.setOffsetState(Cloud.TableStreamOffsetStatePB.TABLE_STREAM_OFFSET_UNKNOWN)
-                                        .setEndTso(200)
-                                        .setVisibleVersion(1)))
-                        .build();
+                                        .setLastConsumptionTimeMs(999)))
+                                .build();
             });
 
-            List<TRow> rows = new ArrayList<>();
-            
Env.getCurrentEnv().getTableStreamManager().fillStreamConsumptionValuesMetadataResult(rows);
+            MonitoredReentrantReadWriteLock tableLock = 
Deencapsulation.getField(table, "rwLock");
+            MonitoredReentrantReadWriteLock streamLock = 
Deencapsulation.getField(stream, "rwLock");
+            AtomicBoolean unitSelected = new AtomicBoolean(false);
+            List<TRow> cloudRows = new ArrayList<>();
+            
Env.getCurrentEnv().getTableStreamManager().fillStreamConsumptionValuesMetadataResult(
+                    cloudRows, (dbName, streamName, streamId, unit) -> {
+                        if (unit != null) {
+                            unitSelected.set(true);
+                            Assertions.assertEquals(0, 
streamLock.getReadHoldCount());
+                            Assertions.assertEquals(0, 
tableLock.getReadHoldCount());
+                        }
+                        return streamName.equals("s1") && (unit == null || 
unit.equals("p1"));
+                    });
+            Assertions.assertTrue(unitSelected.get());
+            Assertions.assertEquals(1, cloudRows.size());
+            Assertions.assertEquals("p1", 
cloudRows.get(0).getColumnValue().get(3).getStringVal());
+
+            TFetchSchemaTableDataRequest request = 
newConsumptionRequest("test_cloud_stream_consumption", "s1");
+            TFetchSchemaTableDataResult result = 
MetadataGenerator.getSchemaTableData(request);
+            Assertions.assertEquals(TStatusCode.OK, 
result.getStatus().getStatusCode());
+            List<TRow> rows = new ArrayList<>(result.getDataBatch());
             rows.sort(Comparator.comparing(row -> 
row.getColumnValue().get(3).getStringVal()));
-            Assertions.assertEquals(2, rows.size());
+            Assertions.assertEquals(1, rows.size());
             Assertions.assertEquals("p1", 
rows.get(0).getColumnValue().get(3).getStringVal());
             Assertions.assertEquals("100", 
rows.get(0).getColumnValue().get(4).getStringVal());
             Assertions.assertEquals("30", 
rows.get(0).getColumnValue().get(5).getStringVal());
             Assertions.assertEquals(999, 
rows.get(0).getColumnValue().get(6).getLongVal());
-            Assertions.assertEquals("p2", 
rows.get(1).getColumnValue().get(3).getStringVal());
-            Assertions.assertEquals("N/A", 
rows.get(1).getColumnValue().get(4).getStringVal());
-            Assertions.assertEquals("0", 
rows.get(1).getColumnValue().get(5).getStringVal());
-            Assertions.assertEquals(-1, 
rows.get(1).getColumnValue().get(6).getLongVal());
+
+            result = 
MetadataGenerator.getSchemaTableData(newStreamIdConsumptionRequest(secondStream.getId()));
+            Assertions.assertEquals(TStatusCode.OK, 
result.getStatus().getStatusCode());
+            Assertions.assertEquals(2, result.getDataBatchSize());
+            Assertions.assertTrue(result.getDataBatch().stream()
+                    .allMatch(row -> row.getColumnValue().get(2).getLongVal() 
== secondStream.getId()));
+
+            TSchemaTableRequestParams invalidParams = new 
TSchemaTableRequestParams();
+            invalidParams.setFrontendConjuncts("{");
+            result = MetadataGenerator.getSchemaTableData(new 
TFetchSchemaTableDataRequest()
+                    
.setSchemaTableName(TSchemaTableName.TABLE_STREAM_CONSUMPTION)
+                    .setSchemaTableParams(invalidParams));
+            Assertions.assertEquals(TStatusCode.OK, 
result.getStatus().getStatusCode());
+            Assertions.assertEquals(4, result.getDataBatchSize());
 
             table.writeLock();
             try {
@@ -136,13 +223,39 @@ public class CloudTableStreamConsumptionTest extends 
TestWithFeService {
             } finally {
                 table.writeUnlock();
             }
-            rows.clear();
-            
Env.getCurrentEnv().getTableStreamManager().fillStreamConsumptionValuesMetadataResult(rows);
-            Assertions.assertTrue(rows.isEmpty());
-            Mockito.verify(proxy, 
Mockito.times(1)).getTableStreamOffset(Mockito.any());
+            result = MetadataGenerator.getSchemaTableData(request);
+            Assertions.assertEquals(TStatusCode.OK, 
result.getStatus().getStatusCode());
+            Assertions.assertTrue(result.getDataBatch().isEmpty());
+            Mockito.verify(proxy, 
Mockito.times(4)).getTableStreamOffset(Mockito.any());
         } finally {
             Config.cloud_unique_id = previousCloudUniqueId;
             Config.meta_service_endpoint = previousMetaServiceEndpoint;
         }
     }
+
+    private TFetchSchemaTableDataRequest newConsumptionRequest(String dbName, 
String streamName) {
+        TSchemaTableRequestParams params = new TSchemaTableRequestParams();
+        params.setFrontendConjuncts(GsonUtils.GSON.toJson(List.of(
+                new BinaryPredicate(BinaryPredicate.Operator.EQ,
+                        new SlotRef(null, "DB_NAME"), new 
StringLiteral(dbName)),
+                new BinaryPredicate(BinaryPredicate.Operator.EQ,
+                        new SlotRef(null, "STREAM_NAME"), new 
StringLiteral(streamName)),
+                new BinaryPredicate(BinaryPredicate.Operator.EQ,
+                        new SlotRef(null, "UNIT"), new StringLiteral("p1")),
+                new BinaryPredicate(BinaryPredicate.Operator.GT,
+                        new SlotRef(null, "LAG"), new IntLiteral(100)))));
+        return new TFetchSchemaTableDataRequest()
+                .setSchemaTableName(TSchemaTableName.TABLE_STREAM_CONSUMPTION)
+                .setSchemaTableParams(params);
+    }
+
+    private TFetchSchemaTableDataRequest newStreamIdConsumptionRequest(long 
streamId) {
+        TSchemaTableRequestParams params = new TSchemaTableRequestParams();
+        params.setFrontendConjuncts(GsonUtils.GSON.toJson(List.of(
+                new BinaryPredicate(BinaryPredicate.Operator.EQ,
+                        new SlotRef(null, "STREAM_ID"), new 
IntLiteral(streamId)))));
+        return new TFetchSchemaTableDataRequest()
+                .setSchemaTableName(TSchemaTableName.TABLE_STREAM_CONSUMPTION)
+                .setSchemaTableParams(params);
+    }
 }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/OlapTableStreamConsumptionTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/OlapTableStreamConsumptionTest.java
new file mode 100644
index 00000000000..5b34e8854ba
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/OlapTableStreamConsumptionTest.java
@@ -0,0 +1,69 @@
+// 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.catalog.stream;
+
+import org.apache.doris.catalog.OlapTable;
+import org.apache.doris.catalog.Partition;
+import org.apache.doris.common.jmockit.Deencapsulation;
+import org.apache.doris.thrift.TRow;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+public class OlapTableStreamConsumptionTest {
+
+    @Test
+    public void testUnitSelectorRunsAfterReleasingBaseTableLock() {
+        OlapTable table = Mockito.mock(OlapTable.class);
+        Partition p1 = Mockito.mock(Partition.class);
+        Partition p2 = Mockito.mock(Partition.class);
+        AtomicBoolean tableLocked = new AtomicBoolean(false);
+        Mockito.when(table.readLockIfExist()).thenAnswer(invocation -> {
+            tableLocked.set(true);
+            return true;
+        });
+        Mockito.doAnswer(invocation -> {
+            tableLocked.set(false);
+            return null;
+        }).when(table).readUnlock();
+        Mockito.when(table.getPartitions()).thenReturn(List.of(p1, p2));
+        Mockito.when(p1.getName()).thenReturn("p1");
+        Mockito.when(p2.getName()).thenReturn("p2");
+        Mockito.when(p1.getId()).thenReturn(1L);
+        Mockito.when(p2.getId()).thenReturn(2L);
+
+        OlapTableStream stream = Mockito.spy(new OlapTableStream());
+        Mockito.doReturn(table).when(stream).getBaseTableNullable();
+        Deencapsulation.setField(stream, "partitionOffset", new HashMap<Long, 
Long>());
+        Deencapsulation.setField(stream, "partitionConsumptionTime", new 
HashMap<Long, Long>());
+        List<TRow> rows = new ArrayList<>();
+
+        stream.fillTableStreamConsumptionInfo(rows, unit -> {
+            Assertions.assertFalse(tableLocked.get());
+            return false;
+        });
+
+        Assertions.assertTrue(rows.isEmpty());
+    }
+}
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/RewriteRuleSuiteTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/RewriteRuleSuiteTest.java
index c1024026e02..92319ae4372 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/RewriteRuleSuiteTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/rules/rewrite/RewriteRuleSuiteTest.java
@@ -29,6 +29,7 @@ import org.apache.doris.catalog.MaterializedIndex;
 import org.apache.doris.catalog.OlapTable;
 import org.apache.doris.catalog.Partition;
 import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.SchemaTable;
 import org.apache.doris.catalog.Tablet;
 import org.apache.doris.nereids.CascadesContext;
 import org.apache.doris.nereids.sqltest.SqlTestBase;
@@ -36,13 +37,20 @@ import org.apache.doris.nereids.trees.expressions.EqualTo;
 import org.apache.doris.nereids.trees.expressions.GreaterThanEqual;
 import org.apache.doris.nereids.trees.expressions.InPredicate;
 import org.apache.doris.nereids.trees.expressions.LessThanEqual;
+import org.apache.doris.nereids.trees.expressions.Slot;
+import org.apache.doris.nereids.trees.expressions.functions.scalar.AssertTrue;
+import org.apache.doris.nereids.trees.expressions.literal.BooleanLiteral;
 import org.apache.doris.nereids.trees.expressions.literal.DateLiteral;
 import org.apache.doris.nereids.trees.expressions.literal.Literal;
+import org.apache.doris.nereids.trees.expressions.literal.VarcharLiteral;
 import org.apache.doris.nereids.trees.plans.RelationId;
 import org.apache.doris.nereids.trees.plans.logical.LogicalFilter;
 import org.apache.doris.nereids.trees.plans.logical.LogicalOlapScan;
+import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
+import org.apache.doris.nereids.trees.plans.logical.LogicalSchemaScan;
 import org.apache.doris.nereids.util.MemoTestUtils;
 import org.apache.doris.nereids.util.PlanChecker;
+import org.apache.doris.nereids.util.PlanConstructor;
 import org.apache.doris.planner.PartitionColumnFilter;
 
 import com.google.common.collect.ImmutableList;
@@ -69,6 +77,41 @@ import java.util.Objects;
  */
 public class RewriteRuleSuiteTest extends SqlTestBase {
 
+    @Test
+    void testNonMovableFunctionBlocksSchemaScanPushdown() {
+        LogicalSchemaScan streamScan = new 
LogicalSchemaScan(PlanConstructor.getNextRelationId(),
+                SchemaTable.TABLE_MAP.get("table_stream_consumption"), 
ImmutableList.of("information_schema"));
+        Slot dbName = streamScan.getOutput().stream()
+                .filter(slot -> slot.getName().equalsIgnoreCase("DB_NAME"))
+                .findFirst()
+                .orElseThrow(IllegalStateException::new);
+        LogicalFilter<LogicalSchemaScan> streamFilter = new 
LogicalFilter<>(ImmutableSet.of(
+                new EqualTo(dbName, new VarcharLiteral("__missing__")),
+                new AssertTrue(BooleanLiteral.FALSE, new VarcharLiteral("must 
fail"))), streamScan);
+
+        LogicalPlan rewritten = (LogicalPlan) PlanChecker.from(connectContext, 
streamFilter)
+                .applyTopDown(new PushDownFilterIntoSchemaScan())
+                .getPlan();
+        LogicalSchemaScan rewrittenScan = (LogicalSchemaScan) 
rewritten.child(0);
+        Assertions.assertTrue(rewrittenScan.getFrontendConjuncts().isEmpty());
+
+        LogicalSchemaScan tablesScan = new 
LogicalSchemaScan(PlanConstructor.getNextRelationId(),
+                SchemaTable.TABLE_MAP.get("tables"), 
ImmutableList.of("information_schema"));
+        Slot tableSchema = tablesScan.getOutput().stream()
+                .filter(slot -> 
slot.getName().equalsIgnoreCase("TABLE_SCHEMA"))
+                .findFirst()
+                .orElseThrow(IllegalStateException::new);
+        LogicalFilter<LogicalSchemaScan> tablesFilter = new 
LogicalFilter<>(ImmutableSet.of(
+                new EqualTo(tableSchema, new VarcharLiteral("__missing__")),
+                new AssertTrue(BooleanLiteral.FALSE, new VarcharLiteral("must 
fail"))), tablesScan);
+
+        rewritten = (LogicalPlan) PlanChecker.from(connectContext, 
tablesFilter)
+                .applyTopDown(new PushDownFilterIntoSchemaScan())
+                .getPlan();
+        rewrittenScan = (LogicalSchemaScan) rewritten.child(0);
+        Assertions.assertFalse(rewrittenScan.getSchemaDatabase().isPresent());
+    }
+
     // 
-------------------------------------------------------------------------
     // from PruneOlapScanTabletTest
     // 
-------------------------------------------------------------------------
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/copier/LogicalPlanDeepCopierTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/copier/LogicalPlanDeepCopierTest.java
index b139d05910a..95c3dddeaf7 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/copier/LogicalPlanDeepCopierTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/copier/LogicalPlanDeepCopierTest.java
@@ -17,15 +17,19 @@
 
 package org.apache.doris.nereids.trees.copier;
 
+import org.apache.doris.catalog.SchemaTable;
+import org.apache.doris.nereids.trees.expressions.EqualTo;
 import org.apache.doris.nereids.trees.expressions.Expression;
 import org.apache.doris.nereids.trees.expressions.NamedExpression;
 import org.apache.doris.nereids.trees.expressions.Slot;
 import org.apache.doris.nereids.trees.expressions.SlotReference;
+import org.apache.doris.nereids.trees.expressions.literal.VarcharLiteral;
 import org.apache.doris.nereids.trees.plans.Plan;
 import org.apache.doris.nereids.trees.plans.algebra.Repeat.RepeatType;
 import org.apache.doris.nereids.trees.plans.logical.LogicalAggregate;
 import org.apache.doris.nereids.trees.plans.logical.LogicalOlapScan;
 import org.apache.doris.nereids.trees.plans.logical.LogicalRepeat;
+import org.apache.doris.nereids.trees.plans.logical.LogicalSchemaScan;
 import org.apache.doris.nereids.types.BigIntType;
 import org.apache.doris.nereids.util.PlanConstructor;
 
@@ -91,6 +95,29 @@ public class LogicalPlanDeepCopierTest {
         Assertions.assertEquals(ImmutableList.of(aCopy.getOutput().get(1)), 
aCopy.getOperativeSlots());
     }
 
+    @Test
+    public void testDeepCopySchemaScanCopiesFrontendConjunctSlots() {
+        LogicalSchemaScan scan = new 
LogicalSchemaScan(PlanConstructor.getNextRelationId(),
+                SchemaTable.TABLE_MAP.get("table_stream_consumption"), 
ImmutableList.of("information_schema"));
+        SlotReference dbName = (SlotReference) scan.getOutput().stream()
+                .filter(slot -> slot.getName().equalsIgnoreCase("DB_NAME"))
+                .findFirst()
+                .orElseThrow(IllegalStateException::new);
+        scan = scan.withFrontendConjuncts(Optional.empty(), Optional.empty(), 
Optional.empty(),
+                ImmutableList.of(new EqualTo(dbName, new 
VarcharLiteral("db1"))));
+
+        LogicalSchemaScan copied = (LogicalSchemaScan) scan.accept(
+                LogicalPlanDeepCopier.INSTANCE, new DeepCopierContext());
+        Slot copiedConjunctSlot = 
copied.getFrontendConjuncts().get(0).getInputSlots().iterator().next();
+        Slot copiedDbName = copied.getOutput().stream()
+                .filter(slot -> slot.getName().equalsIgnoreCase("DB_NAME"))
+                .findFirst()
+                .orElseThrow(IllegalStateException::new);
+
+        Assertions.assertNotEquals(dbName.getExprId(), 
copiedConjunctSlot.getExprId());
+        Assertions.assertEquals(copiedDbName.getExprId(), 
copiedConjunctSlot.getExprId());
+    }
+
     @Test
     public void testDeepCopyAggregateWithSourceRepeat() {
         LogicalOlapScan scan = PlanConstructor.newLogicalOlapScan(0, "t", 0);
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/util/FrontendConjunctsUtilsTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/util/FrontendConjunctsUtilsTest.java
index 70bcac72bab..c19f5127542 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/nereids/util/FrontendConjunctsUtilsTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/nereids/util/FrontendConjunctsUtilsTest.java
@@ -17,6 +17,10 @@
 
 package org.apache.doris.nereids.util;
 
+import org.apache.doris.analysis.Expr;
+import org.apache.doris.analysis.IntLiteral;
+import org.apache.doris.analysis.SlotRef;
+import org.apache.doris.catalog.Type;
 import org.apache.doris.nereids.analyzer.UnboundSlot;
 import org.apache.doris.nereids.trees.expressions.And;
 import org.apache.doris.nereids.trees.expressions.EqualTo;
@@ -31,6 +35,8 @@ import org.apache.doris.nereids.trees.expressions.Not;
 import org.apache.doris.nereids.trees.expressions.NullSafeEqual;
 import org.apache.doris.nereids.trees.expressions.Or;
 import org.apache.doris.nereids.trees.expressions.literal.Literal;
+import org.apache.doris.nereids.types.DataType;
+import org.apache.doris.persist.gson.GsonUtils;
 
 import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.Lists;
@@ -111,6 +117,30 @@ public class FrontendConjunctsUtilsTest {
         
Assertions.assertTrue(FrontendConjunctsUtils.isFiltered(Lists.newArrayList(in), 
"c1", "v3"));
     }
 
+    @Test
+    public void testSerializedBigIntInPredicate() throws Exception {
+        List<Expr> options = Lists.newArrayList(
+                new IntLiteral(10000L, Type.BIGINT),
+                new IntLiteral(10001L, Type.BIGINT));
+        Expr in = new org.apache.doris.analysis.InPredicate(new SlotRef(null, 
"STREAM_ID"), options, false);
+        List<Expression> expressions = 
FrontendConjunctsUtils.convertToExpression(
+                GsonUtils.GSON.toJson(Lists.newArrayList(in)));
+
+        Assertions.assertFalse(FrontendConjunctsUtils.isFiltered(expressions, 
"STREAM_ID", 10000L));
+        Assertions.assertTrue(FrontendConjunctsUtils.isFiltered(expressions, 
"STREAM_ID", 10002L));
+    }
+
+    @Test
+    public void testSerializedIntLiteralPreservesType() throws Exception {
+        for (Type type : Lists.newArrayList(Type.TINYINT, Type.SMALLINT, 
Type.INT, Type.BIGINT)) {
+            IntLiteral literal = new IntLiteral(1L, type);
+            Expression expression = FrontendConjunctsUtils.convertToExpression(
+                    GsonUtils.GSON.toJson(Lists.newArrayList(literal))).get(0);
+
+            Assertions.assertEquals(DataType.fromCatalogType(type), 
expression.getDataType(), type.toSql());
+        }
+    }
+
     @Test
     public void testLike() {
         Like like = generateLike("c1", "%value%");
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/tablefunction/StreamConsumptionMetadataGeneratorTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/tablefunction/StreamConsumptionMetadataGeneratorTest.java
new file mode 100644
index 00000000000..2514d0e28d7
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/tablefunction/StreamConsumptionMetadataGeneratorTest.java
@@ -0,0 +1,115 @@
+// 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.tablefunction;
+
+import org.apache.doris.analysis.BinaryPredicate;
+import org.apache.doris.analysis.CompoundPredicate;
+import org.apache.doris.analysis.Expr;
+import org.apache.doris.analysis.SlotRef;
+import org.apache.doris.analysis.StringLiteral;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.stream.TableStreamManager;
+import org.apache.doris.persist.gson.GsonUtils;
+import org.apache.doris.thrift.TFetchSchemaTableDataRequest;
+import org.apache.doris.thrift.TFetchSchemaTableDataResult;
+import org.apache.doris.thrift.TSchemaTableName;
+import org.apache.doris.thrift.TSchemaTableRequestParams;
+import org.apache.doris.thrift.TStatusCode;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.util.List;
+
+public class StreamConsumptionMetadataGeneratorTest {
+
+    @Test
+    public void testNoConjunctsUseUnfilteredScanPath() throws Exception {
+        Env env = Mockito.mock(Env.class);
+        TableStreamManager manager = Mockito.mock(TableStreamManager.class);
+        Mockito.when(env.getTableStreamManager()).thenReturn(manager);
+
+        try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+            mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+
+            TFetchSchemaTableDataRequest request = new 
TFetchSchemaTableDataRequest()
+                    
.setSchemaTableName(TSchemaTableName.TABLE_STREAM_CONSUMPTION)
+                    .setSchemaTableParams(new TSchemaTableRequestParams());
+            TFetchSchemaTableDataResult result = 
MetadataGenerator.getSchemaTableData(request);
+
+            Assertions.assertEquals(TStatusCode.OK, 
result.getStatus().getStatusCode());
+            
Mockito.verify(manager).fillStreamConsumptionValuesMetadataResult(Mockito.anyList());
+            Mockito.verify(manager, 
Mockito.never()).fillStreamConsumptionValuesMetadataResult(
+                    Mockito.anyList(), 
Mockito.any(TableStreamManager.StreamConsumptionSelector.class));
+        }
+    }
+
+    @Test
+    public void testStreamConjunctIsNotReevaluatedForUnit() throws Exception {
+        TableStreamManager.StreamConsumptionSelector selector = 
captureSelector(List.of(
+                new BinaryPredicate(BinaryPredicate.Operator.EQ,
+                        new SlotRef(null, "DB_NAME"), new 
StringLiteral("db1"))));
+
+        Assertions.assertFalse(selector.hasUnitFilter());
+        Assertions.assertFalse(selector.test("other_db", "s1", 1, null));
+        Assertions.assertTrue(selector.test("other_db", "s1", 1, "p1"));
+    }
+
+    @Test
+    public void testCrossLevelOrIsEvaluatedWithUnit() throws Exception {
+        Expr dbPredicate = new BinaryPredicate(BinaryPredicate.Operator.EQ,
+                new SlotRef(null, "DB_NAME"), new StringLiteral("db1"));
+        Expr unitPredicate = new BinaryPredicate(BinaryPredicate.Operator.EQ,
+                new SlotRef(null, "UNIT"), new StringLiteral("p1"));
+        TableStreamManager.StreamConsumptionSelector selector = 
captureSelector(List.of(
+                new CompoundPredicate(CompoundPredicate.Operator.OR, 
dbPredicate, unitPredicate)));
+
+        Assertions.assertTrue(selector.hasUnitFilter());
+        Assertions.assertTrue(selector.test("other_db", "s1", 1, null));
+        Assertions.assertTrue(selector.test("db1", "s1", 1, "p2"));
+        Assertions.assertTrue(selector.test("other_db", "s1", 1, "p1"));
+        Assertions.assertFalse(selector.test("other_db", "s1", 1, "p2"));
+    }
+
+    private TableStreamManager.StreamConsumptionSelector 
captureSelector(List<Expr> conjuncts) throws Exception {
+        Env env = Mockito.mock(Env.class);
+        TableStreamManager manager = Mockito.mock(TableStreamManager.class);
+        Mockito.when(env.getTableStreamManager()).thenReturn(manager);
+
+        try (MockedStatic<Env> mockedEnv = Mockito.mockStatic(Env.class)) {
+            mockedEnv.when(Env::getCurrentEnv).thenReturn(env);
+
+            TSchemaTableRequestParams params = new TSchemaTableRequestParams();
+            params.setFrontendConjuncts(GsonUtils.GSON.toJson(conjuncts));
+            TFetchSchemaTableDataResult result = 
MetadataGenerator.getSchemaTableData(
+                    new TFetchSchemaTableDataRequest()
+                            
.setSchemaTableName(TSchemaTableName.TABLE_STREAM_CONSUMPTION)
+                            .setSchemaTableParams(params));
+
+            Assertions.assertEquals(TStatusCode.OK, 
result.getStatus().getStatusCode());
+            ArgumentCaptor<TableStreamManager.StreamConsumptionSelector> 
selectorCaptor =
+                    
ArgumentCaptor.forClass(TableStreamManager.StreamConsumptionSelector.class);
+            Mockito.verify(manager).fillStreamConsumptionValuesMetadataResult(
+                    Mockito.anyList(), selectorCaptor.capture());
+            return selectorCaptor.getValue();
+        }
+    }
+}
diff --git 
a/regression-test/suites/query_p0/schema_table/test_stream_consumption_schema.groovy
 
b/regression-test/suites/query_p0/schema_table/test_stream_consumption_schema.groovy
index f06b488615f..ab5e8cb0fbe 100644
--- 
a/regression-test/suites/query_p0/schema_table/test_stream_consumption_schema.groovy
+++ 
b/regression-test/suites/query_p0/schema_table/test_stream_consumption_schema.groovy
@@ -133,5 +133,40 @@ suite("test_stream_consumption_schema") {
     """
 
     qt_sql "select DB_NAME,STREAM_NAME,UNIT,LAG,LAST_CONSUMPTION_TIME from 
information_schema.table_stream_consumption where DB_NAME = 
'test_stream_consumption_db' order by STREAM_NAME, UNIT;"
+
+    def explain = sql """
+        EXPLAIN SELECT * FROM information_schema.table_stream_consumption
+        WHERE DB_NAME = 'test_stream_consumption_db'
+    """
+    assertTrue(explain.toString().contains("FRONTEND PREDICATES"))
+
+    explain = sql """
+        EXPLAIN SELECT * FROM information_schema.table_stream_consumption
+        WHERE STREAM_ID = 1
+    """
+    assertTrue(explain.toString().contains("FRONTEND PREDICATES"))
+
+    explain = sql """
+        EXPLAIN SELECT * FROM information_schema.table_stream_consumption
+        WHERE UNIT = 'p1'
+    """
+    assertTrue(explain.toString().contains("FRONTEND PREDICATES"))
+
+    explain = sql """
+        EXPLAIN SELECT * FROM information_schema.table_stream_consumption
+        WHERE DB_NAME = '__missing__'
+          AND assert_true(false, 'frontend predicate pushdown must preserve 
assert_true')
+    """
+    assertFalse(explain.toString().contains("FRONTEND PREDICATES"))
+
+    test {
+        sql """
+            SELECT COUNT(*) FROM information_schema.table_stream_consumption
+            WHERE DB_NAME = '__missing__'
+              AND assert_true(false, 'frontend predicate pushdown must 
preserve assert_true')
+        """
+        exception "frontend predicate pushdown must preserve assert_true"
+    }
+
     sql "DROP DATABASE IF EXISTS test_stream_consumption_db"
-}
\ No newline at end of file
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to