github-actions[bot] commented on code in PR #68196:
URL: https://github.com/apache/doris/pull/68196#discussion_r4134257951
##########
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:
[P3] Treat an invalidated row-count load as an expected miss. When a table
or DB refresh marks an in-flight LoadFence, a successful loader now completes
this future with null. The synchronous getCachedRowCount path then dereferences
it at line 307, logs an Unexpected exception while returning row count WARN
with an NPE stack, and returns UNKNOWN; repeated partition events can fill FE
logs during normal refresh/read overlap. Handle this null as an expected
invalidation result without logging an exception, while keeping it
non-cacheable.
##########
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:
[P2] Bound database invalidation to that database. A warm REFRESH DATABASE
calls this fence before and after resetting its table objects. Each call holds
the global publicationLock write lock while removeIf traverses up to 100,000
row-count entries across every catalog (plus all in-flight loads), blocking
unrelated planners that need the read lock. The catalog fence in this PR is
O(1), but this ordinary DB route still makes two full scans. Use a per-DB
generation or scoped index so a DB refresh does not serialize unrelated catalog
reads behind the full cache.
--
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]