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]