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..f768f5afeb9bf9 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 @@ -32,7 +32,9 @@ import org.apache.doris.common.io.Text; import org.apache.doris.common.io.Writable; import org.apache.doris.common.lock.MonitoredReentrantReadWriteLock; +import org.apache.doris.common.util.DebugPointUtil; import org.apache.doris.common.util.MasterDaemon; +import org.apache.doris.persist.EditLog.EditLogItem; import org.apache.doris.persist.TableStreamCleanupInfo; import org.apache.doris.persist.gson.GsonPostProcessable; import org.apache.doris.persist.gson.GsonUtils; @@ -57,6 +59,7 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.LockSupport; public class TableStreamManager extends MasterDaemon implements Writable, GsonPostProcessable { private static final Logger LOG = LogManager.getLogger(TableStreamManager.class); @@ -166,7 +169,7 @@ public List getCloudTableStreamsForBaseTable( public void cleanupStalePartitionOffsets() { List staleDbIds = new ArrayList<>(); List> staleStreamIds = new ArrayList<>(); - List pruneEntries = new ArrayList<>(); + List editLogItems = new ArrayList<>(); for (Map.Entry> entry : copyDbStreamMap().entrySet()) { Optional db = Env.getCurrentInternalCatalog().getDb(entry.getKey()); if (!db.isPresent()) { @@ -183,18 +186,18 @@ public void cleanupStalePartitionOffsets() { staleStreamIds.add(Pair.of(db.get().getId(), tableId)); continue; } - cleanupStalePartitionOffsets((OlapTableStream) table.get()).ifPresent(pruneEntries::add); + cleanupStalePartitionOffsets((OlapTableStream) table.get()).ifPresent(editLogItems::add); } } removeStaleDbAndStream(staleDbIds, staleStreamIds); - if (!pruneEntries.isEmpty() || !staleDbIds.isEmpty() || !staleStreamIds.isEmpty()) { - Env.getCurrentEnv().getEditLog().logTableStreamCleanup( - new TableStreamCleanupInfo(pruneEntries, staleDbIds, staleStreamIds)); + if (!staleDbIds.isEmpty() || !staleStreamIds.isEmpty()) { + editLogItems.add(Env.getCurrentEnv().getEditLog().logTableStreamCleanup( + new TableStreamCleanupInfo(Collections.emptyList(), staleDbIds, staleStreamIds))); } + editLogItems.forEach(EditLogItem::await); } - private Optional cleanupStalePartitionOffsets( - OlapTableStream stream) { + private Optional cleanupStalePartitionOffsets(OlapTableStream stream) { if (!stream.tryReadLock(Table.TRY_LOCK_TIMEOUT_MS, TimeUnit.MILLISECONDS)) { if (LOG.isDebugEnabled()) { LOG.debug("skip cleaning stream {} because stream read lock is busy", stream.getName()); @@ -215,7 +218,6 @@ private Optional cleanupStaleP stream.readUnlock(); } // stream read lock is released - // base table read lock is held if (!baseTable.tryReadLock(Table.TRY_LOCK_TIMEOUT_MS, TimeUnit.MILLISECONDS)) { if (LOG.isDebugEnabled()) { LOG.debug("skip cleaning stream {} because base table {} read lock is busy", @@ -223,43 +225,47 @@ private Optional cleanupStaleP } return Optional.empty(); } - Set validPartitionIds; + Set stalePartitionIds; + EditLogItem editLogItem; try { if (baseTable.isDropped) { return Optional.empty(); } - validPartitionIds = new HashSet<>(baseTable.getPartitionIds()); - } finally { - baseTable.readUnlock(); - } - // base table read lock is released - // stream write lock is held - if (!stream.tryWriteLock(Table.TRY_LOCK_TIMEOUT_MS, TimeUnit.MILLISECONDS)) { - if (LOG.isDebugEnabled()) { - LOG.debug("skip cleaning stream {} because stream write lock is busy", stream.getName()); + Set validPartitionIds = new HashSet<>(baseTable.getPartitionIds()); + while (DebugPointUtil.getDebugParamOrDefault( + "TableStreamManager.cleanupStalePartitionOffsets.blockAfterPartitionSnapshot", -1L) + == stream.getId()) { + LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(10)); } - return Optional.empty(); - } - Set stalePartitionIds; - try { - if (stream.isDisabled() || stream.isStale()) { + if (!stream.tryWriteLockIfExist(Table.TRY_LOCK_TIMEOUT_MS, TimeUnit.MILLISECONDS)) { + if (LOG.isDebugEnabled()) { + LOG.debug("skip cleaning stream {} because it is busy or dropped", stream.getName()); + } return Optional.empty(); } - stalePartitionIds = stream.unprotectedCollectStalePartitionOffsetIds(validPartitionIds); - if (stalePartitionIds.isEmpty()) { - return Optional.empty(); + try { + if (stream.isDisabled() || stream.isStale()) { + return Optional.empty(); + } + stalePartitionIds = stream.unprotectedCollectStalePartitionOffsetIds(validPartitionIds); + if (stalePartitionIds.isEmpty()) { + return Optional.empty(); + } + stream.unprotectedPrunePartitionOffsets(stalePartitionIds); + TableStreamCleanupInfo.PartitionOffsetPruneEntry pruneEntry = + new TableStreamCleanupInfo.PartitionOffsetPruneEntry( + stream.getDatabase().getId(), stream.getId(), stalePartitionIds); + editLogItem = Env.getCurrentEnv().getEditLog().logTableStreamCleanup( + new TableStreamCleanupInfo(Collections.singletonList(pruneEntry))); + } finally { + stream.writeUnlock(); } - stream.unprotectedPrunePartitionOffsets(stalePartitionIds); } finally { - stream.writeUnlock(); - } - // stream write lock is released - if (stalePartitionIds.size() > 0) { - LOG.info("cleaned {} stale partition offset entries from stream {}.{} ({})", - stalePartitionIds.size(), stream.getDatabase().getFullName(), stream.getName(), stream.getId()); + baseTable.readUnlock(); } - return Optional.of(new TableStreamCleanupInfo.PartitionOffsetPruneEntry( - stream.getDatabase().getId(), stream.getId(), stalePartitionIds)); + LOG.info("cleaned {} stale partition offset entries from stream {}.{} ({})", + stalePartitionIds.size(), stream.getDatabase().getFullName(), stream.getName(), stream.getId()); + return Optional.of(editLogItem); } public void replayTableStreamCleanup(TableStreamCleanupInfo info) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/persist/EditLog.java b/fe/fe-core/src/main/java/org/apache/doris/persist/EditLog.java index 1c7c9738415d4b..950be715a607ca 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/persist/EditLog.java +++ b/fe/fe-core/src/main/java/org/apache/doris/persist/EditLog.java @@ -2308,8 +2308,8 @@ public void logDynamicPartition(ModifyTablePropertyOperationLog info) { logModifyTableProperty(OperationType.OP_DYNAMIC_PARTITION, info); } - public void logTableStreamCleanup(TableStreamCleanupInfo info) { - logEdit(OperationType.OP_TABLE_STREAM_CLEANUP, info); + public EditLogItem logTableStreamCleanup(TableStreamCleanupInfo info) { + return submitEdit(OperationType.OP_TABLE_STREAM_CLEANUP, info); } public long logModifyReplicationNum(ModifyTablePropertyOperationLog info) { diff --git a/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/TableStreamManagerCleanupTest.java b/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/TableStreamManagerCleanupTest.java index 0d6b3350582300..45a8ccb725fc64 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/TableStreamManagerCleanupTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/catalog/stream/TableStreamManagerCleanupTest.java @@ -25,16 +25,31 @@ import org.apache.doris.common.FeConstants; import org.apache.doris.common.Pair; import org.apache.doris.common.jmockit.Deencapsulation; +import org.apache.doris.common.lock.MonitoredReentrantReadWriteLock; +import org.apache.doris.common.util.DebugPointUtil; +import org.apache.doris.common.util.DebugPointUtil.DebugPoint; +import org.apache.doris.persist.DropInfo; +import org.apache.doris.persist.EditLog; +import org.apache.doris.persist.RecoverInfo; import org.apache.doris.persist.TableStreamCleanupInfo; import org.apache.doris.utframe.TestWithFeService; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import org.mockito.Mockito; +import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; 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; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.BooleanSupplier; public class TableStreamManagerCleanupTest extends TestWithFeService { @@ -60,6 +75,201 @@ public void testCleanupRemovedPartitionOffsets() throws Exception { assertPartitionState(context.stream, keptPartitionId, removedPartitionId, true); } + @Test + public void testCleanupRetainsOffsetForPartitionAddedAfterSnapshot() throws Exception { + StreamContext context = createStreamContext("cleanup_partition_race"); + String debugPointName = + "TableStreamManager.cleanupStalePartitionOffsets.blockAfterPartitionSnapshot"; + DebugPoint debugPoint = new DebugPoint(); + debugPoint.executeLimit = Integer.MAX_VALUE; + debugPoint.params.put("value", String.valueOf(context.stream.getId())); + boolean debugPointsEnabled = Config.enable_debug_points; + Config.enable_debug_points = true; + DebugPointUtil.addDebugPoint(debugPointName, debugPoint); + + MonitoredReentrantReadWriteLock baseTableLock = Deencapsulation.getField(context.baseTable, "rwLock"); + AtomicReference addPartitionThread = new AtomicReference<>(); + ExecutorService executor = Executors.newFixedThreadPool(2); + Future cleanup = executor.submit( + () -> Env.getCurrentEnv().getTableStreamManager().cleanupStalePartitionOffsets()); + Future addPartition = null; + try { + try { + await(() -> debugPoint.executeNum.get() > 0); + addPartition = executor.submit(() -> { + connectContext.setThreadLocalInfo(); + addPartitionThread.set(Thread.currentThread()); + alterTableSync("alter table test_stream_cleanup." + context.baseTable.getName() + + " add partition p3 values less than (\"300\")"); + long partitionId = context.baseTable.getPartition("p3").getId(); + updatePartitionOffset(context.stream, partitionId, 33L, 333L); + return null; + }); + Future addPartitionResult = addPartition; + await(() -> addPartitionResult.isDone() + || addPartitionThread.get() != null + && baseTableLock.hasQueuedThread(addPartitionThread.get())); + } finally { + DebugPointUtil.removeDebugPoint(debugPointName); + Config.enable_debug_points = debugPointsEnabled; + } + + cleanup.get(10, TimeUnit.SECONDS); + addPartition.get(10, TimeUnit.SECONDS); + } finally { + executor.shutdownNow(); + Assertions.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS)); + } + + long partitionId = context.baseTable.getPartition("p3").getId(); + Assertions.assertTrue(context.stream.hasConsumedData(partitionId)); + } + + @Test + public void testCleanupJournalOrderMatchesLeaderState() throws Exception { + StreamContext context = createStreamContext("cleanup_journal_race"); + long keptPartitionId = context.baseTable.getPartition("p1").getId(); + long removedPartitionId = context.baseTable.getPartition("p2").getId(); + setPartitionState(context.stream, keptPartitionId, removedPartitionId); + alterTableSync("alter table test_stream_cleanup." + context.baseTable.getName() + " drop partition p2"); + + List journalOrder = Collections.synchronizedList(new ArrayList<>()); + CountDownLatch allowUpdate = new CountDownLatch(1); + CountDownLatch updateDone = new CountDownLatch(1); + OlapTableStreamUpdate streamUpdate = new OlapTableStreamUpdate( + Collections.emptyMap(), Collections.singletonMap(removedPartitionId, 44L)); + ExecutorService executor = Executors.newSingleThreadExecutor(); + Future update = executor.submit(() -> { + Assertions.assertTrue(allowUpdate.await(10, TimeUnit.SECONDS)); + context.stream.writeLock(); + try { + journalOrder.add(streamUpdate); + context.stream.unprotectedUpdateStreamUpdate(streamUpdate, 444L); + } finally { + context.stream.writeUnlock(); + updateDone.countDown(); + } + return null; + }); + + EditLog editLog = Env.getCurrentEnv().getEditLog(); + EditLog spyEditLog = Mockito.spy(editLog); + Mockito.doAnswer(invocation -> { + TableStreamCleanupInfo cleanupInfo = invocation.getArgument(0); + allowUpdate.countDown(); + if (!context.stream.isWriteLockHeldByCurrentThread()) { + Assertions.assertTrue(updateDone.await(10, TimeUnit.SECONDS)); + } + journalOrder.add(cleanupInfo); + return Mockito.mock(EditLog.EditLogItem.class); + }).when(spyEditLog).logTableStreamCleanup(Mockito.any(TableStreamCleanupInfo.class)); + Env.getCurrentEnv().setEditLog(spyEditLog); + try { + Env.getCurrentEnv().getTableStreamManager().cleanupStalePartitionOffsets(); + update.get(10, TimeUnit.SECONDS); + } finally { + Env.getCurrentEnv().setEditLog(editLog); + executor.shutdownNow(); + Assertions.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS)); + } + + Map leaderOffsets = new HashMap<>( + Deencapsulation.getField(context.stream, "partitionOffset")); + setPartitionState(context.stream, keptPartitionId, removedPartitionId); + for (Object journal : journalOrder) { + if (journal instanceof OlapTableStreamUpdate) { + updatePartitionOffset(context.stream, removedPartitionId, 44L, 444L); + } else { + Env.getCurrentEnv().getTableStreamManager() + .replayTableStreamCleanup((TableStreamCleanupInfo) journal); + } + } + + Map replayedOffsets = Deencapsulation.getField(context.stream, "partitionOffset"); + Assertions.assertEquals(leaderOffsets, replayedOffsets); + } + + @Test + public void testDropStreamDuringCleanupReplaysDeterministically() throws Exception { + StreamContext context = createStreamContext("cleanup_drop_race"); + long keptPartitionId = context.baseTable.getPartition("p1").getId(); + long removedPartitionId = context.baseTable.getPartition("p2").getId(); + setPartitionState(context.stream, keptPartitionId, removedPartitionId); + alterTableSync("alter table test_stream_cleanup." + context.baseTable.getName() + " drop partition p2"); + + String debugPointName = + "TableStreamManager.cleanupStalePartitionOffsets.blockAfterPartitionSnapshot"; + DebugPoint debugPoint = new DebugPoint(); + debugPoint.executeLimit = Integer.MAX_VALUE; + debugPoint.params.put("value", String.valueOf(context.stream.getId())); + boolean debugPointsEnabled = Config.enable_debug_points; + Config.enable_debug_points = true; + DebugPointUtil.addDebugPoint(debugPointName, debugPoint); + + List journalOrder = Collections.synchronizedList(new ArrayList<>()); + EditLog editLog = Env.getCurrentEnv().getEditLog(); + EditLog spyEditLog = Mockito.spy(editLog); + Mockito.doAnswer(invocation -> { + journalOrder.add(invocation.getArgument(0)); + return null; + }).when(spyEditLog).logDropTable(Mockito.any(DropInfo.class)); + Mockito.doAnswer(invocation -> { + journalOrder.add(invocation.getArgument(0)); + return Mockito.mock(EditLog.EditLogItem.class); + }).when(spyEditLog).logTableStreamCleanup(Mockito.any(TableStreamCleanupInfo.class)); + Mockito.doAnswer(invocation -> { + journalOrder.add(invocation.getArgument(0)); + return null; + }).when(spyEditLog).logRecoverTable(Mockito.any(RecoverInfo.class)); + Env.getCurrentEnv().setEditLog(spyEditLog); + + ExecutorService executor = Executors.newSingleThreadExecutor(); + Future cleanup = executor.submit( + () -> Env.getCurrentEnv().getTableStreamManager().cleanupStalePartitionOffsets()); + try { + try { + await(() -> debugPoint.executeNum.get() > 0); + Env.getCurrentInternalCatalog().dropTable( + "test_stream_cleanup", context.stream.getName(), false, false, true, + false, false, false); + } finally { + DebugPointUtil.removeDebugPoint(debugPointName); + Config.enable_debug_points = debugPointsEnabled; + } + cleanup.get(10, TimeUnit.SECONDS); + Env.getCurrentEnv().recoverTable( + "test_stream_cleanup", context.stream.getName(), "", -1L); + } finally { + executor.shutdownNow(); + try { + Assertions.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS)); + } finally { + Env.getCurrentEnv().setEditLog(editLog); + } + } + + Map leaderOffsets = new HashMap<>( + Deencapsulation.getField(context.stream, "partitionOffset")); + setPartitionState(context.stream, keptPartitionId, removedPartitionId); + Database db = (Database) Env.getCurrentInternalCatalog().getDbOrMetaException("test_stream_cleanup"); + for (Object journal : journalOrder) { + if (journal instanceof DropInfo) { + DropInfo dropInfo = (DropInfo) journal; + Env.getCurrentEnv().replayDropTable( + db, dropInfo.getTableId(), dropInfo.isForceDrop(), dropInfo.getRecycleTime()); + } else if (journal instanceof TableStreamCleanupInfo) { + Env.getCurrentEnv().getTableStreamManager() + .replayTableStreamCleanup((TableStreamCleanupInfo) journal); + } else { + Env.getCurrentEnv().replayRecoverTable((RecoverInfo) journal); + } + } + + Map replayedOffsets = Deencapsulation.getField(context.stream, "partitionOffset"); + Assertions.assertEquals(leaderOffsets, replayedOffsets); + Assertions.assertTrue(replayedOffsets.containsKey(removedPartitionId)); + } + @Test public void testCleanupSkipsDisabledStream() throws Exception { StreamContext context = createStreamContext("cleanup_disabled"); @@ -207,6 +417,25 @@ private void assertPartitionState(OlapTableStream stream, long keptPartitionId, Assertions.assertEquals(!removedExpected, historicalPartitionTSO.containsKey(removedPartitionId)); } + private static void await(BooleanSupplier condition) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + while (!condition.getAsBoolean() && System.nanoTime() < deadline) { + Thread.sleep(10); + } + Assertions.assertTrue(condition.getAsBoolean()); + } + + private static void updatePartitionOffset( + OlapTableStream stream, long partitionId, long offset, long commitTimeMs) { + stream.writeLock(); + try { + stream.unprotectedUpdateStreamUpdate(new OlapTableStreamUpdate( + Collections.emptyMap(), Collections.singletonMap(partitionId, offset)), commitTimeMs); + } finally { + stream.writeUnlock(); + } + } + private static class StreamContext { private final OlapTable baseTable; private final OlapTableStream stream;