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]

Reply via email to