From 28af68e4896f385ca8f04911bf4f0e865b18a897 Mon Sep 17 00:00:00 2001 From: seawinde Date: Thu, 30 Jul 2026 17:25:35 +0800 Subject: [PATCH 1/8] [fix](stream) Mark streams stale when base tables are dropped Issue Number: N/A Related PR: https://github.com/yujun777/doris/pull/29 Problem Summary: Dropping a stream base table left the stream enabled and non-stale because its runtime cache continued returning the dropped table. Derive stream availability from the persisted base table identity, preserve diagnostic qualifiers when the table is unavailable, and allow recovery only when the original table ID returns. Add unit and regression coverage for drop, recover, and force-drop behavior. Streams whose base tables are unavailable are now reported disabled and stale. - Test: Unit Test - ./run-fe-ut.sh --run org.apache.doris.catalog.DropTableStreamTest (4 tests) - env DISABLE_BUILD_UI=ON ./build.sh --fe - Behavior changed: Yes. Streams become disabled and stale while their base tables are unavailable. - Does this need documentation: No --- .../doris/catalog/stream/BaseTableStream.java | 11 ++- .../stream/TableStreamBaseTableInfo.java | 2 +- .../catalog/stream/TableStreamManager.java | 30 ++---- .../doris/catalog/DropTableStreamTest.java | 96 +++++++++++++++++++ .../test_ivm_drop_base_table_stream_state.out | 19 ++++ ...st_ivm_drop_base_table_stream_state.groovy | 96 +++++++++++++++++++ 6 files changed, 230 insertions(+), 24 deletions(-) create mode 100644 regression-test/data/mtmv_p0/ivm/test_ivm_drop_base_table_stream_state.out create mode 100644 regression-test/suites/mtmv_p0/ivm/test_ivm_drop_base_table_stream_state.groovy 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 951df48b20b065..a114d37cf69a13 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 @@ -37,6 +37,8 @@ import java.util.Map; public abstract class BaseTableStream extends Table { + private static final String BASE_TABLE_NOT_FOUND_STALE_REASON = "Base table does not exist"; + public enum StreamScanType { APPEND_ONLY, MIN_DELTA, @@ -113,6 +115,9 @@ public BaseTableStream(String streamName, List fullSchema, TableIf baseT } public TableIf getBaseTableNullable() { + if (baseTable instanceof Table && ((Table) baseTable).isDropped) { + baseTable = null; + } if (baseTable == null) { baseTable = baseTableInfo.getTableNullable(); } @@ -139,7 +144,7 @@ public StreamScanType getStreamScanType() { } public boolean isDisabled() { - return disabled; + return disabled || getBaseTableNullable() == null; } public void setDisabled(boolean disabled) { @@ -147,7 +152,7 @@ public void setDisabled(boolean disabled) { } public boolean isStale() { - return stale; + return stale || getBaseTableNullable() == null; } public void setStale(boolean stale) { @@ -155,7 +160,7 @@ public void setStale(boolean stale) { } public String getStaleReason() { - return staleReason; + return getBaseTableNullable() == null ? BASE_TABLE_NOT_FOUND_STALE_REASON : staleReason; } public void setStaleReason(String staleReason) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/TableStreamBaseTableInfo.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/TableStreamBaseTableInfo.java index 3571b7e3592edd..c6181663832f78 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/TableStreamBaseTableInfo.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/stream/TableStreamBaseTableInfo.java @@ -111,7 +111,7 @@ public TableIf getTableNullable() { } } } - LOG.warn("invalid base table: {}", this); + LOG.debug("invalid base table: {}", this); return null; } 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 052cf16dcad995..635547a77846ac 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 @@ -349,27 +349,17 @@ public void fillTableStreamValuesMetadataResult(List dataBatch) { trow.addToColumnValue(new TCell().setStringVal(stream.getScanTypeString())); // STREAM_COMMENT trow.addToColumnValue(new TCell().setStringVal(stream.getComment())); + List baseTableQualifiers = stream.getBaseTableFullQualifiers(); + // BASE_TABLE_NAME + trow.addToColumnValue(new TCell().setStringVal(baseTableQualifiers.get(2))); + // BASE_TABLE_DB + trow.addToColumnValue(new TCell().setStringVal(baseTableQualifiers.get(1))); + // BASE_TABLE_CTL + trow.addToColumnValue(new TCell().setStringVal(baseTableQualifiers.get(0))); + // BASE_TABLE_TYPE TableIf baseTable = stream.getBaseTableNullable(); - if (baseTable == null) { - // BASE_TABLE_NAME - trow.addToColumnValue(new TCell().setStringVal("N/A")); - // BASE_TABLE_DB - trow.addToColumnValue(new TCell().setStringVal("N/A")); - // BASE_TABLE_CTL - trow.addToColumnValue(new TCell().setStringVal("N/A")); - // BASE_TABLE_TYPE - trow.addToColumnValue(new TCell().setStringVal("N/A")); - } else { - List baseTableQualifiers = baseTable.getFullQualifiers(); - // BASE_TABLE_NAME - trow.addToColumnValue(new TCell().setStringVal(baseTableQualifiers.get(2))); - // BASE_TABLE_DB - trow.addToColumnValue(new TCell().setStringVal(baseTableQualifiers.get(1))); - // BASE_TABLE_CTL - trow.addToColumnValue(new TCell().setStringVal(baseTableQualifiers.get(0))); - // BASE_TABLE_TYPE - trow.addToColumnValue(new TCell().setStringVal(baseTable.getType().name())); - } + trow.addToColumnValue(new TCell().setStringVal( + baseTable == null ? "N/A" : baseTable.getType().name())); // ENABLED trow.addToColumnValue(new TCell().setBoolVal(!stream.isDisabled())); // IS_STALE diff --git a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java index 83a43edfb00c03..76b7511c0f6899 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java @@ -17,19 +17,26 @@ package org.apache.doris.catalog; +import org.apache.doris.catalog.stream.OlapTableStream; import org.apache.doris.common.Config; import org.apache.doris.common.DdlException; import org.apache.doris.common.ExceptionChecker; import org.apache.doris.common.FeConstants; +import org.apache.doris.nereids.exceptions.AnalysisException; import org.apache.doris.nereids.parser.NereidsParser; import org.apache.doris.nereids.trees.plans.commands.DropStreamCommand; import org.apache.doris.nereids.trees.plans.logical.LogicalPlan; +import org.apache.doris.persist.gson.GsonUtils; import org.apache.doris.qe.StmtExecutor; +import org.apache.doris.thrift.TRow; import org.apache.doris.utframe.TestWithFeService; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.util.ArrayList; +import java.util.List; + public class DropTableStreamTest extends TestWithFeService { @Override @@ -67,6 +74,15 @@ private void dropStream(String sql) throws Exception { } } + private void createBaseTableAndStream(String tableName, String streamName) throws Exception { + createTable("create table test_stream." + tableName + " (k1 int, k2 int) " + + "unique key(k1) distributed by hash(k1) buckets 1 " + + "properties('replication_num' = '1', 'binlog.enable' = 'true', 'binlog.format' = 'ROW', " + + "'binlog.need_historical_value' = 'true')"); + createTable("create stream test_stream." + streamName + " on table test_stream." + tableName + + " properties('show_initial_rows' = 'true')"); + } + @Test public void testNormalDropStream() throws Exception { // test drop @@ -104,6 +120,86 @@ public void testCloudDropRequiresForce() { } } + @Test + public void testStreamStateFollowsRecoverableBaseTableDrop() throws Exception { + createBaseTableAndStream("tbl_recover", "s_recover"); + Database db = Env.getCurrentInternalCatalog().getDbOrMetaException("test_stream"); + OlapTable baseTable = (OlapTable) db.getTableOrMetaException("tbl_recover"); + OlapTableStream stream = (OlapTableStream) db.getTableOrMetaException("s_recover"); + + dropTableWithSql("drop table test_stream.tbl_recover"); + + Assertions.assertTrue(baseTable.isDropped); + Assertions.assertNull(stream.getBaseTableNullable()); + Assertions.assertTrue(stream.isDisabled()); + Assertions.assertTrue(stream.isStale()); + Assertions.assertEquals("Base table does not exist", stream.getStaleReason()); + Assertions.assertTrue(Env.getCurrentEnv().getTableStreamManager().getTableStreamIds(db) + .contains(stream.getId())); + + List rows = new ArrayList<>(); + Env.getCurrentEnv().getTableStreamManager().fillTableStreamValuesMetadataResult(rows); + TRow streamRow = rows.stream() + .filter(row -> "s_recover".equals(row.getColumnValue().get(1).getStringVal())) + .findFirst() + .orElseThrow(AssertionError::new); + List baseTableQualifiers = stream.getBaseTableFullQualifiers(); + Assertions.assertEquals(baseTableQualifiers.get(2), streamRow.getColumnValue().get(6).getStringVal()); + Assertions.assertEquals(baseTableQualifiers.get(1), streamRow.getColumnValue().get(7).getStringVal()); + Assertions.assertEquals(baseTableQualifiers.get(0), streamRow.getColumnValue().get(8).getStringVal()); + Assertions.assertEquals("N/A", streamRow.getColumnValue().get(9).getStringVal()); + Assertions.assertFalse(streamRow.getColumnValue().get(10).isBoolVal()); + Assertions.assertTrue(streamRow.getColumnValue().get(11).isBoolVal()); + Assertions.assertEquals("Base table does not exist", + streamRow.getColumnValue().get(12).getStringVal()); + + ExceptionChecker.expectThrowsWithMsg(AnalysisException.class, "Unknown base table 'tbl_recover'", + stream::getBaseTableOrNereidsAnalysisException); + ExceptionChecker.expectThrowsWithMsg(IllegalStateException.class, "Table [tbl_recover] does not exist", + () -> executeSql("select * from test_stream.s_recover")); + + recoverTable("recover table test_stream.tbl_recover"); + + Assertions.assertSame(baseTable, stream.getBaseTableNullable()); + Assertions.assertEquals(baseTable.getId(), stream.getBaseTableNullable().getId()); + Assertions.assertFalse(stream.isDisabled()); + Assertions.assertFalse(stream.isStale()); + Assertions.assertEquals("N/A", stream.getStaleReason()); + } + + @Test + public void testForceDropAndSameNameTableDoNotRestoreStream() throws Exception { + createBaseTableAndStream("tbl_force", "s_force"); + Database db = Env.getCurrentInternalCatalog().getDbOrMetaException("test_stream"); + long oldBaseTableId = db.getTableOrMetaException("tbl_force").getId(); + OlapTableStream stream = (OlapTableStream) db.getTableOrMetaException("s_force"); + OlapTableStream deserializedStream = GsonUtils.GSON.fromJson( + GsonUtils.GSON.toJson(stream), OlapTableStream.class); + + dropTableWithSql("drop table test_stream.tbl_force force"); + + Assertions.assertNull(stream.getBaseTableNullable()); + Assertions.assertNull(deserializedStream.getBaseTableNullable()); + Assertions.assertTrue(deserializedStream.isDisabled()); + Assertions.assertTrue(deserializedStream.isStale()); + ExceptionChecker.expectThrowsWithMsg(DdlException.class, + "Unknown table 'tbl_force' or table id '-1' in test_stream", + () -> recoverTable("recover table test_stream.tbl_force")); + + createTable("create table test_stream.tbl_force (k1 int, k2 int) " + + "unique key(k1) distributed by hash(k1) buckets 1 " + + "properties('replication_num' = '1', 'binlog.enable' = 'true', 'binlog.format' = 'ROW', " + + "'binlog.need_historical_value' = 'true')"); + + Assertions.assertNotEquals(oldBaseTableId, db.getTableOrMetaException("tbl_force").getId()); + Assertions.assertNull(stream.getBaseTableNullable()); + Assertions.assertTrue(stream.isDisabled()); + Assertions.assertTrue(stream.isStale()); + Assertions.assertEquals("Base table does not exist", stream.getStaleReason()); + Assertions.assertTrue(Env.getCurrentEnv().getTableStreamManager().getTableStreamIds(db) + .contains(stream.getId())); + } + @Override protected void runAfterAll() throws Exception { dropDatabase("test_stream"); diff --git a/regression-test/data/mtmv_p0/ivm/test_ivm_drop_base_table_stream_state.out b/regression-test/data/mtmv_p0/ivm/test_ivm_drop_base_table_stream_state.out new file mode 100644 index 00000000000000..c2023463324825 --- /dev/null +++ b/regression-test/data/mtmv_p0/ivm/test_ivm_drop_base_table_stream_state.out @@ -0,0 +1,19 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !initial_mv_rows -- +1 10 +2 20 + +-- !stream_before_drop -- +internal regression_test_mtmv_p0_ivm ivm_drop_base_stream_state_base OLAP true false N/A + +-- !stream_after_drop -- +internal regression_test_mtmv_p0_ivm ivm_drop_base_stream_state_base N/A false true Base table does not exist + +-- !stream_after_recover -- +internal regression_test_mtmv_p0_ivm ivm_drop_base_stream_state_base OLAP true false N/A + +-- !mv_rows_after_recover -- +1 10 +2 20 +3 30 + diff --git a/regression-test/suites/mtmv_p0/ivm/test_ivm_drop_base_table_stream_state.groovy b/regression-test/suites/mtmv_p0/ivm/test_ivm_drop_base_table_stream_state.groovy new file mode 100644 index 00000000000000..0d875f04abb97f --- /dev/null +++ b/regression-test/suites/mtmv_p0/ivm/test_ivm_drop_base_table_stream_state.groovy @@ -0,0 +1,96 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +suite("test_ivm_drop_base_table_stream_state") { + if (isCloudMode()) { + return + } + + sql "DROP STREAM IF EXISTS ivm_drop_base_stream_state_stream" + sql "DROP TABLE IF EXISTS ivm_drop_base_stream_state_base" + + sql """ + CREATE TABLE ivm_drop_base_stream_state_base ( + id BIGINT NOT NULL, + value BIGINT + ) UNIQUE KEY(id) + DISTRIBUTED BY HASH(id) BUCKETS 1 + PROPERTIES ( + "replication_num" = "1", + "enable_unique_key_merge_on_write" = "true", + "binlog.enable" = "true", + "binlog.format" = "ROW", + "binlog.need_historical_value" = "true" + ) + """ + sql "INSERT INTO ivm_drop_base_stream_state_base VALUES (1, 10), (2, 20)" + + sql """ + CREATE STREAM ivm_drop_base_stream_state_stream + ON TABLE ivm_drop_base_stream_state_base + PROPERTIES ("show_initial_rows" = "true") + """ + + order_qt_initial_mv_rows """ + SELECT id, value + FROM ivm_drop_base_stream_state_stream + ORDER BY id + """ + + order_qt_stream_before_drop """ + SELECT BASE_TABLE_CTL, BASE_TABLE_DB, BASE_TABLE_NAME, BASE_TABLE_TYPE, + ENABLED, IS_STALE, STALE_REASON + FROM information_schema.table_streams + WHERE DB_NAME = '${context.dbName}' + AND BASE_TABLE_NAME = 'ivm_drop_base_stream_state_base' + AND STREAM_NAME = 'ivm_drop_base_stream_state_stream' + ORDER BY STREAM_NAME + """ + + sql "DROP TABLE ivm_drop_base_stream_state_base" + + order_qt_stream_after_drop """ + SELECT BASE_TABLE_CTL, BASE_TABLE_DB, BASE_TABLE_NAME, BASE_TABLE_TYPE, + ENABLED, IS_STALE, STALE_REASON + FROM information_schema.table_streams + WHERE DB_NAME = '${context.dbName}' + AND BASE_TABLE_NAME = 'ivm_drop_base_stream_state_base' + AND STREAM_NAME = 'ivm_drop_base_stream_state_stream' + ORDER BY STREAM_NAME + """ + + sql "RECOVER TABLE ivm_drop_base_stream_state_base" + + order_qt_stream_after_recover """ + SELECT BASE_TABLE_CTL, BASE_TABLE_DB, BASE_TABLE_NAME, BASE_TABLE_TYPE, + ENABLED, IS_STALE, STALE_REASON + FROM information_schema.table_streams + WHERE DB_NAME = '${context.dbName}' + AND BASE_TABLE_NAME = 'ivm_drop_base_stream_state_base' + AND STREAM_NAME = 'ivm_drop_base_stream_state_stream' + ORDER BY STREAM_NAME + """ + + sql "INSERT INTO ivm_drop_base_stream_state_base VALUES (3, 30)" + sql "sync" + + order_qt_mv_rows_after_recover """ + SELECT id, value + FROM ivm_drop_base_stream_state_stream + ORDER BY id + """ +} From 6282603aac09e3873343cba14dba4fdcd1799a61 Mon Sep 17 00:00:00 2001 From: seawinde Date: Mon, 3 Aug 2026 15:20:58 +0800 Subject: [PATCH 2/8] [fix](stream) Keep base table metadata consistent ### What problem does this PR solve? Issue Number: N/A Related PR: #61382 Problem Summary: Stream metadata could dereference a concurrently cleared base table, mix availability states within one table_streams row, and keep creation-time base table names after same-ID rename or recovery. Resolve the base table once per metadata row, use live qualifiers when available, and recover the latest dropped-table identity from the recycle bin for SHOW CREATE without adding persistent state. ### Release note Stream metadata and SHOW CREATE STREAM now preserve consistent, current base table identity across rename, drop, and recovery. ### Check List (For Author) - Test: Unit Test - ./run-fe-ut.sh --run org.apache.doris.catalog.DropTableStreamTest (5 tests passed before the final SHOW CREATE assertion; the final source and test compilation passed, but rerunning was blocked by an external IntelliJ JPS error-stub class in target/test-classes) - env DISABLE_BUILD_UI=ON ./build.sh --fe --clean (passed with 0 checkstyle violations before the external JPS process began overwriting target classes) - Behavior changed: Yes. Stream metadata uses one base-table availability snapshot and preserves the latest same-ID base table qualifiers. - Does this need documentation: No --- .../doris/catalog/CatalogRecycleBin.java | 10 +++ .../java/org/apache/doris/catalog/Env.java | 7 +- .../doris/catalog/stream/BaseTableStream.java | 41 +++++++++--- .../catalog/stream/TableStreamManager.java | 10 +-- .../doris/catalog/DropTableStreamTest.java | 65 +++++++++++++++++-- 5 files changed, 108 insertions(+), 25 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/CatalogRecycleBin.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/CatalogRecycleBin.java index ca6f332a99531f..a185071fce2e7b 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/CatalogRecycleBin.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/CatalogRecycleBin.java @@ -275,6 +275,16 @@ public boolean isRecycleTable(long dbId, long tableId) { return isRecycleDatabase(dbId) || idToTable.containsKey(tableId); } + public Table getRecycledTableNullable(long dbId, long tableId) { + readLock(); + try { + RecycleTableInfo tableInfo = idToTable.get(tableId); + return tableInfo != null && tableInfo.getDbId() == dbId ? tableInfo.getTable() : null; + } finally { + readUnlock(); + } + } + public boolean isRecyclePartition(long dbId, long tableId, long partitionId) { return isRecycleTable(dbId, tableId) || idToPartition.containsKey(partitionId); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java index 50c1bfc217192d..41f0cf395a53ba 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java @@ -4590,12 +4590,7 @@ public static void getDdlStmt(Command command, String dbName, TableIf table, Lis sb.append("CREATE STREAM "); sb.append('`').append(table.getName()).append('`').append('\n'); - TableIf baseTable = stream.getBaseTableNullable(); - if (baseTable != null) { - sb.append("ON TABLE ").append(baseTable.getNameWithFullQualifiers()); - } else { - sb.append("ON TABLE ").append("UNKNOWN"); - } + sb.append("ON TABLE ").append(String.join(".", stream.getBaseTableFullQualifiers())); // (COMMENT STRING_LITERAL)? addTableComment(table, sb); // properties=propertyClause? 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 a114d37cf69a13..db395f1cf22794 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 @@ -18,6 +18,7 @@ package org.apache.doris.catalog.stream; import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.Env; import org.apache.doris.catalog.Table; import org.apache.doris.catalog.TableIf; import org.apache.doris.common.UserException; @@ -115,13 +116,16 @@ public BaseTableStream(String streamName, List fullSchema, TableIf baseT } public TableIf getBaseTableNullable() { - if (baseTable instanceof Table && ((Table) baseTable).isDropped) { + TableIf cachedBaseTable = baseTable; + if (cachedBaseTable instanceof Table && ((Table) cachedBaseTable).isDropped) { baseTable = null; + return null; } - if (baseTable == null) { - baseTable = baseTableInfo.getTableNullable(); + if (cachedBaseTable == null) { + cachedBaseTable = baseTableInfo.getTableNullable(); + baseTable = cachedBaseTable; } - return baseTable; + return cachedBaseTable; } public void setProperties(Map properties) throws org.apache.doris.common.AnalysisException { @@ -144,7 +148,11 @@ public StreamScanType getStreamScanType() { } public boolean isDisabled() { - return disabled || getBaseTableNullable() == null; + return isDisabled(getBaseTableNullable()); + } + + boolean isDisabled(TableIf availableBaseTable) { + return disabled || availableBaseTable == null; } public void setDisabled(boolean disabled) { @@ -152,7 +160,11 @@ public void setDisabled(boolean disabled) { } public boolean isStale() { - return stale || getBaseTableNullable() == null; + return isStale(getBaseTableNullable()); + } + + boolean isStale(TableIf availableBaseTable) { + return stale || availableBaseTable == null; } public void setStale(boolean stale) { @@ -160,7 +172,11 @@ public void setStale(boolean stale) { } public String getStaleReason() { - return getBaseTableNullable() == null ? BASE_TABLE_NOT_FOUND_STALE_REASON : staleReason; + return getStaleReason(getBaseTableNullable()); + } + + String getStaleReason(TableIf availableBaseTable) { + return availableBaseTable == null ? BASE_TABLE_NOT_FOUND_STALE_REASON : staleReason; } public void setStaleReason(String staleReason) { @@ -203,7 +219,16 @@ public TableIf getBaseTableOrNereidsAnalysisException() throws AnalysisException } public List getBaseTableFullQualifiers() { - return baseTableInfo.getFullQualifiers(); + return getBaseTableFullQualifiers(getBaseTableNullable()); + } + + List getBaseTableFullQualifiers(TableIf availableBaseTable) { + TableIf displayBaseTable = availableBaseTable; + if (displayBaseTable == null && baseTableInfo.isInternalTable()) { + displayBaseTable = Env.getCurrentRecycleBin().getRecycledTableNullable( + baseTableInfo.getDbId(), baseTableInfo.getTableId()); + } + return displayBaseTable == null ? baseTableInfo.getFullQualifiers() : displayBaseTable.getFullQualifiers(); } public TableStreamBaseTableInfo getBaseTableInfo() { 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 635547a77846ac..5ebb4ecc57153a 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 @@ -349,7 +349,8 @@ public void fillTableStreamValuesMetadataResult(List dataBatch) { trow.addToColumnValue(new TCell().setStringVal(stream.getScanTypeString())); // STREAM_COMMENT trow.addToColumnValue(new TCell().setStringVal(stream.getComment())); - List baseTableQualifiers = stream.getBaseTableFullQualifiers(); + TableIf baseTable = stream.getBaseTableNullable(); + List baseTableQualifiers = stream.getBaseTableFullQualifiers(baseTable); // BASE_TABLE_NAME trow.addToColumnValue(new TCell().setStringVal(baseTableQualifiers.get(2))); // BASE_TABLE_DB @@ -357,15 +358,14 @@ public void fillTableStreamValuesMetadataResult(List dataBatch) { // BASE_TABLE_CTL trow.addToColumnValue(new TCell().setStringVal(baseTableQualifiers.get(0))); // BASE_TABLE_TYPE - TableIf baseTable = stream.getBaseTableNullable(); trow.addToColumnValue(new TCell().setStringVal( baseTable == null ? "N/A" : baseTable.getType().name())); // ENABLED - trow.addToColumnValue(new TCell().setBoolVal(!stream.isDisabled())); + trow.addToColumnValue(new TCell().setBoolVal(!stream.isDisabled(baseTable))); // IS_STALE - trow.addToColumnValue(new TCell().setBoolVal(stream.isStale())); + trow.addToColumnValue(new TCell().setBoolVal(stream.isStale(baseTable))); // STALE_REASON - trow.addToColumnValue(new TCell().setStringVal(stream.getStaleReason())); + trow.addToColumnValue(new TCell().setStringVal(stream.getStaleReason(baseTable))); dataBatch.add(trow); } finally { stream.readUnlock(); diff --git a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java index 76b7511c0f6899..15474599212f28 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java @@ -83,6 +83,21 @@ private void createBaseTableAndStream(String tableName, String streamName) throw + " properties('show_initial_rows' = 'true')"); } + private TRow getStreamMetadataRow(String streamName) { + List rows = new ArrayList<>(); + Env.getCurrentEnv().getTableStreamManager().fillTableStreamValuesMetadataResult(rows); + return rows.stream() + .filter(row -> streamName.equals(row.getColumnValue().get(1).getStringVal())) + .findFirst() + .orElseThrow(AssertionError::new); + } + + private String getStreamDdl(OlapTableStream stream) { + List createStreamStmt = new ArrayList<>(); + Env.getDdlStmt(stream, createStreamStmt, null, null, false, true, -1L); + return createStreamStmt.get(0); + } + @Test public void testNormalDropStream() throws Exception { // test drop @@ -134,15 +149,11 @@ public void testStreamStateFollowsRecoverableBaseTableDrop() throws Exception { Assertions.assertTrue(stream.isDisabled()); Assertions.assertTrue(stream.isStale()); Assertions.assertEquals("Base table does not exist", stream.getStaleReason()); + Assertions.assertTrue(getStreamDdl(stream).contains("ON TABLE internal.test_stream.tbl_recover")); Assertions.assertTrue(Env.getCurrentEnv().getTableStreamManager().getTableStreamIds(db) .contains(stream.getId())); - List rows = new ArrayList<>(); - Env.getCurrentEnv().getTableStreamManager().fillTableStreamValuesMetadataResult(rows); - TRow streamRow = rows.stream() - .filter(row -> "s_recover".equals(row.getColumnValue().get(1).getStringVal())) - .findFirst() - .orElseThrow(AssertionError::new); + TRow streamRow = getStreamMetadataRow("s_recover"); List baseTableQualifiers = stream.getBaseTableFullQualifiers(); Assertions.assertEquals(baseTableQualifiers.get(2), streamRow.getColumnValue().get(6).getStringVal()); Assertions.assertEquals(baseTableQualifiers.get(1), streamRow.getColumnValue().get(7).getStringVal()); @@ -167,6 +178,48 @@ public void testStreamStateFollowsRecoverableBaseTableDrop() throws Exception { Assertions.assertEquals("N/A", stream.getStaleReason()); } + @Test + public void testBaseTableQualifiersFollowRenameAndRecovery() throws Exception { + createBaseTableAndStream("tbl_rename", "s_rename"); + Database db = Env.getCurrentInternalCatalog().getDbOrMetaException("test_stream"); + OlapTable baseTable = (OlapTable) db.getTableOrMetaException("tbl_rename"); + OlapTableStream stream = (OlapTableStream) db.getTableOrMetaException("s_rename"); + + executeSql("alter table test_stream.tbl_rename rename tbl_renamed"); + + TRow streamRow = getStreamMetadataRow("s_rename"); + Assertions.assertEquals("tbl_renamed", streamRow.getColumnValue().get(6).getStringVal()); + Assertions.assertEquals("OLAP", streamRow.getColumnValue().get(9).getStringVal()); + Assertions.assertTrue(streamRow.getColumnValue().get(10).isBoolVal()); + Assertions.assertTrue(getStreamDdl(stream).contains("ON TABLE internal.test_stream.tbl_renamed")); + + dropTableWithSql("drop table test_stream.tbl_renamed"); + + streamRow = getStreamMetadataRow("s_rename"); + Assertions.assertEquals("tbl_renamed", streamRow.getColumnValue().get(6).getStringVal()); + Assertions.assertEquals("N/A", streamRow.getColumnValue().get(9).getStringVal()); + Assertions.assertFalse(streamRow.getColumnValue().get(10).isBoolVal()); + Assertions.assertTrue(streamRow.getColumnValue().get(11).isBoolVal()); + Assertions.assertTrue(getStreamDdl(stream).contains("ON TABLE internal.test_stream.tbl_renamed")); + + OlapTableStream deserializedStream = GsonUtils.GSON.fromJson( + GsonUtils.GSON.toJson(stream), OlapTableStream.class); + Assertions.assertEquals("tbl_renamed", deserializedStream.getBaseTableFullQualifiers().get(2)); + Assertions.assertTrue(getStreamDdl(deserializedStream).contains("ON TABLE internal.test_stream.tbl_renamed")); + + recoverTable("recover table test_stream.tbl_renamed as tbl_recovered"); + + Assertions.assertSame(baseTable, stream.getBaseTableNullable()); + Assertions.assertEquals("tbl_recovered", stream.getBaseTableFullQualifiers().get(2)); + Assertions.assertTrue(getStreamDdl(stream).contains("ON TABLE internal.test_stream.tbl_recovered")); + + dropTableWithSql("drop table test_stream.tbl_recovered"); + + deserializedStream = GsonUtils.GSON.fromJson(GsonUtils.GSON.toJson(stream), OlapTableStream.class); + Assertions.assertEquals("tbl_recovered", deserializedStream.getBaseTableFullQualifiers().get(2)); + Assertions.assertTrue(getStreamDdl(deserializedStream).contains("ON TABLE internal.test_stream.tbl_recovered")); + } + @Test public void testForceDropAndSameNameTableDoNotRestoreStream() throws Exception { createBaseTableAndStream("tbl_force", "s_force"); From cc88b0656fc1d376c0a7baf1cab230b2fc1154ca Mon Sep 17 00:00:00 2001 From: seawinde Date: Tue, 4 Aug 2026 15:02:36 +0800 Subject: [PATCH 3/8] [fix](stream) Resolve recycled base table qualifiers by ID ### What problem does this PR solve? Issue Number: N/A Related PR: #66287 Problem Summary: A recycled base table retained its pre-rename database name. Resolving full qualifiers through that detached table could dereference a missing database after the database was renamed. Resolve internal database qualifiers by stable database ID and read only the table name from live or recycled table metadata, with persisted names as fallback. ### Release note Streams with recycled base tables continue to expose valid qualifiers after database rename. ### Check List (For Author) - Test: Unit Test - ./run-fe-ut.sh --run org.apache.doris.catalog.DropTableStreamTest - Behavior changed: Yes. Stream metadata uses the current internal database name after database rename. - Does this need documentation: No --- .../doris/catalog/stream/BaseTableStream.java | 7 ++++ .../doris/catalog/DropTableStreamTest.java | 37 +++++++++++++++++++ 2 files changed, 44 insertions(+) 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 db395f1cf22794..45ab47b76d38bd 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 @@ -228,6 +228,13 @@ List getBaseTableFullQualifiers(TableIf availableBaseTable) { displayBaseTable = Env.getCurrentRecycleBin().getRecycledTableNullable( baseTableInfo.getDbId(), baseTableInfo.getTableId()); } + if (baseTableInfo.isInternalTable()) { + return ImmutableList.of( + baseTableInfo.getCtlName(), + Env.getCurrentInternalCatalog().getDb(baseTableInfo.getDbId()) + .map(db -> db.getFullName()).orElse(baseTableInfo.getDbName()), + displayBaseTable == null ? baseTableInfo.getTableName() : displayBaseTable.getName()); + } return displayBaseTable == null ? baseTableInfo.getFullQualifiers() : displayBaseTable.getFullQualifiers(); } diff --git a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java index 15474599212f28..961696e7520de3 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java @@ -220,6 +220,43 @@ public void testBaseTableQualifiersFollowRenameAndRecovery() throws Exception { Assertions.assertTrue(getStreamDdl(deserializedStream).contains("ON TABLE internal.test_stream.tbl_recovered")); } + @Test + public void testBaseTableQualifiersFollowDatabaseRenameAfterDrop() throws Exception { + createDatabase("test_stream_db_rename"); + createTable("create table test_stream_db_rename.tbl_db_rename (k1 int, k2 int) " + + "unique key(k1) distributed by hash(k1) buckets 1 " + + "properties('replication_num' = '1', 'binlog.enable' = 'true', 'binlog.format' = 'ROW', " + + "'binlog.need_historical_value' = 'true')"); + createTable("create stream test_stream_db_rename.s_db_rename " + + "on table test_stream_db_rename.tbl_db_rename " + + "properties('show_initial_rows' = 'true')"); + Database db = Env.getCurrentInternalCatalog().getDbOrMetaException("test_stream_db_rename"); + OlapTable baseTable = (OlapTable) db.getTableOrMetaException("tbl_db_rename"); + OlapTableStream stream = (OlapTableStream) db.getTableOrMetaException("s_db_rename"); + + dropTableWithSql("drop table test_stream_db_rename.tbl_db_rename"); + executeSql("alter database test_stream_db_rename rename test_stream_db_renamed"); + + TRow streamRow = getStreamMetadataRow("s_db_rename"); + Assertions.assertEquals("tbl_db_rename", streamRow.getColumnValue().get(6).getStringVal()); + Assertions.assertEquals("test_stream_db_renamed", streamRow.getColumnValue().get(7).getStringVal()); + Assertions.assertEquals("internal", streamRow.getColumnValue().get(8).getStringVal()); + Assertions.assertEquals("N/A", streamRow.getColumnValue().get(9).getStringVal()); + Assertions.assertFalse(streamRow.getColumnValue().get(10).isBoolVal()); + Assertions.assertTrue(streamRow.getColumnValue().get(11).isBoolVal()); + Assertions.assertTrue(getStreamDdl(stream) + .contains("ON TABLE internal.test_stream_db_renamed.tbl_db_rename")); + + OlapTableStream deserializedStream = GsonUtils.GSON.fromJson( + GsonUtils.GSON.toJson(stream), OlapTableStream.class); + Env.getCurrentRecycleBin().eraseTableInstantly(baseTable.getId()); + + Assertions.assertEquals(List.of("internal", "test_stream_db_renamed", "tbl_db_rename"), + deserializedStream.getBaseTableFullQualifiers()); + Assertions.assertTrue(getStreamDdl(deserializedStream) + .contains("ON TABLE internal.test_stream_db_renamed.tbl_db_rename")); + } + @Test public void testForceDropAndSameNameTableDoNotRestoreStream() throws Exception { createBaseTableAndStream("tbl_force", "s_force"); From 7ea00646773495a9cd64d714d27441575e8539af Mon Sep 17 00:00:00 2001 From: seawinde Date: Wed, 5 Aug 2026 14:57:28 +0800 Subject: [PATCH 4/8] [fix](stream) Lock stream base table by stable ID ### What problem does this PR solve? Issue Number: None Related PR: #66287 Problem Summary: Stream relation collection previously re-resolved a renamed or recovered base table by display name, which could cache and lock a same-name replacement while binding later resolved the original table by stable ID. Cache the ID-resolved base table directly so the planner locks the object used during stream binding, while preserving the existing missing-table error message. ### Release note None ### Check List (For Author) - Test: FE unit test - ./run-fe-ut.sh --run org.apache.doris.catalog.DropTableStreamTest - Behavior changed: No - Does this need documentation: No --- .../rules/analysis/CollectRelation.java | 28 ++++++++++++++++--- 1 file changed, 24 insertions(+), 4 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java index 34819b568213da..b2b09ae3b06a63 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java @@ -237,7 +237,7 @@ private void collectFromUnboundRelation(CascadesContext cascadesContext, } // we need to collect stream table's base table as well if (table instanceof BaseTableStream) { - collectFromTableStream((BaseTableStream) table, cascadesContext, tableFrom, unboundRelation); + collectFromTableStream((BaseTableStream) table, cascadesContext, tableFrom); } } @@ -314,9 +314,29 @@ protected void parseAndCollectFromView(List tableQualifier, View view, C } private void collectFromTableStream(BaseTableStream tableStream, CascadesContext cascadesContext, - TableFrom tableFrom, Optional unboundRelation) { + TableFrom tableFrom) { StatementContext statementContext = cascadesContext.getConnectContext().getStatementContext(); - List tableQualifier = tableStream.getBaseTableFullQualifiers(); - statementContext.getAndCacheTable(tableQualifier, tableFrom, unboundRelation); + TableIf baseTable = tableStream.getBaseTableNullable(); + if (baseTable == null) { + throw new AnalysisException("Table [" + + tableStream.getBaseTableFullQualifiers().get(2) + "] does not exist"); + } + + // Cache the ID-resolved base table so planner locks the object used during binding. + Map, TableIf> tables; + switch (tableFrom) { + case QUERY: + tables = statementContext.getTables(); + break; + case INSERT_TARGET: + tables = statementContext.getInsertTargetTables(); + break; + case MTMV: + tables = statementContext.getMtmvRelatedTables(); + break; + default: + throw new AnalysisException("Unknown table from " + tableFrom); + } + tables.put(baseTable.getFullQualifiers(), baseTable); } } From a7850adc43aef51b33514ac9c14d2572afd2a702 Mon Sep 17 00:00:00 2001 From: seawinde Date: Tue, 11 Aug 2026 07:46:28 +0800 Subject: [PATCH 5/8] [fix](stream) Keep stream lock dependencies out of relation cache ### What problem does this PR solve? Issue Number: close #65418 Related PR: #66287 Problem Summary: The previous stable-ID locking fix cached a stream base table under its current name qualifier. A pre-lock rename could therefore overwrite an explicit relation already resolved under the same qualifier and make binding scan the wrong table. Keep the relation cache limited to SQL relations and expand each stream base table by stable ID only when constructing the existing ordered lock queue. ### Release note None ### Check List (For Author) - Test: FE unit test and FE build - ./run-fe-ut.sh --run org.apache.doris.nereids.StatementContextTest,org.apache.doris.nereids.trees.plans.ExplainTableStreamPlanTest - ./run-fe-ut.sh --run org.apache.doris.catalog.DropTableStreamTest - env DISABLE_BUILD_UI=ON ./build.sh --fe - Behavior changed: Yes. Explicit relation bindings are preserved while stream base tables are locked by stable ID. - Does this need documentation: No --- .../doris/nereids/StatementContext.java | 22 +++++- .../rules/analysis/CollectRelation.java | 26 +------ .../doris/nereids/StatementContextTest.java | 73 +++++++++++++++++++ .../plans/ExplainTableStreamPlanTest.java | 69 ++++++++++++++++++ 4 files changed, 164 insertions(+), 26 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java index ae37ee567fed47..56855a30a8c43e 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java @@ -27,6 +27,7 @@ import org.apache.doris.catalog.Partition; import org.apache.doris.catalog.TableIf; import org.apache.doris.catalog.View; +import org.apache.doris.catalog.stream.BaseTableStream; import org.apache.doris.common.Id; import org.apache.doris.common.IdGenerator; import org.apache.doris.common.Pair; @@ -966,9 +967,9 @@ public synchronized void lock() { PriorityQueue tableIfs = new PriorityQueue<>( tables.size() + mtmvRelatedTables.size() + insertTargetTables.size(), Comparator.comparing(TableIf::getId)); - tableIfs.addAll(tables.values()); - tableIfs.addAll(mtmvRelatedTables.values()); - tableIfs.addAll(insertTargetTables.values()); + addTablesToLock(tableIfs, tables.values()); + addTablesToLock(tableIfs, mtmvRelatedTables.values()); + addTablesToLock(tableIfs, insertTargetTables.values()); while (!tableIfs.isEmpty()) { TableIf tableIf = tableIfs.poll(); if (!tableIf.needReadLockWhenPlan()) { @@ -1292,10 +1293,25 @@ private boolean containsPlanReadLockTable(Collection tableIfs) { if (tableIf.needReadLockWhenPlan()) { return true; } + if (tableIf instanceof BaseTableStream) { + TableIf baseTable = ((BaseTableStream) tableIf).getBaseTableNullable(); + if (baseTable != null && baseTable.needReadLockWhenPlan()) { + return true; + } + } } return false; } + private void addTablesToLock(PriorityQueue tableIfs, Collection tables) { + tableIfs.addAll(tables); + for (TableIf tableIf : tables) { + if (tableIf instanceof BaseTableStream) { + tableIfs.add(((BaseTableStream) tableIf).getBaseTableOrNereidsAnalysisException()); + } + } + } + private static class CloseableResource implements Closeable { public final String resourceName; public final String threadName; diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java index b2b09ae3b06a63..c7ce36d804c01d 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java @@ -237,7 +237,7 @@ private void collectFromUnboundRelation(CascadesContext cascadesContext, } // we need to collect stream table's base table as well if (table instanceof BaseTableStream) { - collectFromTableStream((BaseTableStream) table, cascadesContext, tableFrom); + collectFromTableStream((BaseTableStream) table); } } @@ -313,30 +313,10 @@ protected void parseAndCollectFromView(List tableQualifier, View view, C parentContext.addPlanProcesses(viewContext.getPlanProcesses()); } - private void collectFromTableStream(BaseTableStream tableStream, CascadesContext cascadesContext, - TableFrom tableFrom) { - StatementContext statementContext = cascadesContext.getConnectContext().getStatementContext(); - TableIf baseTable = tableStream.getBaseTableNullable(); - if (baseTable == null) { + private void collectFromTableStream(BaseTableStream tableStream) { + if (tableStream.getBaseTableNullable() == null) { throw new AnalysisException("Table [" + tableStream.getBaseTableFullQualifiers().get(2) + "] does not exist"); } - - // Cache the ID-resolved base table so planner locks the object used during binding. - Map, TableIf> tables; - switch (tableFrom) { - case QUERY: - tables = statementContext.getTables(); - break; - case INSERT_TARGET: - tables = statementContext.getInsertTargetTables(); - break; - case MTMV: - tables = statementContext.getMtmvRelatedTables(); - break; - default: - throw new AnalysisException("Unknown table from " + tableFrom); - } - tables.put(baseTable.getFullQualifiers(), baseTable); } } diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/StatementContextTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/StatementContextTest.java index 70a29a23f05008..ea0643ef203b95 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/StatementContextTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/StatementContextTest.java @@ -21,6 +21,7 @@ import org.apache.doris.analysis.TableSnapshot; import org.apache.doris.catalog.DatabaseIf; import org.apache.doris.catalog.TableIf; +import org.apache.doris.catalog.stream.BaseTableStream; import org.apache.doris.datasource.CatalogIf; import org.apache.doris.datasource.mvcc.MvccSnapshot; import org.apache.doris.datasource.mvcc.PluginDrivenMvccExternalTable; @@ -170,6 +171,78 @@ public void testSkipPreloadWhenNoInternalTableNeedsPlanReadLock() { } } + @Test + public void testLockIncludesStreamBaseTableWithoutReplacingRelationCache() { + ConnectContext connectContext = Mockito.mock(ConnectContext.class); + TableIf explicitTable = Mockito.mock(TableIf.class); + BaseTableStream stream = Mockito.mock(BaseTableStream.class); + TableIf baseTable = Mockito.mock(TableIf.class); + Mockito.when(explicitTable.getId()).thenReturn(11L); + Mockito.when(explicitTable.getName()).thenReturn("explicit"); + Mockito.when(explicitTable.getNameWithFullQualifiers()).thenReturn("internal.db.explicit"); + Mockito.when(explicitTable.needReadLockWhenPlan()).thenReturn(true); + Mockito.when(explicitTable.tryReadLock(Mockito.anyLong(), Mockito.any())).thenReturn(true); + Mockito.when(stream.getId()).thenReturn(12L); + Mockito.when(stream.needReadLockWhenPlan()).thenReturn(false); + Mockito.when(stream.getBaseTableOrNereidsAnalysisException()).thenReturn(baseTable); + Mockito.when(baseTable.getId()).thenReturn(13L); + Mockito.when(baseTable.getName()).thenReturn("base"); + Mockito.when(baseTable.getNameWithFullQualifiers()).thenReturn("internal.db.base"); + Mockito.when(baseTable.needReadLockWhenPlan()).thenReturn(true); + Mockito.when(baseTable.tryReadLock(Mockito.anyLong(), Mockito.any())).thenReturn(true); + + StatementContext statementContext = new StatementContext(connectContext, + new OriginStatement("select * from db.explicit join db.stream", 0)); + try { + statementContext.getTables().put(ImmutableList.of("internal", "db", "explicit"), explicitTable); + statementContext.getTables().put(ImmutableList.of("internal", "db", "stream"), stream); + + statementContext.lock(); + + Mockito.verify(explicitTable).tryReadLock(Mockito.anyLong(), Mockito.any()); + Mockito.verify(baseTable).tryReadLock(Mockito.anyLong(), Mockito.any()); + Mockito.verify(stream, Mockito.never()).tryReadLock(Mockito.anyLong(), Mockito.any()); + org.junit.jupiter.api.Assertions.assertSame( + explicitTable, statementContext.getTables().get(ImmutableList.of("internal", "db", "explicit"))); + } finally { + statementContext.close(); + } + } + + @Test + public void testPreloadRecognizesStreamBaseTablePlanReadLock() { + ConnectContext connectContext = Mockito.mock(ConnectContext.class); + BaseTableStream stream = Mockito.mock(BaseTableStream.class); + TableIf baseTable = Mockito.mock(TableIf.class); + PluginDrivenExternalTable externalTable = Mockito.mock(PluginDrivenExternalTable.class); + SessionVariable sessionVariable = new SessionVariable(); + sessionVariable.setEnablePreloadExternalMetadata(true); + + Mockito.when(connectContext.getSessionVariable()).thenReturn(sessionVariable); + Mockito.when(connectContext.getQueryIdentifier()).thenReturn("stream-preload"); + Mockito.when(stream.needReadLockWhenPlan()).thenReturn(false); + Mockito.when(stream.getBaseTableNullable()).thenReturn(baseTable); + Mockito.when(baseTable.needReadLockWhenPlan()).thenReturn(true); + Mockito.when(externalTable.getId()).thenReturn(21L); + Mockito.when(externalTable.supportsExternalMetadataPreload()).thenReturn(true); + Mockito.when(externalTable.getBaseSchema()).thenReturn(Collections.emptyList()); + Mockito.when(externalTable.supportInternalPartitionPruned()).thenReturn(false); + + StatementContext statementContext = new StatementContext(connectContext, + new OriginStatement("select * from db.stream join ext", 0)); + try { + statementContext.getTables().put(ImmutableList.of("internal", "db", "stream"), stream); + statementContext.registerExternalTableForPreload(externalTable, Optional.empty(), Optional.empty()); + + ExternalMetadataPreloadResult result = executePreload(statementContext); + + org.junit.jupiter.api.Assertions.assertTrue(result.isExecuted()); + Mockito.verify(externalTable).getBaseSchema(); + } finally { + statementContext.close(); + } + } + @Test public void testPreloadIcebergLatestSnapshotBeforeLock() { ConnectContext connectContext = Mockito.mock(ConnectContext.class); diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java index 7e9288721e433d..e6bb3b783a05c4 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java @@ -67,10 +67,16 @@ import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import org.mockito.Mockito; import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; /** * UTs for table stream query plan, including @@ -652,6 +658,69 @@ public void testRecordPlanForMvPreRewriteNormalizesStreamScanInsideCte() throws } } + @Test + public void testPreLockRenameDoesNotReplaceCachedExplicitTable() throws Exception { + String baseTableName = "stream_cache_base"; + String streamName = "stream_cache_stream"; + String explicitTableName = "stream_cache_explicit"; + createTable("create table test_stream." + baseTableName + " (k1 int, k2 int) " + + "unique key(k1) distributed by hash(k1) buckets 1 " + + "properties('replication_num' = '1', 'binlog.enable' = 'true', " + + "'binlog.format' = 'ROW', 'binlog.need_historical_value' = 'true')"); + createTable("create stream test_stream." + streamName + " on table test_stream." + baseTableName + + " properties('show_initial_rows' = 'true')"); + createTable("create table test_stream." + explicitTableName + " (k1 int, k2 int) " + + "unique key(k1) distributed by hash(k1) buckets 1 " + + "properties('replication_num' = '1')"); + + Database db = (Database) Env.getCurrentInternalCatalog().getDbOrMetaException("test_stream"); + OlapTable baseTable = (OlapTable) db.getTableOrMetaException(baseTableName); + OlapTable explicitTable = (OlapTable) db.getTableOrMetaException(explicitTableName); + OlapTableStream stream = (OlapTableStream) db.getTableOrMetaException(streamName); + OlapTableStream blockingStream = Mockito.spy(stream); + CountDownLatch streamCollectionStarted = new CountDownLatch(1); + CountDownLatch renamesFinished = new CountDownLatch(1); + Mockito.doAnswer(invocation -> { + streamCollectionStarted.countDown(); + Assertions.assertTrue(renamesFinished.await(10, TimeUnit.SECONDS)); + return invocation.callRealMethod(); + }).when(blockingStream).getBaseTableNullable(); + + ConnectContext ctx = createDefaultCtx(); + ctx.setDatabase("test_stream"); + String sql = "select b.k1, s.k1 from test_stream." + explicitTableName + " b join test_stream." + + streamName + " s on b.k1 = s.k1"; + StatementContext statementContext = MemoTestUtils.createStatementContext(ctx, sql); + statementContext.getTables().put(List.of("internal", "test_stream", streamName), blockingStream); + ExecutorService plannerExecutor = Executors.newSingleThreadExecutor(); + try { + Future planFuture = plannerExecutor.submit(() -> { + ctx.setThreadLocalInfo(); + new NereidsPlanner(statementContext).planWithLock( + (LogicalPlan) parser.parseSingle(sql), PhysicalProperties.ANY); + }); + + Assertions.assertTrue(streamCollectionStarted.await(10, TimeUnit.SECONDS)); + Assertions.assertSame(explicitTable, + statementContext.getTables().get(List.of("internal", "test_stream", explicitTableName))); + + connectContext.setThreadLocalInfo(); + executeSql("alter table test_stream." + explicitTableName + + " rename " + explicitTableName + "_renamed"); + executeSql("alter table test_stream." + baseTableName + " rename " + explicitTableName); + renamesFinished.countDown(); + planFuture.get(30, TimeUnit.SECONDS); + + Assertions.assertSame(explicitTable, + statementContext.getTables().get(List.of("internal", "test_stream", explicitTableName))); + Assertions.assertSame(baseTable, blockingStream.getBaseTableNullable()); + } finally { + renamesFinished.countDown(); + plannerExecutor.shutdownNow(); + statementContext.close(); + } + } + @Test public void testMowTimeTravelQualifiedColumnCanBind() { // MOW time-travel goes through a union whose outputs are rebuilt with empty qualifiers. From 9c8100c412da6993e72ab0f7d450f338dff87870 Mon Sep 17 00:00:00 2001 From: seawinde Date: Tue, 11 Aug 2026 11:54:29 +0800 Subject: [PATCH 6/8] [chore](stream) Document stream base lock handling ### What problem does this PR solve? Issue Number: close #65418 Related PR: #66287 Problem Summary: Document why stream base tables are validated by stable ID without entering the qualifier relation cache, and why implicit base dependencies are expanded only in the local ID-ordered planner lock queue. The comments preserve the concurrency invariant behind the relation-cache collision fix. ### Release note None ### Check List (For Author) - Test: FE unit test - ./run-fe-ut.sh --run org.apache.doris.nereids.StatementContextTest,org.apache.doris.nereids.trees.plans.ExplainTableStreamPlanTest - Behavior changed: No - Does this need documentation: No --- .../java/org/apache/doris/nereids/StatementContext.java | 6 ++++++ .../doris/nereids/rules/analysis/CollectRelation.java | 6 +++++- 2 files changed, 11 insertions(+), 1 deletion(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java index 56855a30a8c43e..b43381c44aba24 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java @@ -1294,6 +1294,7 @@ private boolean containsPlanReadLockTable(Collection tableIfs) { return true; } if (tableIf instanceof BaseTableStream) { + // Mirror addTablesToLock(): a stream needs no plan lock itself, but its stable-ID base may need one. TableIf baseTable = ((BaseTableStream) tableIf).getBaseTableNullable(); if (baseTable != null && baseTable.needReadLockWhenPlan()) { return true; @@ -1303,6 +1304,11 @@ private boolean containsPlanReadLockTable(Collection tableIfs) { return false; } + /** + * Add explicit relations and stable-ID stream bases to the local ID-ordered lock queue. Stream bases must not be + * added to the qualifier relation maps because a concurrent rename can make a base qualifier collide with an + * explicitly resolved relation. + */ private void addTablesToLock(PriorityQueue tableIfs, Collection tables) { tableIfs.addAll(tables); for (TableIf tableIf : tables) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java index c7ce36d804c01d..fa2338a269cdda 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java @@ -235,7 +235,6 @@ private void collectFromUnboundRelation(CascadesContext cascadesContext, if (table instanceof View) { parseAndCollectFromView(tableQualifier, (View) table, cascadesContext); } - // we need to collect stream table's base table as well if (table instanceof BaseTableStream) { collectFromTableStream((BaseTableStream) table); } @@ -313,6 +312,11 @@ protected void parseAndCollectFromView(List tableQualifier, View view, C parentContext.addPlanProcesses(viewContext.getPlanProcesses()); } + /** + * Validate the stream base by stable ID without caching it under display qualifiers. A concurrent rename can + * make that qualifier belong to an explicitly resolved relation; stream lock dependencies are expanded separately + * in {@link StatementContext#lock()}. + */ private void collectFromTableStream(BaseTableStream tableStream) { if (tableStream.getBaseTableNullable() == null) { throw new AnalysisException("Table [" From 25d0c27f52c20efa8496732fcb0ab457e21e9377 Mon Sep 17 00:00:00 2001 From: seawinde Date: Tue, 11 Aug 2026 17:58:00 +0800 Subject: [PATCH 7/8] [fix](stream) Track stream base dependencies by identity ### What problem does this PR solve? Issue Number: close #65418 Related PR: #66287 Problem Summary: Stream base tables are implicit planning dependencies resolved by stable table ID. Keeping them in the qualifier-keyed relation cache can overwrite an explicitly resolved relation after concurrent renames, while resolving them again during preload or locking can select a different object snapshot. MTMV relation generation also omitted the underlying stream base after it was removed from the relation cache. Store exact implicit dependency objects separately, reuse them for preload and locking, include them in all-level MTMV dependencies, and reject dropped tables during cache-miss recovery resolution. ### Release note Streams continue to track and lock their original base tables by stable identity without corrupting explicit relation bindings. ### Check List (For Author) - Test: Unit Test - 43 FE tests passed across DropTableStreamTest, StatementContextTest, ExplainTableStreamPlanTest, and MTMVRelationTest - StatementContextTest rerun: 11 tests passed - Java compilation and Checkstyle passed - Behavior changed: Yes. Stream base dependencies use a stable per-statement object snapshot for planning locks and MTMV dependency tracking. - Does this need documentation: No --- .../doris/catalog/stream/BaseTableStream.java | 19 ++++--- .../org/apache/doris/mtmv/MTMVPlanUtil.java | 7 ++- .../doris/nereids/StatementContext.java | 52 ++++++++----------- .../rules/analysis/CollectRelation.java | 14 +++-- .../doris/catalog/DropTableStreamTest.java | 42 +++++++++++++++ .../apache/doris/mtmv/MTMVRelationTest.java | 50 ++++++++++++++++++ .../doris/nereids/StatementContextTest.java | 24 +++++++-- .../plans/ExplainTableStreamPlanTest.java | 2 +- 8 files changed, 159 insertions(+), 51 deletions(-) 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 45ab47b76d38bd..9ff289cd52c48b 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 @@ -117,15 +117,20 @@ public BaseTableStream(String streamName, List fullSchema, TableIf baseT public TableIf getBaseTableNullable() { TableIf cachedBaseTable = baseTable; - if (cachedBaseTable instanceof Table && ((Table) cachedBaseTable).isDropped) { - baseTable = null; - return null; + if (cachedBaseTable != null) { + if (cachedBaseTable instanceof Table && ((Table) cachedBaseTable).isDropped) { + baseTable = null; + return null; + } + return cachedBaseTable; } - if (cachedBaseTable == null) { - cachedBaseTable = baseTableInfo.getTableNullable(); - baseTable = cachedBaseTable; + TableIf resolvedBaseTable = baseTableInfo.getTableNullable(); + // Recovery publishes the table into database maps before clearing its dropped flag. + if (resolvedBaseTable instanceof Table && ((Table) resolvedBaseTable).isDropped) { + return null; } - return cachedBaseTable; + baseTable = resolvedBaseTable; + return resolvedBaseTable; } public void setProperties(Map properties) throws org.apache.doris.common.AnalysisException { diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPlanUtil.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPlanUtil.java index 4be21db3008033..ef6e3d8df4a018 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPlanUtil.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVPlanUtil.java @@ -224,7 +224,10 @@ public static Pair, Set> getBaseTableFromQuery(String quer try { NereidsPlanner planner = new NereidsPlanner(ctx.getStatementContext()); planner.planWithLock(logicalPlan, PhysicalProperties.ANY, ExplainLevel.ANALYZED_PLAN); - return Pair.of(Sets.newHashSet(ctx.getStatementContext().getTables().values()), + Set baseTables = Sets.newHashSet(ctx.getStatementContext().getTables().values()); + // Implicit dependencies are all-level tables, not relations written at the first query level. + baseTables.addAll(ctx.getStatementContext().getImplicitTableDependencies()); + return Pair.of(baseTables, Sets.newHashSet(ctx.getStatementContext().getOneLevelTables().values())); } finally { ctx.setStatementContext(original); @@ -448,6 +451,8 @@ public static MTMVAnalyzeQueryInfo analyzeQuery(ConnectContext ctx, Map baseTables = Sets.newHashSet(statementContext.getTables().values()); + // Implicit dependencies are all-level tables, not relations written at the first query level. + baseTables.addAll(statementContext.getImplicitTableDependencies()); Set oneLevelTables = Sets.newHashSet(statementContext.getOneLevelTables().values()); for (TableIf table : baseTables) { if (table.isTemporary()) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java index b43381c44aba24..058315a7ec95b4 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/StatementContext.java @@ -27,7 +27,6 @@ import org.apache.doris.catalog.Partition; import org.apache.doris.catalog.TableIf; import org.apache.doris.catalog.View; -import org.apache.doris.catalog.stream.BaseTableStream; import org.apache.doris.common.Id; import org.apache.doris.common.IdGenerator; import org.apache.doris.common.Pair; @@ -220,6 +219,9 @@ public enum TableFrom { // tables in this query directly private final Map, TableIf> tables = Maps.newHashMap(); + // Underlying tables resolved while collecting explicit relations. They are not independently named by SQL, so + // keep them out of the qualifier-keyed relation maps and reuse the same snapshot through planning. + private final Set implicitTableDependencies = Sets.newIdentityHashSet(); // onelevel tables in this query directly, // if // create v1 as select * from t1 @@ -443,6 +445,14 @@ public Map, TableIf> getTables() { return tables; } + public void addImplicitTableDependency(TableIf table) { + implicitTableDependencies.add(table); + } + + public Set getImplicitTableDependencies() { + return implicitTableDependencies; + } + public Map, TableIf> getOneLevelTables() { return oneLevelTables; } @@ -960,16 +970,20 @@ public Map getRelationIdToStatisticsMap() { */ public synchronized void lock() { if (!needLockTables - || (tables.isEmpty() && mtmvRelatedTables.isEmpty() && insertTargetTables.isEmpty()) + || (tables.isEmpty() && mtmvRelatedTables.isEmpty() && insertTargetTables.isEmpty() + && implicitTableDependencies.isEmpty()) || !plannerResources.isEmpty()) { return; } + // The same object can be both an explicit relation and an implicit dependency; lock it only once. + Set tablesToLock = Sets.newIdentityHashSet(); + tablesToLock.addAll(tables.values()); + tablesToLock.addAll(mtmvRelatedTables.values()); + tablesToLock.addAll(insertTargetTables.values()); + tablesToLock.addAll(implicitTableDependencies); PriorityQueue tableIfs = new PriorityQueue<>( - tables.size() + mtmvRelatedTables.size() + insertTargetTables.size(), - Comparator.comparing(TableIf::getId)); - addTablesToLock(tableIfs, tables.values()); - addTablesToLock(tableIfs, mtmvRelatedTables.values()); - addTablesToLock(tableIfs, insertTargetTables.values()); + tablesToLock.size(), Comparator.comparing(TableIf::getId)); + tableIfs.addAll(tablesToLock); while (!tableIfs.isEmpty()) { TableIf tableIf = tableIfs.poll(); if (!tableIf.needReadLockWhenPlan()) { @@ -1277,7 +1291,8 @@ public int getExternalTablePreloadCandidateCount() { public boolean hasAnyPlanReadLockTable() { return containsPlanReadLockTable(tables.values()) || containsPlanReadLockTable(mtmvRelatedTables.values()) - || containsPlanReadLockTable(insertTargetTables.values()); + || containsPlanReadLockTable(insertTargetTables.values()) + || containsPlanReadLockTable(implicitTableDependencies); } public Optional getExternalMetadataPreloadResult() { @@ -1293,31 +1308,10 @@ private boolean containsPlanReadLockTable(Collection tableIfs) { if (tableIf.needReadLockWhenPlan()) { return true; } - if (tableIf instanceof BaseTableStream) { - // Mirror addTablesToLock(): a stream needs no plan lock itself, but its stable-ID base may need one. - TableIf baseTable = ((BaseTableStream) tableIf).getBaseTableNullable(); - if (baseTable != null && baseTable.needReadLockWhenPlan()) { - return true; - } - } } return false; } - /** - * Add explicit relations and stable-ID stream bases to the local ID-ordered lock queue. Stream bases must not be - * added to the qualifier relation maps because a concurrent rename can make a base qualifier collide with an - * explicitly resolved relation. - */ - private void addTablesToLock(PriorityQueue tableIfs, Collection tables) { - tableIfs.addAll(tables); - for (TableIf tableIf : tables) { - if (tableIf instanceof BaseTableStream) { - tableIfs.add(((BaseTableStream) tableIf).getBaseTableOrNereidsAnalysisException()); - } - } - } - private static class CloseableResource implements Closeable { public final String resourceName; public final String threadName; diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java index fa2338a269cdda..c4d1e62148fda9 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java @@ -236,7 +236,7 @@ private void collectFromUnboundRelation(CascadesContext cascadesContext, parseAndCollectFromView(tableQualifier, (View) table, cascadesContext); } if (table instanceof BaseTableStream) { - collectFromTableStream((BaseTableStream) table); + collectFromTableStream((BaseTableStream) table, cascadesContext.getStatementContext()); } } @@ -312,15 +312,13 @@ protected void parseAndCollectFromView(List tableQualifier, View view, C parentContext.addPlanProcesses(viewContext.getPlanProcesses()); } - /** - * Validate the stream base by stable ID without caching it under display qualifiers. A concurrent rename can - * make that qualifier belong to an explicitly resolved relation; stream lock dependencies are expanded separately - * in {@link StatementContext#lock()}. - */ - private void collectFromTableStream(BaseTableStream tableStream) { - if (tableStream.getBaseTableNullable() == null) { + private void collectFromTableStream(BaseTableStream tableStream, StatementContext statementContext) { + // Capture the stable-ID result once so preload, locking, and dependency tracking use the same table object. + TableIf baseTable = tableStream.getBaseTableNullable(); + if (baseTable == null) { throw new AnalysisException("Table [" + tableStream.getBaseTableFullQualifiers().get(2) + "] does not exist"); } + statementContext.addImplicitTableDependency(baseTable); } } diff --git a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java index 961696e7520de3..28aeee3935f6e6 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java @@ -33,9 +33,15 @@ import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import org.mockito.Mockito; import java.util.ArrayList; import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; public class DropTableStreamTest extends TestWithFeService { @@ -178,6 +184,42 @@ public void testStreamStateFollowsRecoverableBaseTableDrop() throws Exception { Assertions.assertEquals("N/A", stream.getStaleReason()); } + @Test + public void testRecoveringBaseTableRemainsUnavailableUntilUnmarkedDropped() throws Exception { + createBaseTableAndStream("tbl_recovery_publish", "s_recovery_publish"); + Database db = Env.getCurrentInternalCatalog().getDbOrMetaException("test_stream"); + OlapTable baseTable = (OlapTable) db.getTableOrMetaException("tbl_recovery_publish"); + OlapTableStream stream = (OlapTableStream) db.getTableOrMetaException("s_recovery_publish"); + + baseTable.markDropped(); + db.unregisterTable(baseTable.getId()); + Assertions.assertNull(stream.getBaseTableNullable()); + + OlapTable recoveringBaseTable = Mockito.spy(baseTable); + CountDownLatch tablePublished = new CountDownLatch(1); + CountDownLatch allowUnmarkDropped = new CountDownLatch(1); + Mockito.doAnswer(invocation -> { + tablePublished.countDown(); + Assertions.assertTrue(allowUnmarkDropped.await(10, TimeUnit.SECONDS)); + return invocation.callRealMethod(); + }).when(recoveringBaseTable).unmarkDropped(); + + ExecutorService executor = Executors.newSingleThreadExecutor(); + try { + Future registerFuture = executor.submit(() -> db.registerTable(recoveringBaseTable)); + Assertions.assertTrue(tablePublished.await(10, TimeUnit.SECONDS)); + + Assertions.assertNull(stream.getBaseTableNullable()); + + allowUnmarkDropped.countDown(); + Assertions.assertTrue(registerFuture.get(10, TimeUnit.SECONDS)); + Assertions.assertSame(recoveringBaseTable, stream.getBaseTableNullable()); + } finally { + allowUnmarkDropped.countDown(); + executor.shutdownNow(); + } + } + @Test public void testBaseTableQualifiersFollowRenameAndRecovery() throws Exception { createBaseTableAndStream("tbl_rename", "s_rename"); diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelationTest.java b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelationTest.java index 3b93daca9b7453..43f1f75e3c1313 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelationTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelationTest.java @@ -20,13 +20,25 @@ import org.apache.doris.catalog.Database; import org.apache.doris.catalog.Env; import org.apache.doris.catalog.MTMV; +import org.apache.doris.catalog.TableIf; +import org.apache.doris.common.Config; +import org.apache.doris.common.Pair; import org.apache.doris.utframe.TestWithFeService; import com.google.common.collect.Sets; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.util.Set; + public class MTMVRelationTest extends TestWithFeService { + + @Override + protected void runBeforeAll() throws Exception { + Config.enable_table_stream = true; + Config.enable_feature_binlog = true; + } + // t1 => v1 => v2 // t2 => mv1 // mv1 join v2 => mv2 @@ -114,4 +126,42 @@ public void testMTMVRelation() throws Exception { Assertions.assertEquals(Sets.newHashSet(), relationManager.getMtmvsByBaseView(v1)); Assertions.assertEquals(Sets.newHashSet(), relationManager.getMtmvsByBaseView(v2)); } + + @Test + public void testMTMVRelationIncludesStreamBaseDependency() throws Exception { + createDatabaseAndUse("stream_mtmv_db"); + createTables( + "CREATE TABLE stream_base (k1 int, k2 int)\n" + + "UNIQUE KEY(k1)\n" + + "DISTRIBUTED BY HASH(k1) BUCKETS 1\n" + + "PROPERTIES ('replication_num' = '1', 'binlog.enable' = 'true',\n" + + "'binlog.format' = 'ROW', 'binlog.need_historical_value' = 'true')", + "CREATE STREAM stream_source ON TABLE stream_base\n" + + "PROPERTIES ('show_initial_rows' = 'true')"); + createMvByNereids("CREATE MATERIALIZED VIEW stream_mv BUILD DEFERRED\n" + + "REFRESH COMPLETE ON MANUAL\n" + + "DISTRIBUTED BY RANDOM BUCKETS 1\n" + + "PROPERTIES ('replication_num' = '1')\n" + + "AS SELECT k1, k2 FROM stream_source"); + + Database db = Env.getCurrentEnv().getInternalCatalog().getDbOrAnalysisException("stream_mtmv_db"); + MTMV mtmv = (MTMV) db.getTableOrAnalysisException("stream_mv"); + TableIf stream = db.getTableOrAnalysisException("stream_source"); + TableIf baseTable = db.getTableOrAnalysisException("stream_base"); + BaseTableInfo streamInfo = new BaseTableInfo(stream); + BaseTableInfo baseTableInfo = new BaseTableInfo(baseTable); + BaseTableInfo mtmvInfo = new BaseTableInfo(mtmv); + + Assertions.assertEquals(Sets.newHashSet(streamInfo, baseTableInfo), mtmv.getRelation().getBaseTables()); + Assertions.assertEquals(Sets.newHashSet(streamInfo), mtmv.getRelation().getBaseTablesOneLevel()); + Assertions.assertEquals(Sets.newHashSet(streamInfo, baseTableInfo), + mtmv.getRelation().getBaseTablesOneLevelAndFromView()); + Assertions.assertEquals(Sets.newHashSet(mtmvInfo), Env.getCurrentEnv().getMtmvService().getRelationManager() + .getMtmvsByBaseTable(baseTableInfo)); + + Pair, Set> queryTables = MTMVPlanUtil.getBaseTableFromQuery( + mtmv.getQuerySql(), connectContext); + Assertions.assertEquals(Sets.newHashSet(stream, baseTable), queryTables.first); + Assertions.assertEquals(Sets.newHashSet(stream), queryTables.second); + } } diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/StatementContextTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/StatementContextTest.java index ea0643ef203b95..02e145ab40fa9e 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/StatementContextTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/StatementContextTest.java @@ -184,7 +184,6 @@ public void testLockIncludesStreamBaseTableWithoutReplacingRelationCache() { Mockito.when(explicitTable.tryReadLock(Mockito.anyLong(), Mockito.any())).thenReturn(true); Mockito.when(stream.getId()).thenReturn(12L); Mockito.when(stream.needReadLockWhenPlan()).thenReturn(false); - Mockito.when(stream.getBaseTableOrNereidsAnalysisException()).thenReturn(baseTable); Mockito.when(baseTable.getId()).thenReturn(13L); Mockito.when(baseTable.getName()).thenReturn("base"); Mockito.when(baseTable.getNameWithFullQualifiers()).thenReturn("internal.db.base"); @@ -196,12 +195,17 @@ public void testLockIncludesStreamBaseTableWithoutReplacingRelationCache() { try { statementContext.getTables().put(ImmutableList.of("internal", "db", "explicit"), explicitTable); statementContext.getTables().put(ImmutableList.of("internal", "db", "stream"), stream); + statementContext.getTables().put(ImmutableList.of("internal", "db", "base"), baseTable); + statementContext.addImplicitTableDependency(baseTable); + statementContext.addImplicitTableDependency(baseTable); statementContext.lock(); - Mockito.verify(explicitTable).tryReadLock(Mockito.anyLong(), Mockito.any()); - Mockito.verify(baseTable).tryReadLock(Mockito.anyLong(), Mockito.any()); + InOrder lockOrder = Mockito.inOrder(explicitTable, baseTable); + lockOrder.verify(explicitTable, Mockito.times(1)).tryReadLock(Mockito.anyLong(), Mockito.any()); + lockOrder.verify(baseTable, Mockito.times(1)).tryReadLock(Mockito.anyLong(), Mockito.any()); Mockito.verify(stream, Mockito.never()).tryReadLock(Mockito.anyLong(), Mockito.any()); + Mockito.verify(stream, Mockito.never()).getBaseTableNullable(); org.junit.jupiter.api.Assertions.assertSame( explicitTable, statementContext.getTables().get(ImmutableList.of("internal", "db", "explicit"))); } finally { @@ -210,10 +214,11 @@ public void testLockIncludesStreamBaseTableWithoutReplacingRelationCache() { } @Test - public void testPreloadRecognizesStreamBaseTablePlanReadLock() { + public void testPreloadAndLockUseSameStreamBaseSnapshot() { ConnectContext connectContext = Mockito.mock(ConnectContext.class); BaseTableStream stream = Mockito.mock(BaseTableStream.class); TableIf baseTable = Mockito.mock(TableIf.class); + TableIf recoveredBaseTable = Mockito.mock(TableIf.class); PluginDrivenExternalTable externalTable = Mockito.mock(PluginDrivenExternalTable.class); SessionVariable sessionVariable = new SessionVariable(); sessionVariable.setEnablePreloadExternalMetadata(true); @@ -221,8 +226,12 @@ public void testPreloadRecognizesStreamBaseTablePlanReadLock() { Mockito.when(connectContext.getSessionVariable()).thenReturn(sessionVariable); Mockito.when(connectContext.getQueryIdentifier()).thenReturn("stream-preload"); Mockito.when(stream.needReadLockWhenPlan()).thenReturn(false); - Mockito.when(stream.getBaseTableNullable()).thenReturn(baseTable); + Mockito.when(stream.getBaseTableNullable()).thenReturn(recoveredBaseTable); + Mockito.when(baseTable.getId()).thenReturn(20L); + Mockito.when(baseTable.getName()).thenReturn("base"); + Mockito.when(baseTable.getNameWithFullQualifiers()).thenReturn("internal.db.base"); Mockito.when(baseTable.needReadLockWhenPlan()).thenReturn(true); + Mockito.when(baseTable.tryReadLock(Mockito.anyLong(), Mockito.any())).thenReturn(true); Mockito.when(externalTable.getId()).thenReturn(21L); Mockito.when(externalTable.supportsExternalMetadataPreload()).thenReturn(true); Mockito.when(externalTable.getBaseSchema()).thenReturn(Collections.emptyList()); @@ -232,12 +241,17 @@ public void testPreloadRecognizesStreamBaseTablePlanReadLock() { new OriginStatement("select * from db.stream join ext", 0)); try { statementContext.getTables().put(ImmutableList.of("internal", "db", "stream"), stream); + statementContext.addImplicitTableDependency(baseTable); statementContext.registerExternalTableForPreload(externalTable, Optional.empty(), Optional.empty()); ExternalMetadataPreloadResult result = executePreload(statementContext); + statementContext.lock(); org.junit.jupiter.api.Assertions.assertTrue(result.isExecuted()); Mockito.verify(externalTable).getBaseSchema(); + Mockito.verify(stream, Mockito.never()).getBaseTableNullable(); + Mockito.verify(baseTable, Mockito.times(1)).tryReadLock(Mockito.anyLong(), Mockito.any()); + Mockito.verifyNoInteractions(recoveredBaseTable); } finally { statementContext.close(); } diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java index e6bb3b783a05c4..c997f54f4235c5 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/ExplainTableStreamPlanTest.java @@ -713,7 +713,7 @@ public void testPreLockRenameDoesNotReplaceCachedExplicitTable() throws Exceptio Assertions.assertSame(explicitTable, statementContext.getTables().get(List.of("internal", "test_stream", explicitTableName))); - Assertions.assertSame(baseTable, blockingStream.getBaseTableNullable()); + Assertions.assertTrue(statementContext.getImplicitTableDependencies().contains(baseTable)); } finally { renamesFinished.countDown(); plannerExecutor.shutdownNow(); From 6f451cc996e75bcaa75052ad9ce7570e4bc8e062 Mon Sep 17 00:00:00 2001 From: seawinde Date: Wed, 12 Aug 2026 10:34:42 +0800 Subject: [PATCH 8/8] [fix](stream) Close stream recovery compatibility gaps ### What problem does this PR solve? Issue Number: close #65418 Related PR: #66287 Problem Summary: A cached stream base table could become visible while its dropped database was only partially recovered. Existing MTMV images could also omit the stream's stable base dependency, leaving no invalidation edge or freshness snapshot after upgrade. Require the owning internal database and its stable table mapping to be published before exposing a stream base, and migrate persisted MTMV relations to include stable stream bases without fabricating historical snapshots. ### Release note Streams remain unavailable until their owning database is fully recovered, and existing MTMVs migrate stable stream base dependencies during upgrade. ### Check List (For Author) - Test: Unit Test - DropTableStreamTest and MTMVRelationTest: 12 tests passed - Java compilation and Checkstyle passed - Behavior changed: Yes. Stream availability now follows database recovery lifecycle, and old MTMV relations gain stable stream base dependencies. - Does this need documentation: No --- .../org/apache/doris/catalog/Database.java | 4 ++ .../doris/catalog/stream/BaseTableStream.java | 27 ++++++++--- .../org/apache/doris/mtmv/MTMVRelation.java | 33 +++++++++++++ .../doris/catalog/DropTableStreamTest.java | 48 +++++++++++++++++++ .../apache/doris/mtmv/MTMVRelationTest.java | 47 ++++++++++++++++++ 5 files changed, 152 insertions(+), 7 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/Database.java b/fe/fe-core/src/main/java/org/apache/doris/catalog/Database.java index e59955c515b282..e1048ce916ff6d 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/catalog/Database.java +++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/Database.java @@ -176,6 +176,10 @@ public void unmarkDropped() { isDropped = false; } + public boolean isDropped() { + return isDropped; + } + public void readLock() { this.rwLock.readLock().lock(); } 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 9ff289cd52c48b..67ef34d6c06dbb 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 @@ -18,6 +18,7 @@ package org.apache.doris.catalog.stream; import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.Database; import org.apache.doris.catalog.Env; import org.apache.doris.catalog.Table; import org.apache.doris.catalog.TableIf; @@ -117,22 +118,34 @@ public BaseTableStream(String streamName, List fullSchema, TableIf baseT public TableIf getBaseTableNullable() { TableIf cachedBaseTable = baseTable; - if (cachedBaseTable != null) { - if (cachedBaseTable instanceof Table && ((Table) cachedBaseTable).isDropped) { - baseTable = null; - return null; - } + if (isBaseTableAvailable(cachedBaseTable)) { return cachedBaseTable; } + if (cachedBaseTable != null) { + baseTable = null; + } TableIf resolvedBaseTable = baseTableInfo.getTableNullable(); - // Recovery publishes the table into database maps before clearing its dropped flag. - if (resolvedBaseTable instanceof Table && ((Table) resolvedBaseTable).isDropped) { + if (!isBaseTableAvailable(resolvedBaseTable)) { return null; } baseTable = resolvedBaseTable; return resolvedBaseTable; } + private boolean isBaseTableAvailable(TableIf candidate) { + if (candidate == null || candidate instanceof Table && ((Table) candidate).isDropped) { + return false; + } + if (!baseTableInfo.isInternalTable()) { + return true; + } + // Table recovery publishes into database maps before clearing the table flag, while database recovery clears + // member-table flags before publishing the database. Require both catalog mappings to reject either window. + Database database = Env.getCurrentInternalCatalog().getDbNullable(baseTableInfo.getDbId()); + return database != null && !database.isDropped() + && database.getTable(baseTableInfo.getTableId()).orElse(null) == candidate; + } + public void setProperties(Map properties) throws org.apache.doris.common.AnalysisException { showInitialRows = PropertyAnalyzer.analyzeBooleanProp(properties, PropertyAnalyzer.PROPERTIES_STREAM_SHOW_INITIAL_ROWS, diff --git a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelation.java b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelation.java index 452733cf7a7486..c1387c1dede462 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelation.java +++ b/fe/fe-core/src/main/java/org/apache/doris/mtmv/MTMVRelation.java @@ -17,6 +17,9 @@ package org.apache.doris.mtmv; +import org.apache.doris.catalog.TableIf; +import org.apache.doris.catalog.stream.BaseTableStream; +import org.apache.doris.common.AnalysisException; import org.apache.doris.datasource.CatalogMgr; import org.apache.doris.persist.gson.GsonPostProcessable; @@ -104,6 +107,12 @@ public void compatible(CatalogMgr catalogMgr) throws Exception { compatible(catalogMgr, baseTables); compatible(catalogMgr, baseViews); compatible(catalogMgr, baseTablesOneLevel); + addStreamBaseTables(baseTables); + if (CollectionUtils.isEmpty(baseTablesOneLevelAndFromView)) { + // Preserve the existing fallback for older images in a separate set before adding implicit stream bases. + baseTablesOneLevelAndFromView = new HashSet<>(getBaseTablesOneLevel()); + } + addStreamBaseTables(baseTablesOneLevelAndFromView); } private void compatible(CatalogMgr catalogMgr, Set infos) throws Exception { @@ -114,4 +123,28 @@ private void compatible(CatalogMgr catalogMgr, Set infos) throws baseTableInfo.compatible(catalogMgr); } } + + private void addStreamBaseTables(Set infos) throws Exception { + if (CollectionUtils.isEmpty(infos)) { + return; + } + // Older images may contain only the stream relation; add its stable base so freshness and invalidation survive + // an upgrade without inventing a historical snapshot for the newly discovered dependency. + for (BaseTableInfo info : new HashSet<>(infos)) { + TableIf table; + try { + table = MTMVUtil.getTable(info); + } catch (AnalysisException e) { + continue; + } + if (table instanceof BaseTableStream) { + TableIf baseTable = ((BaseTableStream) table).getBaseTableNullable(); + if (baseTable == null) { + throw new AnalysisException( + "Failed to resolve stream base table during MTMV compatibility: " + info); + } + infos.add(new BaseTableInfo(baseTable)); + } + } + } } diff --git a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java index 28aeee3935f6e6..b70d89adec58b7 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java @@ -220,6 +220,53 @@ public void testRecoveringBaseTableRemainsUnavailableUntilUnmarkedDropped() thro } } + @Test + public void testBaseTableRemainsUnavailableUntilDatabaseRecovers() throws Exception { + Database baseDb = Mockito.spy(new Database(Env.getCurrentEnv().getNextId(), "test_stream_recovery_db")); + Env.getCurrentInternalCatalog().unprotectCreateDb(baseDb); + createTable("create table test_stream_recovery_db.tbl_recovery_db (k1 int, k2 int) " + + "unique key(k1) distributed by hash(k1) buckets 1 " + + "properties('replication_num' = '1', 'binlog.enable' = 'true', 'binlog.format' = 'ROW', " + + "'binlog.need_historical_value' = 'true')"); + createTable("create stream test_stream.s_recovery_db on table test_stream_recovery_db.tbl_recovery_db " + + "properties('show_initial_rows' = 'true')"); + OlapTable baseTable = (OlapTable) baseDb.getTableOrMetaException("tbl_recovery_db"); + OlapTableStream stream = (OlapTableStream) Env.getCurrentInternalCatalog() + .getDbOrMetaException("test_stream").getTableOrMetaException("s_recovery_db"); + Assertions.assertSame(baseTable, stream.getBaseTableNullable()); + + dropDatabaseWithSql("drop database test_stream_recovery_db"); + CountDownLatch tablesRecovered = new CountDownLatch(1); + CountDownLatch allowDatabaseRecovery = new CountDownLatch(1); + Mockito.doAnswer(invocation -> { + boolean registered = (boolean) invocation.callRealMethod(); + tablesRecovered.countDown(); + Assertions.assertTrue(allowDatabaseRecovery.await(10, TimeUnit.SECONDS)); + return registered; + }).when(baseDb).registerTable(Mockito.any()); + + ExecutorService executor = Executors.newSingleThreadExecutor(); + try { + Future recoverFuture = executor.submit(() -> { + Env.getCurrentInternalCatalog().recoverDatabase("test_stream_recovery_db", -1, ""); + return null; + }); + Assertions.assertTrue(tablesRecovered.await(10, TimeUnit.SECONDS)); + + Assertions.assertFalse(baseTable.isDropped); + Assertions.assertTrue(baseDb.isDropped()); + Assertions.assertNull(Env.getCurrentInternalCatalog().getDbNullable(baseDb.getId())); + Assertions.assertNull(stream.getBaseTableNullable()); + + allowDatabaseRecovery.countDown(); + recoverFuture.get(10, TimeUnit.SECONDS); + Assertions.assertSame(baseTable, stream.getBaseTableNullable()); + } finally { + allowDatabaseRecovery.countDown(); + executor.shutdownNow(); + } + } + @Test public void testBaseTableQualifiersFollowRenameAndRecovery() throws Exception { createBaseTableAndStream("tbl_rename", "s_rename"); @@ -334,6 +381,7 @@ public void testForceDropAndSameNameTableDoNotRestoreStream() throws Exception { @Override protected void runAfterAll() throws Exception { + dropDatabase("test_stream_recovery_db"); dropDatabase("test_stream"); } } diff --git a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelationTest.java b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelationTest.java index 43f1f75e3c1313..ab9dc1fd112861 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelationTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/mtmv/MTMVRelationTest.java @@ -23,6 +23,7 @@ import org.apache.doris.catalog.TableIf; import org.apache.doris.common.Config; import org.apache.doris.common.Pair; +import org.apache.doris.persist.gson.GsonUtils; import org.apache.doris.utframe.TestWithFeService; import com.google.common.collect.Sets; @@ -164,4 +165,50 @@ public void testMTMVRelationIncludesStreamBaseDependency() throws Exception { Assertions.assertEquals(Sets.newHashSet(stream, baseTable), queryTables.first); Assertions.assertEquals(Sets.newHashSet(stream), queryTables.second); } + + @Test + public void testCompatibleAddsPersistedStreamBaseDependency() throws Exception { + createDatabaseAndUse("stream_mtmv_compatible_db"); + createTables( + "CREATE TABLE stream_base (k1 int, k2 int)\n" + + "UNIQUE KEY(k1)\n" + + "DISTRIBUTED BY HASH(k1) BUCKETS 1\n" + + "PROPERTIES ('replication_num' = '1', 'binlog.enable' = 'true',\n" + + "'binlog.format' = 'ROW', 'binlog.need_historical_value' = 'true')", + "CREATE STREAM stream_source ON TABLE stream_base\n" + + "PROPERTIES ('show_initial_rows' = 'true')"); + createMvByNereids("CREATE MATERIALIZED VIEW stream_mv BUILD DEFERRED\n" + + "REFRESH COMPLETE ON MANUAL\n" + + "DISTRIBUTED BY RANDOM BUCKETS 1\n" + + "PROPERTIES ('replication_num' = '1')\n" + + "AS SELECT k1, k2 FROM stream_source"); + + Database db = Env.getCurrentEnv().getInternalCatalog() + .getDbOrAnalysisException("stream_mtmv_compatible_db"); + MTMV mtmv = (MTMV) db.getTableOrAnalysisException("stream_mv"); + TableIf stream = db.getTableOrAnalysisException("stream_source"); + TableIf baseTable = db.getTableOrAnalysisException("stream_base"); + BaseTableInfo streamInfo = new BaseTableInfo(stream); + BaseTableInfo baseTableInfo = new BaseTableInfo(baseTable); + BaseTableInfo mtmvInfo = new BaseTableInfo(mtmv); + + // Model an image written before stream bases were persisted as MTMV dependencies. + MTMVRelation oldRelation = new MTMVRelation(Sets.newHashSet(streamInfo), Sets.newHashSet(streamInfo), null, + Sets.newHashSet(), Sets.newHashSet()); + mtmv.setRelation(GsonUtils.GSON.fromJson(GsonUtils.GSON.toJson(oldRelation), MTMVRelation.class)); + MTMVRelationManager relationManager = Env.getCurrentEnv().getMtmvService().getRelationManager(); + relationManager.refreshMTMVCache(mtmv.getRelation(), mtmvInfo); + Assertions.assertEquals(Sets.newHashSet(), relationManager.getMtmvsByBaseTable(baseTableInfo)); + Assertions.assertTrue(MTMVPartitionUtil.isMTMVSync(mtmv)); + + mtmv.compatible(Env.getCurrentEnv().getCatalogMgr()); + mtmv.compatible(Env.getCurrentEnv().getCatalogMgr()); + + Assertions.assertEquals(Sets.newHashSet(streamInfo, baseTableInfo), mtmv.getRelation().getBaseTables()); + Assertions.assertEquals(Sets.newHashSet(streamInfo), mtmv.getRelation().getBaseTablesOneLevel()); + Assertions.assertEquals(Sets.newHashSet(streamInfo, baseTableInfo), + mtmv.getRelation().getBaseTablesOneLevelAndFromView()); + Assertions.assertEquals(Sets.newHashSet(mtmvInfo), relationManager.getMtmvsByBaseTable(baseTableInfo)); + Assertions.assertFalse(MTMVPartitionUtil.isMTMVSync(mtmv)); + } }