FrankChen021 commented on code in PR #20403:
URL: https://github.com/apache/druid/pull/20403#discussion_r4165693068
##########
processing/src/main/java/org/apache/druid/java/util/http/client/pool/ResourcePool.java:
##########
@@ -75,6 +89,44 @@ public PooledResources<V> load(K input)
}
}
);
+ final long unusedTimeoutMillis = config.getUnusedConnectionTimeoutMillis()
> 0
+ ?
config.getUnusedConnectionTimeoutMillis()
+ :
DEFAULT_UNUSED_CONNECTION_TIMEOUT_MILLIS;
+ this.sweepIntervalNanos =
TimeUnit.MILLISECONDS.toNanos(unusedTimeoutMillis);
+ this.abandonmentTimeoutNanos = sweepIntervalNanos > Long.MAX_VALUE /
KEY_ABANDONMENT_TIMEOUT_MULTIPLIER
+ ? Long.MAX_VALUE
+ : sweepIntervalNanos *
KEY_ABANDONMENT_TIMEOUT_MULTIPLIER;
+ this.nextSweepNanos.set(System.nanoTime() + sweepIntervalNanos);
+ }
+
+ /**
+ * Evicts and closes every key that has no resources on loan and has seen no
activity for
+ * {@link #abandonmentTimeoutNanos}, at most once per {@link
#sweepIntervalNanos}. Closing a holder is safe even
+ * if a taker races with it: {@link PooledResources#get()} then returns null
and {@link #take} retries, and
+ * anything lent out is closed when it is given back.
+ */
+ private void reapAbandonedKeys()
+ {
+ final long now = System.nanoTime();
+ final long next = nextSweepNanos.get();
+ // Compared by subtraction, as nanoTime() values may be negative or wrap.
+ if (now - next < 0 || !nextSweepNanos.compareAndSet(next, now +
sweepIntervalNanos)) {
+ return;
+ }
+ try {
+ for (Map.Entry<K, PooledResources<V>> e : pool.asMap().entrySet()) {
+ final PooledResources<V> holder = e.getValue();
+ if (now - holder.lastActivityNanos >= abandonmentTimeoutNanos
+ && holder.getUsedCount() == 0
+ && pool.asMap().remove(e.getKey(), holder)) {
Review Comment:
## Follow-up assessment
I rechecked 2 changed files for this follow-up. I agree this is better
classified as P2 given the narrow race window and bounded, temporary overage.
The correctness issue remains: after `getUsedCount() == 0`, `take()` can pass
the holder's acquisition check before `reapAbandonedKeys()` removes and closes
it, so the old holder can still lend while a new holder is created for the same
key. Please make holder invalidation atomic with acquisition, or otherwise make
the racing acquisition reliably retry, and add a concurrency regression test;
the current tests cover idle eviction and on-loan protection but not this
interleaving.
<!-- mergelens:review -->
##########
processing/src/main/java/org/apache/druid/java/util/http/client/pool/ResourcePool.java:
##########
@@ -75,6 +89,44 @@ public PooledResources<V> load(K input)
}
}
);
+ final long unusedTimeoutMillis = config.getUnusedConnectionTimeoutMillis()
> 0
+ ?
config.getUnusedConnectionTimeoutMillis()
+ :
DEFAULT_UNUSED_CONNECTION_TIMEOUT_MILLIS;
+ this.sweepIntervalNanos =
TimeUnit.MILLISECONDS.toNanos(unusedTimeoutMillis);
+ this.abandonmentTimeoutNanos = sweepIntervalNanos > Long.MAX_VALUE /
KEY_ABANDONMENT_TIMEOUT_MULTIPLIER
+ ? Long.MAX_VALUE
+ : sweepIntervalNanos *
KEY_ABANDONMENT_TIMEOUT_MULTIPLIER;
+ this.nextSweepNanos.set(System.nanoTime() + sweepIntervalNanos);
+ }
+
+ /**
+ * Evicts and closes every key that has no resources on loan and has seen no
activity for
+ * {@link #abandonmentTimeoutNanos}, at most once per {@link
#sweepIntervalNanos}. Closing a holder is safe even
+ * if a taker races with it: {@link PooledResources#get()} then returns null
and {@link #take} retries, and
+ * anything lent out is closed when it is given back.
+ */
+ private void reapAbandonedKeys()
+ {
+ final long now = System.nanoTime();
+ final long next = nextSweepNanos.get();
+ // Compared by subtraction, as nanoTime() values may be negative or wrap.
+ if (now - next < 0 || !nextSweepNanos.compareAndSet(next, now +
sweepIntervalNanos)) {
+ return;
+ }
+ try {
+ for (Map.Entry<K, PooledResources<V>> e : pool.asMap().entrySet()) {
+ final PooledResources<V> holder = e.getValue();
+ if (now - holder.lastActivityNanos >= abandonmentTimeoutNanos
+ && holder.getUsedCount() == 0
+ && pool.asMap().remove(e.getKey(), holder)) {
+ holder.close();
Review Comment:
## Follow-up assessment
I rechecked 2 changed files for this follow-up. Agreed that RETAINING
cleanup should continue after a per-resource close failure, but catching and
then clearing the deque would still lose the failed resource after the holder
has already been removed from the cache. Please keep failed resources retryable
(or otherwise record retryable cleanup) while continuing to process later
resources and holders, and add a test that exercises abandonment with a failing
close. This remains a P2 cleanup-leak risk.
<!-- mergelens:review -->
--
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]