linrrzqqq commented on code in PR #68227:
URL: https://github.com/apache/doris/pull/68227#discussion_r4068948823
##########
fe/fe-core/src/main/java/org/apache/doris/tablefunction/MetadataGenerator.java:
##########
@@ -2329,8 +2337,68 @@ private static TFetchSchemaTableDataResult
streamMetadataResult(TSchemaTableRequ
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);
Review Comment:
irrelavant
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]