morningman commented on code in PR #68711: URL: https://github.com/apache/doris/pull/68711#discussion_r4176091600
########## 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); Review Comment: No change here: the setting cannot take effect in BE today, and the shared state predates the pool. The mechanism is as described. `FlussConnection`'s constructor re-initializes fluss's process-wide `FileSystem` with the connection's `client.fs.*` options, and kv snapshot and remote log files are resolved through it (`RemoteFileDownloader#downloadFile` -> `FsPath#getFileSystem` -> `FileSystem.get`). What BE's fluss plugin lacks is a filesystem that reads those options: fluss-client 1.0.0 registers no `FileSystemPlugin`, `FlussConnection` passes a null `PluginManager`, and `be/plugins/jni/fluss` deploys no fluss-fs-* jar, so the only filesystem `FileSystem` loads is its built-in local one, and `LocalFileSystemPlugin#create` ignores the configuration. Whichever catalog initialized last, `client.fs.*` reaches nothing BE uses -- which is also why an A/B/A test would pass without testing anything. Nor is the shared state the pool's. On master every range opened a connection of its own and re-initialized the same static configuration, so ranges of two catalogs with different `client.fs.*` reading at the same time already resolved files with whichever had opened last; pooling widens that from concurrent to back to back. Settings per catalog need fluss to hold the filesystem configuration per connection, or Doris to accept one set per BE. That belongs with support for remote filesystems in BE's fluss plugin rather than in this PR. -- 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]
