morningman commented on code in PR #68711: URL: https://github.com/apache/doris/pull/68711#discussion_r4176191912
########## fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionPool.java: ########## @@ -0,0 +1,237 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.fluss; + +import org.apache.fluss.client.Connection; +import org.apache.fluss.client.ConnectionFactory; +import org.apache.fluss.config.Configuration; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.ArrayDeque; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.function.Function; +import java.util.function.LongSupplier; + +/** + * Fluss connections kept open from one scan range to the next, so that a range borrows a connection + * instead of opening one of its own. + * + * <p>A connection is a netty client with its own event loop threads, a metadata cache, and - once a + * range has read a kv snapshot or a remote log segment through it - a pool of download threads. Opening + * one per range cost every range the connect, and a close that waits out netty's two-second graceful + * shutdown. {@link FlussConnectionCloser} took that close off the scanning thread, but only up to the + * number of closes it runs at once: back-to-back queries over many small ranges, or one query over a + * partitioned table with hundreds of them, still left ranges closing their own connections, two + * seconds each. A range read to its end now gives its connection back instead; only one that failed or + * was closed early (a LIMIT, a cancel) has its connection closed, for the reasons {@link #discard} gives. + * + * <p><b>One range per connection at a time.</b> A fluss connection is thread-safe and could serve every + * range at once, but then its download threads ({@code client.remote-file.download-thread-num}, three by + * default) would be shared by every primary-key range copying its kv snapshot at the same moment, where + * each such range used to have three of its own. Lending a connection to one range at a time keeps every + * range's client resources what they were; what changes is that a connection outlives its range and + * serves the next one. The pool never holds more connections than ranges were read at once. + * + * <p><b>Keyed by the whole client configuration.</b> A connection goes back to the ranges that would have + * opened it with the same settings, so ranges of catalogs with different servers, credentials or client + * options never share one. The key holds whatever the configuration holds, credentials included, and is + * never logged. + * + * <p><b>Idle connections are closed after {@link #IDLE_TIMEOUT_NANOS}.</b> What one burst of ranges + * opened is there for the next burst and closed, through the closer, once nothing has borrowed it for + * that long - a BE that stops reading fluss does not keep the threads. + */ +final class FlussConnectionPool { + + private static final Logger LOG = LoggerFactory.getLogger(FlussConnectionPool.class); + + /** How long a connection nobody borrows is kept. */ + static final long IDLE_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(60); + + /** How often idle connections are looked at: one lives at most this much past its timeout. */ + private static final long REAP_INTERVAL_SECONDS = 10; + + static final FlussConnectionPool INSTANCE = + new FlussConnectionPool(ConnectionFactory::createConnection, System::nanoTime); + + static { + ScheduledExecutorService reaper = Executors.newSingleThreadScheduledExecutor(runnable -> { + Thread thread = new Thread(runnable, "fluss-connection-reaper"); + // Never what keeps BE's JVM alive. + thread.setDaemon(true); + // The closer may hand a close back to this thread, and closing can load fluss classes. + thread.setContextClassLoader(FlussConnectionPool.class.getClassLoader()); + return thread; + }); + reaper.scheduleWithFixedDelay(() -> sweep(INSTANCE), REAP_INTERVAL_SECONDS, REAP_INTERVAL_SECONDS, + TimeUnit.SECONDS); + } + + private final Function<Configuration, Connection> factory; + private final LongSupplier nanoTime; + + /** + * Idle connections by the configuration they were opened with, the most recently returned last. A + * configuration with no idle connection has no entry, which is what {@link #borrow} relies on. + */ + private final Map<Map<String, String>, ArrayDeque<Idle>> idle = new HashMap<>(); + + FlussConnectionPool(Function<Configuration, Connection> factory, LongSupplier nanoTime) { + this.factory = factory; + this.nanoTime = nanoTime; + } + + /** + * A connection opened with {@code config}: the one given back most recently if any is idle, else a + * new one. The most recent rather than the oldest, so that after a burst the connections it no longer + * needs are the ones left idle long enough to be closed. + */ + Lease borrow(Configuration config) { + Map<String, String> key = config.toMap(); + synchronized (idle) { + ArrayDeque<Idle> connections = idle.get(key); + if (connections != null) { + Connection connection = connections.pollLast().connection; + if (connections.isEmpty()) { + idle.remove(key); + } + return new Lease(key, connection, false); + } + } + // Outside the lock: opening a connection talks to the cluster. + return new Lease(key, factory.apply(config), true); + } + + /** Takes back the connection of a range that was read to its end, for the next range to borrow. */ + void giveBack(Lease lease) { + long now = nanoTime.getAsLong(); + synchronized (idle) { + idle.computeIfAbsent(lease.key, key -> new ArrayDeque<>()).addLast(new Idle(lease.connection, now)); + } + } + + /** + * Closes, without waiting for it, the connection of a range that failed or was closed before its end, + * instead of lending it again. Such a connection may be broken by what failed it: an + * {@code OutOfMemoryError} on one of its netty threads ends that thread's event loop for good. The + * most recently returned connection is lent first, so one like that, given back, would be lent to + * every range that came next. A primary-key range closed before its kv snapshot arrived hands its + * connection here only once the snapshot copy running on the connection's download threads is over + * ({@code FlussJniScanner#closeInternal}); closed under the copy, the connection would strand it. + */ + void discard(Lease lease) { + FlussConnectionCloser.close(lease.connection); + } + + /** Closes, without waiting for them, the connections nobody has borrowed for {@code idleNanos}. */ + void closeIdleLongerThan(long idleNanos) { + long now = nanoTime.getAsLong(); + List<Connection> expired = new ArrayList<>(); + synchronized (idle) { + Iterator<ArrayDeque<Idle>> keys = idle.values().iterator(); + while (keys.hasNext()) { + ArrayDeque<Idle> connections = keys.next(); + // A deque is in the order its connections came back, so the expired ones are at its head. + while (!connections.isEmpty() && now - connections.peekFirst().since >= idleNanos) { + expired.add(connections.pollFirst().connection); Review Comment: Fixed in bb33f27a1ad, for the sweep. It no longer takes the expired connections out together: it takes one at a time and hands it to the closer at once, with nothing allocated in between (the deque unlinks it, and an emptied entry is removed through the iterator, which reuses the hash the map stored), and logs after the loop. An error mid-sweep now costs at most the connection in hand; the rest stay in the pool for the next sweep. `FlussConnectionCloser.close` closes a connection on the calling thread when building the task or a closer thread fails, instead of dropping it. Test: `FlussConnectionPoolTest.sweepThatFailsMidwayKeepsTheConnectionsItHadNotHandedOver` -- a closer that fails its second hand-over; the third connection stays idle and the next sweep closes it. It fails with the previous sweep. The existing test threw its `OutOfMemoryError` from `nanoTime`, before anything was taken out, so it covered only the next sweeps running. `borrow()` and `giveBack()` stay as they are: each allocates one small object while it moves a connection, so an error there loses at most that connection. No transfer can be made failure-proof under a full heap -- HotSpot can even skip a frame's catch and finally when deoptimization fails to reallocate -- and after BE's JVM runs out of heap the fluss client's own netty event loops can die as well, leaving requests that never complete. What the pool must not do is lose a whole batch, or stop cleaning up for good, and it now does neither. ########## fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussConnectionPool.java: ########## @@ -0,0 +1,237 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.fluss; + +import org.apache.fluss.client.Connection; +import org.apache.fluss.client.ConnectionFactory; +import org.apache.fluss.config.Configuration; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.ArrayDeque; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.function.Function; +import java.util.function.LongSupplier; + +/** + * Fluss connections kept open from one scan range to the next, so that a range borrows a connection + * instead of opening one of its own. + * + * <p>A connection is a netty client with its own event loop threads, a metadata cache, and - once a + * range has read a kv snapshot or a remote log segment through it - a pool of download threads. Opening + * one per range cost every range the connect, and a close that waits out netty's two-second graceful + * shutdown. {@link FlussConnectionCloser} took that close off the scanning thread, but only up to the + * number of closes it runs at once: back-to-back queries over many small ranges, or one query over a + * partitioned table with hundreds of them, still left ranges closing their own connections, two + * seconds each. A range read to its end now gives its connection back instead; only one that failed or + * was closed early (a LIMIT, a cancel) has its connection closed, for the reasons {@link #discard} gives. + * + * <p><b>One range per connection at a time.</b> A fluss connection is thread-safe and could serve every + * range at once, but then its download threads ({@code client.remote-file.download-thread-num}, three by + * default) would be shared by every primary-key range copying its kv snapshot at the same moment, where + * each such range used to have three of its own. Lending a connection to one range at a time keeps every + * range's client resources what they were; what changes is that a connection outlives its range and + * serves the next one. The pool never holds more connections than ranges were read at once. + * + * <p><b>Keyed by the whole client configuration.</b> A connection goes back to the ranges that would have + * opened it with the same settings, so ranges of catalogs with different servers, credentials or client + * options never share one. The key holds whatever the configuration holds, credentials included, and is + * never logged. + * + * <p><b>Idle connections are closed after {@link #IDLE_TIMEOUT_NANOS}.</b> What one burst of ranges + * opened is there for the next burst and closed, through the closer, once nothing has borrowed it for + * that long - a BE that stops reading fluss does not keep the threads. + */ +final class FlussConnectionPool { + + private static final Logger LOG = LoggerFactory.getLogger(FlussConnectionPool.class); + + /** How long a connection nobody borrows is kept. */ + static final long IDLE_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(60); + + /** How often idle connections are looked at: one lives at most this much past its timeout. */ + private static final long REAP_INTERVAL_SECONDS = 10; + + static final FlussConnectionPool INSTANCE = + new FlussConnectionPool(ConnectionFactory::createConnection, System::nanoTime); + + static { + ScheduledExecutorService reaper = Executors.newSingleThreadScheduledExecutor(runnable -> { + Thread thread = new Thread(runnable, "fluss-connection-reaper"); + // Never what keeps BE's JVM alive. + thread.setDaemon(true); + // The closer may hand a close back to this thread, and closing can load fluss classes. + thread.setContextClassLoader(FlussConnectionPool.class.getClassLoader()); + return thread; + }); + reaper.scheduleWithFixedDelay(() -> sweep(INSTANCE), REAP_INTERVAL_SECONDS, REAP_INTERVAL_SECONDS, + TimeUnit.SECONDS); + } + + private final Function<Configuration, Connection> factory; + private final LongSupplier nanoTime; + + /** + * Idle connections by the configuration they were opened with, the most recently returned last. A + * configuration with no idle connection has no entry, which is what {@link #borrow} relies on. + */ + private final Map<Map<String, String>, ArrayDeque<Idle>> idle = new HashMap<>(); + + FlussConnectionPool(Function<Configuration, Connection> factory, LongSupplier nanoTime) { + this.factory = factory; + this.nanoTime = nanoTime; + } + + /** + * A connection opened with {@code config}: the one given back most recently if any is idle, else a + * new one. The most recent rather than the oldest, so that after a burst the connections it no longer + * needs are the ones left idle long enough to be closed. + */ + Lease borrow(Configuration config) { + Map<String, String> key = config.toMap(); + synchronized (idle) { + ArrayDeque<Idle> connections = idle.get(key); + if (connections != null) { + Connection connection = connections.pollLast().connection; + if (connections.isEmpty()) { + idle.remove(key); + } + return new Lease(key, connection, false); + } + } + // Outside the lock: opening a connection talks to the cluster. + return new Lease(key, factory.apply(config), true); + } + + /** Takes back the connection of a range that was read to its end, for the next range to borrow. */ + void giveBack(Lease lease) { + long now = nanoTime.getAsLong(); + synchronized (idle) { + idle.computeIfAbsent(lease.key, key -> new ArrayDeque<>()).addLast(new Idle(lease.connection, now)); Review Comment: The javadoc was wrong across configurations and is corrected in bb33f27a1ad; there is no global cap, on purpose. The bound is per client configuration, and that is the design: connections opened with different settings must not be shared. In practice a catalog has at most two configurations -- `PK_FULL` ranges leave the log read preference at fluss's default and the other ranges set it (`FlussJniScanner#clientConfig`) -- so a BE reading one fluss catalog keeps at most the peaks of those two, for up to about 70 s (the 60 s idle timeout plus the 10 s reaper interval). An idle connection holds one network thread with its selector and the sockets to the servers it reached; its buffers come from netty's shared `PooledByteBufAllocator.DEFAULT`. While a burst runs there is one connection per open range, as on master, where each connection also carried four network threads; a cap on idle connections does not bound that, and below the per-configuration peak it would only turn pooled connections back into a connect and a close per range. Ten catalogs each scanned 128 ranges wide within one minute woul d hold what the example says until the reaper gets to them; I don't see that as an ordinary workload to size a constant for. ########## fe/be-java-extensions/fluss-scanner/src/main/java/org/apache/doris/fluss/FlussJniScanner.java: ########## @@ -348,14 +369,36 @@ protected int getNext() throws IOException { @Override protected void closeInternal() throws IOException { IOException failure = null; + CompletableFuture<Void> released = fullScanner == null + ? CompletableFuture.completedFuture(null) : fullScanner.released(); // Close everything even if an earlier close throws: a leaked fluss connection keeps its netty // and metadata-updater threads alive for the life of the BE process. failure = closeQuietly(scanner, "scanner", failure); scanner = null; + fullScanner = null; failure = closeQuietly(table, "table", failure); table = null; - failure = closeQuietly(connection, "connection", failure); - connection = null; + if (lease != null) { + FlussConnectionPool.Lease returned = lease; + // A range read to its end gives the connection back to serve the next one (see + // FlussConnectionPool); one that failed or was closed before its end (a LIMIT, a cancel) has it + // closed instead, for the reasons FlussConnectionPool#discard gives. + Runnable handBack = finished && failure == null + ? () -> FlussConnectionPool.INSTANCE.giveBack(returned) + : () -> FlussConnectionPool.INSTANCE.discard(returned); + if (released.isDone()) { + handBack.run(); + } else { Review Comment: Fixed in bb33f27a1ad, short of a separate owner: - The waiter keeps waiting after an `OutOfMemoryError` on its own thread, parking one poll interval between attempts. It still closes the reader only after publication or an initializer failure; initializer failures, an `OutOfMemoryError` on fluss's download side included, reach it as exceptions from `pollBatch` (`KvSnapshotBatchScanner#ensureNoException`), so it was only an error raised on the waiter's own thread that ended it. - Building the waiter thread moved inside the `try` that already covered starting it, so a failure to construct it also falls back to waiting on the closing thread. - The deferred hand-back is always a discard (a range whose snapshot never arrived was not read to its end), and `FlussConnectionCloser.close` now closes a connection on the calling thread when it cannot hand it to a closer thread, so the future `thenRun` returns has nothing left to lose. Test: `SafeKvSnapshotAndLogBatchScannerTest.publicationWaiterOutlastsAnOutOfMemoryErrorOnItsOwnThread` -- the waiter's first poll throws `OutOfMemoryError`; it polls again, closes the reader once after publication, and `released()` completes. Without the change the waiter dies and never polls again. For the record, the comparison with master: there, closing a primary-key range before its snapshot arrived closed the connection under the copy, and the copy thread, the waiter and the half-copied directory leaked every time. The deferral in this PR is what fixes that; what remains needs an error on the waiter's own thread. A registry swept by the reaper would also cover a deoptimization that skips the waiter's catch, but it would replace the waiter rather than harden it, for a case in which the fluss client's own threads do not recover either. -- 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]
