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]

Reply via email to