FrankChen021 commented on code in PR #20403:
URL: https://github.com/apache/druid/pull/20403#discussion_r4155409049
##########
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:
[P2] Do not abort reaping on a close failure
**Finding:** For RETAINING, ResourceHolderPerKey.close() calls factory.close
directly and can throw. Because this call is inside one try around the entire
sweep, a single close exception exits the loop after the holder has already
been removed; remaining idle resources in that holder are never cleared and
future sweeps cannot find it, while later abandoned keys wait for another
sweep. This turns a close failure into a permanent cleanup leak.
**Suggestion:** Isolate failures per holder/resource and continue draining,
while preserving a way to retry any resource whose cleanup failed; add a
retaining-implementation test with a failing close during abandonment.
##########
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:
[P1] Serialize eviction with concurrent borrows
**Finding:** The sweep reads getUsedCount() and then removes the holder in a
separate operation. A take() can fetch the same holder after the zero-count
read and enter get() before holder.close() runs; the sweep then detaches and
closes a holder that has just lent a resource. The next take for the key
creates another holder, so a long-lived request can keep the old resource while
the new holder lends up to maxPerKey, violating the pool's per-key limit and
the retry assumption in the comment.
**Suggestion:** Coordinate the idle check, removal, and close with holder
acquisition, or mark/close the holder so a racing get() reliably returns null
and take() retries; add a concurrent sweep-versus-take test for both
implementations.
--
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]