Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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);
Expand Down Expand Up @@ -166,7 +169,7 @@ public List<Cloud.TableStreamIdentityPB> getCloudTableStreamsForBaseTable(
public void cleanupStalePartitionOffsets() {
List<Long> staleDbIds = new ArrayList<>();
List<Pair<Long, Long>> staleStreamIds = new ArrayList<>();
List<TableStreamCleanupInfo.PartitionOffsetPruneEntry> pruneEntries = new ArrayList<>();
List<EditLogItem> editLogItems = new ArrayList<>();
for (Map.Entry<Long, Set<Long>> entry : copyDbStreamMap().entrySet()) {
Optional<Database> db = Env.getCurrentInternalCatalog().getDb(entry.getKey());
if (!db.isPresent()) {
Expand All @@ -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<TableStreamCleanupInfo.PartitionOffsetPruneEntry> cleanupStalePartitionOffsets(
OlapTableStream stream) {
private Optional<EditLogItem> 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());
Expand All @@ -215,51 +218,54 @@ private Optional<TableStreamCleanupInfo.PartitionOffsetPruneEntry> 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",
stream.getName(), baseTable.getName());
}
return Optional.empty();
}
Set<Long> validPartitionIds;
Set<Long> 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<Long> validPartitionIds = new HashSet<>(baseTable.getPartitionIds());
while (DebugPointUtil.getDebugParamOrDefault(
"TableStreamManager.cleanupStalePartitionOffsets.blockAfterPartitionSnapshot", -1L)
== stream.getId()) {
LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(10));
}
return Optional.empty();
}
Set<Long> 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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Loading