924060929 commented on code in PR #68196:
URL: https://github.com/apache/doris/pull/68196#discussion_r4139872920
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/ExternalRowCountCache.java:
##########
@@ -87,6 +122,136 @@ protected Optional<Long> doLoad(RowCountKey rowCountKey) {
}
}
+ private final class InvalidationAwareLoader implements
AsyncCacheLoader<RowCountKey, Optional<Long>> {
+ private final RowCountCacheLoader delegate;
+
+ private InvalidationAwareLoader(RowCountCacheLoader delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public CompletableFuture<Optional<Long>> asyncLoad(RowCountKey key,
Executor executor) {
+ return loadWithInvalidationFence(key, executor, () ->
delegate.doLoad(key));
+ }
+ }
+
+ private static final class LoadFence {
+ private boolean invalidated;
+ }
+
+ // A cache generation keeps the existing table-ID identity within one
catalog incarnation.
+ // In-flight loads need the complete scope so database/table invalidation
can fence a
+ // same-tableId replacement without touching another catalog's load.
+ private static final class LoadKey {
+ private final long catalogId;
+ private final long dbId;
+ private final long tableId;
+
+ private LoadKey(RowCountKey key) {
+ catalogId = key.catalogId;
+ dbId = key.dbId;
+ tableId = key.tableId;
+ }
+
+ @Override
+ public boolean equals(Object obj) {
+ if (this == obj) {
+ return true;
+ }
+ if (!(obj instanceof LoadKey)) {
+ return false;
+ }
+ LoadKey other = (LoadKey) obj;
+ return catalogId == other.catalogId && dbId == other.dbId &&
tableId == other.tableId;
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(catalogId, dbId, tableId);
+ }
+ }
+
+ private CompletableFuture<Optional<Long>> loadWithInvalidationFence(
+ RowCountKey key, Executor executor, Supplier<Optional<Long>>
loader) {
+ LoadFence fence = new LoadFence();
+ LoadKey loadKey = new LoadKey(key);
+ publicationLock.readLock().lock();
+ try {
+ inFlightLoads.compute(loadKey, (ignored, fences) -> {
+ Set<LoadFence> currentFences = fences == null ?
ConcurrentHashMap.newKeySet() : fences;
+ currentFences.add(fence);
+ return currentFences;
+ });
+ } finally {
+ publicationLock.readLock().unlock();
+ }
+
+ CompletableFuture<Optional<Long>> publishedFuture = new
CompletableFuture<>();
+ CompletableFuture<Optional<Long>> loadFuture;
+ try {
+ loadFuture = CompletableFuture.supplyAsync(loader, executor);
+ } catch (RuntimeException e) {
+ publicationLock.readLock().lock();
+ try {
+ removeInFlightLoad(loadKey, fence);
+ } finally {
+ publicationLock.readLock().unlock();
+ }
+ throw e;
+ }
+ loadFuture.whenComplete((value, throwable) -> {
+ publicationLock.readLock().lock();
+ try {
+ if (throwable != null) {
+ publishedFuture.completeExceptionally(throwable);
+ } else if (fence.invalidated) {
+ publishedFuture.complete(null);
Review Comment:
Fixed in d96208503c0. An invalidated/null load now returns UNKNOWN without
dereferencing Optional or emitting an unexpected NPE warning; it remains
non-cacheable so the next read retries. The latching row-count test covers
invalidation while a load is running.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/ExternalRowCountCache.java:
##########
@@ -162,4 +351,60 @@ public long getCachedRowCountIfPresent(long catalogId,
long dbId, long tableId)
return -1;
}
+ // Catalog invalidation is a constant-time generation fence. Old entries
remain bounded by
+ // Caffeine and cannot be returned or reused after the generation changes.
+ void invalidateCatalog(long catalogId) {
+ publicationLock.writeLock().lock();
+ try {
+ AtomicLong current = catalogGenerations.get(catalogId);
+ if (current != null) {
+ current.set(nextCatalogGeneration.incrementAndGet());
+ }
+ } finally {
+ publicationLock.writeLock().unlock();
+ }
+ }
+
+ /** Catalog IDs are never reused after DROP; old readers still fail the
generation check. */
+ void releaseCatalog(long catalogId) {
+ publicationLock.writeLock().lock();
+ try {
+ catalogGenerations.remove(catalogId);
+ } finally {
+ publicationLock.writeLock().unlock();
+ }
+ }
+
+ void invalidateDb(long catalogId, long dbId) {
+ publicationLock.writeLock().lock();
+ try {
+ inFlightLoads.forEach((key, fences) -> {
+ if (key.catalogId == catalogId && key.dbId == dbId) {
+ fences.forEach(fence -> fence.invalidated = true);
+ }
+ });
+ rowCountCache.asMap().keySet().removeIf(key -> key.catalogId ==
catalogId && key.dbId == dbId);
Review Comment:
Fixed in d96208503c0. Database invalidation now advances a per-database
generation in O(1) under the publication lock instead of scanning the global
cache and in-flight loads. The generation registry is Caffeine-bounded, and
globally unique epochs prevent old entries from becoming reachable after
eviction. A latching test checks that an old DB load cannot return stale data
and a hot sibling DB remains readable.
--
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]