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]