maytasm commented on code in PR #20403:
URL: https://github.com/apache/druid/pull/20403#discussion_r4164103130


##########
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:
   I think this is more of a P2 since
   - It's rare. The key has to have been idle for the whole abandonment timeout 
(40 min by default). Then a take for that exact key has to land in a window of 
microseconds, between the sweep's count check and close(). In ADAPTIVE the 
window is even narrower: between acquirePermit()'s closed check and the 
lentResources increment.
   - Nothing leaks. When that request returns the resource, giveBack() sees the 
holder is closed and closes it, in both implementations.
   - It's bounded and temporary. The overage is +1 for each taker caught in the 
window, and it lasts only as long as that request. That could be minutes for a 
long streaming query, but it never grows.



##########
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:
   I think this should be a simple fix and a nice to have.
   You can add try/catch inside the loop on holder.close() and also Also make 
RETAINING's close() keep going past a failure
   i.e. 
   ```
   public void close()
   {
     synchronized (this) {
       closed = true;
       for (ResourceHolder<V> holder : resourceHolderList) {
         try {
           factory.close(holder.getResource());
         }
         catch (Exception e) {
           log.warn(e, "Failed to close resource at key[%s]", key);
         }
       }
       resourceHolderList.clear();
       this.notifyAll();
     }
   }
   ```



##########
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()

Review Comment:
   Any key evicted and closed here will not be counted in stats. When 
drainStats is called, the pool no longer has the Entry for the evicted Key and 
hence we would miss the Evicted key's stats from the metrics. 
   Maybe add a ConcurrentMap (i.e. `evictedStats`) to hold stats of keys 
evicted since the last drainStats(), so their last closes are still reported.



-- 
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