yujun777 commented on code in PR #66287:
URL: https://github.com/apache/doris/pull/66287#discussion_r3818354645


##########
fe/fe-core/src/main/java/org/apache/doris/catalog/stream/BaseTableStream.java:
##########
@@ -139,23 +166,35 @@ public StreamScanType getStreamScanType() {
     }
 
     public boolean isDisabled() {
-        return disabled;
+        return isDisabled(getBaseTableNullable());
+    }
+
+    boolean isDisabled(TableIf availableBaseTable) {
+        return disabled || availableBaseTable == null;
     }
 
     public void setDisabled(boolean disabled) {
         this.disabled = disabled;
     }
 
     public boolean isStale() {
-        return stale;
+        return isStale(getBaseTableNullable());
+    }
+
+    boolean isStale(TableIf availableBaseTable) {
+        return stale || availableBaseTable == null;
     }
 
     public void setStale(boolean stale) {
         this.stale = stale;
     }
 
     public String getStaleReason() {
-        return staleReason;
+        return getStaleReason(getBaseTableNullable());
+    }
+
+    String getStaleReason(TableIf availableBaseTable) {
+        return availableBaseTable == null ? BASE_TABLE_NOT_FOUND_STALE_REASON 
: staleReason;

Review Comment:
   Minor: when the base table is unavailable this always returns 
`BASE_TABLE_NOT_FOUND_STALE_REASON`, masking an explicitly set `staleReason` 
(e.g. from `setStaleReason` / a stale state set by other means). If a stream is 
both explicitly marked stale with a custom reason and its base is dropped, the 
actionable reason is hidden. Consider returning the explicit `staleReason` when 
it is set and only falling back to `Base table does not exist` when it is unset.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/InternalCatalog.java:
##########
@@ -1012,7 +1012,8 @@ private void dropTableInternal(Database db, Table table, 
boolean isView, boolean
         } finally {
             table.writeUnlock();
         }
-        if (table instanceof OlapTable) {
+        // MTMVs persist streams as explicit relations, so dropping a stream 
must invalidate them immediately.
+        if (table instanceof OlapTable || table instanceof BaseTableStream) {

Review Comment:
   P1 (blocking): this invalidation hook only runs on the direct `DROP TABLE` 
live path. A stream dropped through other paths never invalidates its dependent 
MTMVs:
   
   - **Database drop**: `dropDb`/`replayDropDb` -> `unprotectDropDb` (line 
~557) iterates tables via `unprotectDropTable`, which for a `BaseTableStream` 
only calls `removeTableStream` (lines 1047-1048) and never 
`mtmvService.dropTable`. An MTMV over a stream inside a dropped database stays 
rewrite-eligible and can serve results produced from the old stream's base 
after the database is replaced.
   - **Journal replay**: `replayDropTable` -> `dropTable` -> 
`unprotectDropTable` (lines 1100-1116) also bypasses `dropTableInternal`, so 
after failover the replayed stream drop does not mark dependent MTMVs 
`SCHEMA_CHANGE` (that status is not journaled).
   
   Since the comment states the invariant "dropping a stream must invalidate 
them immediately", the hook should live in `unprotectDropTable` (the shared 
path for direct drop + DB drop + both replay paths), not only in 
`dropTableInternal`. Please also add a test for the database-cascade drop (and 
ideally replay) invalidation of an MTMV over a stream.



##########
fe/fe-core/src/test/java/org/apache/doris/catalog/DropTableStreamTest.java:
##########
@@ -104,8 +141,247 @@ 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(getStreamDdl(stream).contains("ON TABLE 
internal.test_stream.tbl_recover"));
+        
Assertions.assertTrue(Env.getCurrentEnv().getTableStreamManager().getTableStreamIds(db)
+                .contains(stream.getId()));
+
+        TRow streamRow = getStreamMetadataRow("s_recover");
+        List<String> 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 
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<Boolean> 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 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");
+        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 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");
+        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_recovery_db");

Review Comment:
   Minor (test hygiene): `testBaseTableQualifiersFollowDatabaseRenameAfterDrop` 
renames `test_stream_db_rename` -> `test_stream_db_renamed` (line 327) and 
never drops it, but `runAfterAll` only cleans up `test_stream_recovery_db` and 
`test_stream`. This leaves `test_stream_db_renamed` behind for the rest of the 
class. Add `dropDatabase("test_stream_db_renamed")` here for a clean class 
teardown.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/rules/analysis/CollectRelation.java:
##########
@@ -313,10 +312,13 @@ protected void parseAndCollectFromView(List<String> 
tableQualifier, View view, C
         parentContext.addPlanProcesses(viewContext.getPlanProcesses());
     }
 
-    private void collectFromTableStream(BaseTableStream tableStream, 
CascadesContext cascadesContext,
-                                        TableFrom tableFrom, 
Optional<UnboundRelation> unboundRelation) {
-        StatementContext statementContext = 
cascadesContext.getConnectContext().getStatementContext();
-        List<String> tableQualifier = tableStream.getBaseTableFullQualifiers();
-        statementContext.getAndCacheTable(tableQualifier, tableFrom, 
unboundRelation);
+    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 ["

Review Comment:
   Minor: this new path reports a missing base table as `Table [<name>] does 
not exist`, while `getBaseTableOrNereidsAnalysisException` (used elsewhere, 
e.g. `BaseTableStream.getBaseTableOrException`) reports `Unknown base table 
'<name>'`. A user sees one or the other depending on which code path runs 
first. Consider reusing a single message (e.g. delegate to 
`getBaseTableOrNereidsAnalysisException`) so the error is consistent across the 
planner collection and binding paths.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to