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


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

Review Comment:
   why do we have to go to nano precision? 
   other config is `millis` either follow that or add `Duration` or something 
which could be handled better.



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

Review Comment:
   put a method on the `holder` and ask that....
   
   ```
   if(holder.isAbandoned() && ... ) 
   ```
   is pretty much talks about what it checks - it can be as complicated as it 
wants inside that method
   



##########
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.
+   */

Review Comment:
   did you wrote this apdoc? 
   it just explains the function in unfollowable plain text 
   



##########
processing/src/main/java/org/apache/druid/java/util/http/client/pool/ResourcePool.java:
##########
@@ -83,6 +135,7 @@ public PooledResources<V> load(K input)
    */
   public Map<K, Stats> drainStats()
   {
+    reapAbandonedKeys();

Review Comment:
   what are you doing here in a `drainStats`  call?



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

Review Comment:
   do we really need coimments like this? 



##########
processing/src/main/java/org/apache/druid/java/util/http/client/pool/ResourcePool.java:
##########
@@ -55,12 +56,25 @@
  * {@link ResourceFactory#isGood} rejects it. With eagerInitialization a key 
starts out full, otherwise empty.
  *
  * {@link ResourcePoolConfig#getPoolImplementation()} selects which {@link 
Implementation} does the pooling.
+ *
+ * A key with no resources on loan that has been neither taken from nor 
returned to for
+ * {@link #KEY_ABANDONMENT_TIMEOUT_MULTIPLIER} times {@link 
ResourcePoolConfig#getUnusedConnectionTimeoutMillis()} is
+ * evicted and closed, so that destinations which are no longer used (e.g. 
finished indexer tasks) do not keep their
+ * resources cached forever. Eviction is checked opportunistically from {@link 
#take} and {@link #drainStats()}, at
+ * most once per sweep interval, which is enough to bound the pool's size 
since it only grows through {@link #take}.
+ * A key is never evicted while any of its resources are on loan.

Review Comment:
   I think this could be described in a single line.....long docs are not read 
because people thing they are AI generated....(is it?)
   
   ```
   The pool is purged from stale keys periodically.
   ```
   
   if someone is interested; could dive into the code...if not...that's already 
plenty



##########
processing/src/main/java/org/apache/druid/java/util/http/client/pool/ResourcePool.java:
##########
@@ -55,12 +56,25 @@
  * {@link ResourceFactory#isGood} rejects it. With eagerInitialization a key 
starts out full, otherwise empty.
  *
  * {@link ResourcePoolConfig#getPoolImplementation()} selects which {@link 
Implementation} does the pooling.
+ *
+ * A key with no resources on loan that has been neither taken from nor 
returned to for
+ * {@link #KEY_ABANDONMENT_TIMEOUT_MULTIPLIER} times {@link 
ResourcePoolConfig#getUnusedConnectionTimeoutMillis()} is
+ * evicted and closed, so that destinations which are no longer used (e.g. 
finished indexer tasks) do not keep their
+ * resources cached forever. Eviction is checked opportunistically from {@link 
#take} and {@link #drainStats()}, at
+ * most once per sweep interval, which is enough to bound the pool's size 
since it only grows through {@link #take}.
+ * A key is never evicted while any of its resources are on loan.
  */
 public class ResourcePool<K, V> implements Closeable
 {
   private static final Logger log = new Logger(ResourcePool.class);
+  // How many multiples of unusedConnectionTimeoutMillis a key must go unused 
for before it is considered abandoned
+  private static final long KEY_ABANDONMENT_TIMEOUT_MULTIPLIER = 10;

Review Comment:
   put this into the config....even if it remains a static final...



##########
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();
+        }
+      }
+    }
+    catch (Throwable t) {
+      // Housekeeping must never fail the caller's take() or drainStats().

Review Comment:
   let's put a comment inside the method body which is documenting the call 
sites of the function....this is a terrible idea....
   
   unleashed AI?



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

Review Comment:
   the defaulting should happen in the `config` class ....instead of exposing 
it at all call sites



##########
processing/src/main/java/org/apache/druid/java/util/http/client/pool/ResourcePool.java:
##########
@@ -661,6 +731,8 @@ V get()
       if (!acquirePermit()) {
         return null;
       }
+      // Counted as used from the permit on, so a slow createResource() 
doesn't look unused to reapAbandonedKeys().
+      lentResources.incrementAndGet();

Review Comment:
   looks like a pointless optimization - undo this



##########
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:
   not sure if its worth to put so much effort into `drainStats` - it needs an 
extra config; also for a removed pool it might not really be relevant...
   also it opens problem we might not want to handle...like there was pool to X 
which is abandoned and re-created in the same minute....
   also if the stats are not drained by anything....there is a chance that 
those will never be drained and thus accumlate forever...
   
   I think if we really want to correct this - another way to implement stats 
is to utilize a listener pattern; which by design will be resilient for an 
issue like this
   
   when the removeal of these abandoned stuff happens - these should already be 
idle for quite some time....so
   
   I feel like we could just let this slip....



##########
processing/src/main/java/org/apache/druid/java/util/http/client/pool/ResourcePool.java:
##########
@@ -101,18 +154,29 @@ public ResourceContainer<V> take(final K key)
       return null;
     }
 
-    final PooledResources<V> holder;
-    try {
-      holder = pool.get(key);
-    }
-    catch (ExecutionException e) {
-      throw new RuntimeException(e);
-    }
-    final V value = holder.get();
+    reapAbandonedKeys();
+
+    // A holder evicted by reapAbandonedKeys() right after being fetched here 
answers get() with null even though the
+    // pool is open, so retry with a fresh holder. A null from an interrupted 
wait is not retried: the interrupt flag
+    // is restored in that case, so checking it without clearing it avoids 
spinning.
+    PooledResources<V> holder;
+    V value;
+    do {
+      try {
+        holder = pool.get(key);
+      }
+      catch (ExecutionException e) {
+        throw new RuntimeException(e);
+      }
+      holder.lastActivityNanos = System.nanoTime();

Review Comment:
   I think the holder should update this - not an external entity



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

Review Comment:
   is it really necessary to go down to comparision and adding time ? isn't 
there a higher level construct to limit the rate of something?
   there is `RateLimiter`  which could do exactly this



##########
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:
   note: probably `Closeable` and/or `IOUtils.closeQuietly`



##########
processing/src/main/java/org/apache/druid/java/util/http/client/pool/ResourcePool.java:
##########
@@ -101,18 +154,29 @@ public ResourceContainer<V> take(final K key)
       return null;
     }
 
-    final PooledResources<V> holder;
-    try {
-      holder = pool.get(key);
-    }
-    catch (ExecutionException e) {
-      throw new RuntimeException(e);
-    }
-    final V value = holder.get();
+    reapAbandonedKeys();
+
+    // A holder evicted by reapAbandonedKeys() right after being fetched here 
answers get() with null even though the
+    // pool is open, so retry with a fresh holder. A null from an interrupted 
wait is not retried: the interrupt flag
+    // is restored in that case, so checking it without clearing it avoids 
spinning.

Review Comment:
   I think it would be better to have different logics....you are adding a loop 
and whatnot...
   
   for adaptive it would be safe to dispose it if you could allocate all tokens 
from the `Semaphore` without wait....its already a parallel programming 
construct so that'll be safe
   
   I think it would be better to use a ReadWriteLock or some other construct - 
I think this part is not that hot that we should be juggling with primitives...
   



##########
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();
+        }
+      }
+    }
+    catch (Throwable t) {

Review Comment:
   do you really want to catch throwable?
   that needs justification....as it includes `InterruptedException` ; all 
subclasses of `Error` including `VirtualMachineError` which includes things 
like `OOM` and other stuff



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