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, ¶m, &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]