rzo1 commented on code in PR #9097:
URL: https://github.com/apache/storm/pull/9097#discussion_r4168994869


##########
storm-client/src/jvm/org/apache/storm/utils/StormThreadFactory.java:
##########
@@ -0,0 +1,85 @@
+/*
+ * 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.storm.utils;
+
+import java.util.Map;
+import java.util.concurrent.ThreadFactory;
+import org.apache.storm.Config;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Builds the {@link ThreadFactory} used by Storm's blocking I/O thread pools.
+ *
+ * <p>When {@link Config#STORM_VIRTUAL_THREADS_ENABLED} is true the factory 
creates virtual threads.
+ * Otherwise it creates non-daemon platform threads with normal priority, 
matching
+ * {@code Executors.defaultThreadFactory()}. In both modes threads are named
+ * {@code <namePrefix>-<n>} with {@code n} starting at 0.
+ */
+public final class StormThreadFactory {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(StormThreadFactory.class);
+
+    private StormThreadFactory() {
+    }
+
+    /**
+     * Create a thread factory for a pool.
+     *
+     * @param conf       cluster or topology configuration
+     * @param namePrefix prefix for the thread names, e.g. {@code 
"nimbus-thrift-handler"}
+     * @return a factory producing virtual or platform threads depending on 
the configuration
+     */
+    public static ThreadFactory create(Map<String, Object> conf, String 
namePrefix) {
+        String prefix = namePrefix + "-";
+        if (isVirtualEnabled(conf)) {
+            return Thread.ofVirtual().name(prefix, 0).factory();
+        }
+        return Thread.ofPlatform()
+                .name(prefix, 0)
+                .daemon(false)
+                .priority(Thread.NORM_PRIORITY)
+                .factory();
+    }
+
+    /**
+     * Whether virtual threads are enabled for blocking I/O pools.
+     *
+     * <p>Boolean and String values are honoured (strings are parsed as 
boolean). Any other type
+     * is treated as disabled. This method never throws.
+     */
+    public static boolean isVirtualEnabled(Map<String, Object> conf) {
+        if (conf == null) {
+            return false;
+        }
+        Object value = conf.get(Config.STORM_VIRTUAL_THREADS_ENABLED);
+        if (value == null) {
+            return false;
+        }
+        if (value instanceof Boolean) {
+            return (Boolean) value;
+        }
+        if (value instanceof String) {
+            return Boolean.parseBoolean((String) value);

Review Comment:
   `"yes"` or `"1"` end up as `false` here without any warning, only non String 
types hit the WARN below. Daemon and topology confs both go through 
`@IsBoolean` validation anyway, so the `String` case should only come from 
`storm.options`. Fine to keep, but please don't replace it with 
`ObjectReader.getBoolean`, that one throws on a `String`.



##########
storm-client/src/jvm/org/apache/storm/Config.java:
##########
@@ -869,6 +869,33 @@ public class Config extends HashMap<String, Object> {
      */
     @IsInteger
     public static final String TOPOLOGY_WORKER_SHARED_THREAD_POOL_SIZE = 
"topology.worker.shared.thread.pool.size";
+    /**
+     * When true, Storm's blocking I/O thread pools run on Java virtual 
threads instead of platform threads.
+     * Affected pools: the Nimbus and Supervisor Thrift server handlers (SASL, 
TLS and simple transports), the
+     * supervisor blob localizer download and task executors, the Nimbus 
assignment distribution service, the
+     * supervisor heartbeat executor, the DRPC spout background executor and 
the worker shared thread pool
+     * exposed through the TopologyContext. Pool sizes configured through the 
corresponding {@code *.threads}
+     * settings keep their meaning as a bound on concurrency. Spout/bolt 
executor threads, worker transfer
+     * threads and Netty event loops are never affected.
+     *
+     * <p>DRPC server handler pools ({@code drpc.worker.threads}, {@code 
drpc.invocations.threads}) are also
+     * affected, since the DRPC server uses the same transport plugins.
+     *
+     * <p>With the simple (non-authenticated) Thrift transport, Storm builds 
its own handler pool of
+     * {@code *.threads} threads whenever a {@code *.queue.size} is 
configured, which the shipped defaults do
+     * for Nimbus, Supervisor and DRPC. Only when the queue size is explicitly 
unset does {@code THsHaServer}
+     * fall back to its own pool, whose effective concurrency is 5; enabling 
this flag in that case raises it to
+     * the configured {@code *.threads} value.
+     *
+     * <p>All virtual threads in a JVM share one carrier scheduler sized to 
the number of available processors;
+     * blocking file I/O (for example blob downloads in the supervisor 
localizer) occupies a carrier, so
+     * operators enabling this on I/O-heavy supervisors should size {@code 
jdk.virtualThreadScheduler.parallelism}
+     * / {@code jdk.virtualThreadScheduler.maxPoolSize} accordingly.
+     *
+     * <p>A topology may override this key in its own config to control its 
worker shared executor.

Review Comment:
   `DRPCSpout` also builds its factory from the conf passed to `open()`, so its 
pool is topology controlled too.



##########
examples/storm-perf/src/main/java/org/apache/storm/perf/ThriftHandlerVirtualThreadsBench.java:
##########
@@ -0,0 +1,286 @@
+/*
+ * 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.storm.perf;
+
+import java.io.IOException;
+import java.lang.management.ManagementFactory;
+import java.lang.management.ThreadMXBean;
+import java.lang.reflect.InvocationHandler;
+import java.lang.reflect.Method;
+import java.lang.reflect.Proxy;
+import java.nio.file.Files;
+import java.nio.file.Paths;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.apache.storm.Config;
+import org.apache.storm.generated.Nimbus;
+import org.apache.storm.security.auth.SimpleTransportPlugin;
+import org.apache.storm.security.auth.ThriftConnectionType;
+import org.apache.storm.security.auth.ThriftServer;
+import org.apache.storm.utils.NimbusClient;
+import org.apache.storm.utils.Utils;
+
+/**
+ * Benchmark for {@code storm.virtual.threads.enabled} on the Nimbus Thrift 
handler pool.
+ *
+ * <p>Starts an in-process {@link ThriftServer} using {@link 
SimpleTransportPlugin} whose {@code getNimbusConf}
+ * handler sleeps for a configurable time to simulate blocking I/O (ZooKeeper, 
blob store), then hammers it with
+ * N client threads and reports throughput, latency percentiles, peak 
platform-thread count and RSS. Run it with
+ * the flag off and on to compare.
+ *
+ * <p>Build and run (JDK 25, Maven 3.9):
+ * <pre>
+ * mvn -pl storm-client,examples/storm-perf -am install -DskipTests 
-Dcheckstyle.skip=true -q
+ * mvn -pl examples/storm-perf dependency:build-classpath 
-Dmdep.outputFile=target/cp.txt -q
+ * java -cp examples/storm-perf/target/classes:$(cat 
examples/storm-perf/target/cp.txt) \
+ *     org.apache.storm.perf.ThriftHandlerVirtualThreadsBench --mode both 
--clients 200 --calls 50 --io-ms 10
+ * </pre>
+ *
+ * <p>Options: {@code --mode off|on|both}, {@code --clients N}, {@code --calls 
M} (per client),
+ * {@code --io-ms L} (simulated handler I/O), {@code --threads T} ({@code 
nimbus.thrift.threads}),
+ * {@code --queue-size Q|none} ({@code nimbus.queue.size}; defaults to the 
shipped 100000, so both modes use a pool of
+ * {@code --threads} handlers; {@code none} removes the key so the 
platform-thread run falls back to THsHaServer's own
+ * pool, whose effective concurrency is 5), {@code --warmup W} (calls per 
client excluded from statistics).
+ *
+ * <p>For comparable RSS numbers run each mode in its own JVM ({@code --mode 
off} then {@code --mode on}).
+ */
+public final class ThriftHandlerVirtualThreadsBench {
+
+    private static final String HOST = "localhost";
+
+    private ThriftHandlerVirtualThreadsBench() {
+    }
+
+    public static void main(String[] args) throws Exception {
+        Map<String, String> opts = parseArgs(args);
+        String mode = opts.getOrDefault("mode", 
"both").toLowerCase(Locale.ROOT);
+        int clients = Integer.parseInt(opts.getOrDefault("clients", "200"));
+        int calls = Integer.parseInt(opts.getOrDefault("calls", "50"));
+        long ioMs = Long.parseLong(opts.getOrDefault("io-ms", "10"));
+        int threads = Integer.parseInt(opts.getOrDefault("threads", "64"));
+        String queueOpt = opts.get("queue-size");
+        boolean noQueue = "none".equalsIgnoreCase(queueOpt);
+        Integer queueSize = queueOpt == null || noQueue ? null : 
Integer.valueOf(queueOpt);
+        int warmup = Integer.parseInt(opts.getOrDefault("warmup", "5"));
+
+        System.out.printf(Locale.ROOT, "clients=%d calls=%d io-ms=%d 
threads=%d queue-size=%s warmup=%d%n",
+                          clients, calls, ioMs, threads, noQueue ? "none" : 
queueSize, warmup);
+        if (mode.equals("off") || mode.equals("both")) {
+            runPhase(false, clients, calls, ioMs, threads, queueSize, noQueue, 
warmup);
+        }
+        if (mode.equals("on") || mode.equals("both")) {
+            runPhase(true, clients, calls, ioMs, threads, queueSize, noQueue, 
warmup);
+        }
+    }
+
+    private static void runPhase(boolean virtual, int clients, int calls, long 
ioMs, int threads,
+                                 Integer queueSize, boolean noQueue, int 
warmup) throws Exception {
+        Map<String, Object> conf = new HashMap<>(Utils.readDefaultConfig());
+        conf.put(Config.STORM_THRIFT_TRANSPORT_PLUGIN, 
SimpleTransportPlugin.class.getName());
+        conf.put(Config.NIMBUS_THRIFT_PORT, 0);
+        conf.put(Config.NIMBUS_THRIFT_THREADS, threads);
+        conf.put(Config.STORM_VIRTUAL_THREADS_ENABLED, virtual);
+        conf.put(Config.STORM_NIMBUS_RETRY_TIMES, 0);
+        if (noQueue) {
+            conf.remove(Config.NIMBUS_QUEUE_SIZE);
+        } else if (queueSize != null) {
+            conf.put(Config.NIMBUS_QUEUE_SIZE, queueSize);
+        }
+
+        AtomicInteger inFlight = new AtomicInteger();
+        AtomicInteger peakInFlight = new AtomicInteger();
+        Nimbus.Iface handler = sleepingHandler(ioMs, inFlight, peakInFlight);
+        ThriftServer server = new ThriftServer(conf, new 
Nimbus.Processor<>(handler), ThriftConnectionType.NIMBUS);
+        Thread serveThread = new Thread(server::serve, "bench-thrift-serve");
+        serveThread.setDaemon(true);
+        serveThread.start();
+        while (!server.isServing()) {
+            Thread.sleep(10);
+        }
+        int port = server.getPort();
+
+        ThreadMXBean threadMx = ManagementFactory.getThreadMXBean();
+        long rssBeforeKb = rssKb();
+        int threadsBefore = threadMx.getThreadCount();
+        AtomicInteger peakThreads = new AtomicInteger(threadsBefore);
+        Thread sampler = new Thread(() -> {
+            while (!Thread.currentThread().isInterrupted()) {
+                peakThreads.accumulateAndGet(threadMx.getThreadCount(), 
Math::max);
+                try {
+                    Thread.sleep(20);
+                } catch (InterruptedException e) {
+                    return;
+                }
+            }
+        }, "bench-thread-sampler");
+        sampler.setDaemon(true);
+        sampler.start();
+
+        long[][] latenciesNs = new long[clients][calls];
+        List<Throwable> failures = new ArrayList<>();
+        CountDownLatch ready = new CountDownLatch(clients);
+        CountDownLatch go = new CountDownLatch(1);
+        CountDownLatch done = new CountDownLatch(clients);
+        for (int c = 0; c < clients; c++) {
+            final int clientIdx = c;
+            Thread t = new Thread(() -> {
+                try (NimbusClient client = 
NimbusClient.Builder.withConf(conf).withTimeout(60_000)
+                        .buildWithNimbusHostPort(HOST, port)) {
+                    Nimbus.Iface iface = client.getClient();
+                    for (int i = 0; i < warmup; i++) {
+                        iface.getNimbusConf();
+                    }
+                    ready.countDown();
+                    go.await();
+                    for (int i = 0; i < calls; i++) {
+                        long start = System.nanoTime();
+                        iface.getNimbusConf();
+                        latenciesNs[clientIdx][i] = System.nanoTime() - start;
+                    }
+                } catch (Throwable e) {
+                    synchronized (failures) {
+                        failures.add(e);
+                    }
+                    ready.countDown();
+                } finally {
+                    done.countDown();
+                }
+            }, "bench-client-" + c);
+            t.start();
+        }
+        ready.await();
+        long startNs = System.nanoTime();
+        go.countDown();
+        done.await();
+        long elapsedNs = System.nanoTime() - startNs;
+
+        sampler.interrupt();
+        sampler.join(1000);
+        long rssAfterKb = rssKb();
+        server.stop();
+        serveThread.join(5000);
+
+        report(virtual, clients, calls, ioMs, threads, queueSize, latenciesNs, 
elapsedNs, failures.size(),
+               threadsBefore, peakThreads.get(), peakInFlight.get(), 
rssBeforeKb, rssAfterKb);
+    }
+
+    private static void report(boolean virtual, int clients, int calls, long 
ioMs, int threads, Integer queueSize,
+                               long[][] latenciesNs, long elapsedNs, int 
failures, int threadsBefore, int peakThreads,
+                               int peakInFlight, long rssBeforeKb, long 
rssAfterKb) {
+        long[] all = 
Arrays.stream(latenciesNs).flatMapToLong(Arrays::stream).filter(v -> v > 
0).sorted().toArray();
+        int total = clients * calls;

Review Comment:
   Counts failed calls into the throughput, while the latencies filter them out.



##########
examples/storm-perf/src/main/java/org/apache/storm/perf/ThriftHandlerVirtualThreadsBench.java:
##########
@@ -0,0 +1,286 @@
+/*
+ * 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.storm.perf;
+
+import java.io.IOException;
+import java.lang.management.ManagementFactory;
+import java.lang.management.ThreadMXBean;
+import java.lang.reflect.InvocationHandler;
+import java.lang.reflect.Method;
+import java.lang.reflect.Proxy;
+import java.nio.file.Files;
+import java.nio.file.Paths;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.apache.storm.Config;
+import org.apache.storm.generated.Nimbus;
+import org.apache.storm.security.auth.SimpleTransportPlugin;
+import org.apache.storm.security.auth.ThriftConnectionType;
+import org.apache.storm.security.auth.ThriftServer;
+import org.apache.storm.utils.NimbusClient;
+import org.apache.storm.utils.Utils;
+
+/**
+ * Benchmark for {@code storm.virtual.threads.enabled} on the Nimbus Thrift 
handler pool.
+ *
+ * <p>Starts an in-process {@link ThriftServer} using {@link 
SimpleTransportPlugin} whose {@code getNimbusConf}
+ * handler sleeps for a configurable time to simulate blocking I/O (ZooKeeper, 
blob store), then hammers it with
+ * N client threads and reports throughput, latency percentiles, peak 
platform-thread count and RSS. Run it with
+ * the flag off and on to compare.
+ *
+ * <p>Build and run (JDK 25, Maven 3.9):
+ * <pre>
+ * mvn -pl storm-client,examples/storm-perf -am install -DskipTests 
-Dcheckstyle.skip=true -q
+ * mvn -pl examples/storm-perf dependency:build-classpath 
-Dmdep.outputFile=target/cp.txt -q
+ * java -cp examples/storm-perf/target/classes:$(cat 
examples/storm-perf/target/cp.txt) \
+ *     org.apache.storm.perf.ThriftHandlerVirtualThreadsBench --mode both 
--clients 200 --calls 50 --io-ms 10
+ * </pre>
+ *
+ * <p>Options: {@code --mode off|on|both}, {@code --clients N}, {@code --calls 
M} (per client),
+ * {@code --io-ms L} (simulated handler I/O), {@code --threads T} ({@code 
nimbus.thrift.threads}),
+ * {@code --queue-size Q|none} ({@code nimbus.queue.size}; defaults to the 
shipped 100000, so both modes use a pool of
+ * {@code --threads} handlers; {@code none} removes the key so the 
platform-thread run falls back to THsHaServer's own
+ * pool, whose effective concurrency is 5), {@code --warmup W} (calls per 
client excluded from statistics).
+ *
+ * <p>For comparable RSS numbers run each mode in its own JVM ({@code --mode 
off} then {@code --mode on}).
+ */
+public final class ThriftHandlerVirtualThreadsBench {
+
+    private static final String HOST = "localhost";
+
+    private ThriftHandlerVirtualThreadsBench() {
+    }
+
+    public static void main(String[] args) throws Exception {
+        Map<String, String> opts = parseArgs(args);
+        String mode = opts.getOrDefault("mode", 
"both").toLowerCase(Locale.ROOT);

Review Comment:
   Default `both` runs both modes in one JVM (shared heap and JIT), while the 
javadoc recommends separate JVMs. I would default to one mode.



##########
storm-client/src/jvm/org/apache/storm/security/auth/SimpleTransportPlugin.java:
##########
@@ -74,9 +78,18 @@ public TServer getServer(TProcessor processor) throws 
IOException, TTransportExc
 
         serverArgs.maxReadBufferBytes = maxBufferSize;
 
-        if (queueSize != null) {
+        if (queueSize != null || 
StormThreadFactory.isVirtualEnabled(topoConf)) {
+            // THsHaServer builds its own platform-thread pool (core 5, 
unbounded LinkedBlockingQueue) when no
+            // executor is supplied, so when virtual threads are requested we 
always supply an executor;
+            // without a configured queue size we use the same unbounded queue 
so no request is rejected
+            // that would not be rejected today. The shipped defaults always 
configure a queue size.

Review Comment:
   Not true for `DRPC_INVOCATIONS`, its queue key is hard coded `null` in 
`ThriftConnectionType`. So with the flag on this branch is always taken there 
(the 5 -> 64 jump from @reiabreu's review), and also for anyone who removed 
`nimbus.queue.size` or `supervisor.queue.size`. Same statement in the javadoc 
of the new key.



##########
storm-client/src/jvm/org/apache/storm/Config.java:
##########
@@ -869,6 +869,33 @@ public class Config extends HashMap<String, Object> {
      */
     @IsInteger
     public static final String TOPOLOGY_WORKER_SHARED_THREAD_POOL_SIZE = 
"topology.worker.shared.thread.pool.size";
+    /**
+     * When true, Storm's blocking I/O thread pools run on Java virtual 
threads instead of platform threads.
+     * Affected pools: the Nimbus and Supervisor Thrift server handlers (SASL, 
TLS and simple transports), the
+     * supervisor blob localizer download and task executors, the Nimbus 
assignment distribution service, the
+     * supervisor heartbeat executor, the DRPC spout background executor and 
the worker shared thread pool
+     * exposed through the TopologyContext. Pool sizes configured through the 
corresponding {@code *.threads}
+     * settings keep their meaning as a bound on concurrency. Spout/bolt 
executor threads, worker transfer
+     * threads and Netty event loops are never affected.
+     *
+     * <p>DRPC server handler pools ({@code drpc.worker.threads}, {@code 
drpc.invocations.threads}) are also
+     * affected, since the DRPC server uses the same transport plugins.
+     *
+     * <p>With the simple (non-authenticated) Thrift transport, Storm builds 
its own handler pool of
+     * {@code *.threads} threads whenever a {@code *.queue.size} is 
configured, which the shipped defaults do
+     * for Nimbus, Supervisor and DRPC. Only when the queue size is explicitly 
unset does {@code THsHaServer}
+     * fall back to its own pool, whose effective concurrency is 5; enabling 
this flag in that case raises it to
+     * the configured {@code *.threads} value.
+     *
+     * <p>All virtual threads in a JVM share one carrier scheduler sized to 
the number of available processors;
+     * blocking file I/O (for example blob downloads in the supervisor 
localizer) occupies a carrier, so

Review Comment:
   On 25 the scheduler compensates for blocking file I/O (temporarily adds 
carriers up to `maxPoolSize`), and since JEP 491 `synchronized` no longer pins. 
What still pins is native frames (Hadoop native libs, JNI in user code on the 
shared executor). Raising `parallelism` is not the right advice here.
   
   What is missing instead: virtual threads are not time sliced, so CPU heavy 
Nimbus calls (`getTopologyPageInfo`, `submitTopology` on a big cluster) can 
hold all carriers and delay cheap ones. Probably rare, but should be mentioned.
   
   And with the flag on, `jstack` and `Thread.getAllStackTraces()` do not show 
virtual threads. That makes `Utils.threadDump()` blind for these pools (stuck 
slot dump in `ReadClusterState`, the forced halt dump in `Utils`), and the 
jstack action from the UI via `flight.bash` as well. A hung blob download 
simply disappears from the dump. Only `jcmd <pid> Thread.dump_to_file` shows 
them. I think this needs at least a note here, better `flight.bash` using 
`jcmd` when the flag is on.



##########
storm-server/src/test/java/org/apache/storm/security/auth/AuthTest.java:
##########
@@ -647,4 +647,92 @@ public void impersonationAuthorizerTest() throws Exception 
{
     public interface MyBiConsumer<T, U> {
         void accept(T t, U u) throws Exception;
     }
-}
\ No newline at end of file
+
+    @Test
+    public void simpleTransportRunsHandlerOnVirtualThreadWhenEnabled() throws 
Exception {

Review Comment:
   This does not test the no queue size path. `withServer` starts from 
`ConfigUtils.readStormConfig()`, which already has `nimbus.queue.size: 100000` 
from `defaults.yaml`, so `queueSize != null` and it runs the same branch as the 
`WithQueueSize` test. The `LinkedBlockingQueue` arm is not covered at all, and 
that is exactly the DRPC case. Setting `NIMBUS_QUEUE_SIZE` to `null` in `extra` 
or using `DRPC_INVOCATIONS` should do. The description claims both are covered.
   
   TLS has no test either.



##########
examples/storm-perf/src/main/java/org/apache/storm/perf/ThriftHandlerVirtualThreadsBench.java:
##########
@@ -0,0 +1,286 @@
+/*
+ * 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.storm.perf;
+
+import java.io.IOException;
+import java.lang.management.ManagementFactory;
+import java.lang.management.ThreadMXBean;
+import java.lang.reflect.InvocationHandler;
+import java.lang.reflect.Method;
+import java.lang.reflect.Proxy;
+import java.nio.file.Files;
+import java.nio.file.Paths;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.apache.storm.Config;
+import org.apache.storm.generated.Nimbus;
+import org.apache.storm.security.auth.SimpleTransportPlugin;
+import org.apache.storm.security.auth.ThriftConnectionType;
+import org.apache.storm.security.auth.ThriftServer;
+import org.apache.storm.utils.NimbusClient;
+import org.apache.storm.utils.Utils;
+
+/**
+ * Benchmark for {@code storm.virtual.threads.enabled} on the Nimbus Thrift 
handler pool.
+ *
+ * <p>Starts an in-process {@link ThriftServer} using {@link 
SimpleTransportPlugin} whose {@code getNimbusConf}
+ * handler sleeps for a configurable time to simulate blocking I/O (ZooKeeper, 
blob store), then hammers it with
+ * N client threads and reports throughput, latency percentiles, peak 
platform-thread count and RSS. Run it with
+ * the flag off and on to compare.
+ *
+ * <p>Build and run (JDK 25, Maven 3.9):
+ * <pre>
+ * mvn -pl storm-client,examples/storm-perf -am install -DskipTests 
-Dcheckstyle.skip=true -q
+ * mvn -pl examples/storm-perf dependency:build-classpath 
-Dmdep.outputFile=target/cp.txt -q
+ * java -cp examples/storm-perf/target/classes:$(cat 
examples/storm-perf/target/cp.txt) \
+ *     org.apache.storm.perf.ThriftHandlerVirtualThreadsBench --mode both 
--clients 200 --calls 50 --io-ms 10
+ * </pre>
+ *
+ * <p>Options: {@code --mode off|on|both}, {@code --clients N}, {@code --calls 
M} (per client),
+ * {@code --io-ms L} (simulated handler I/O), {@code --threads T} ({@code 
nimbus.thrift.threads}),
+ * {@code --queue-size Q|none} ({@code nimbus.queue.size}; defaults to the 
shipped 100000, so both modes use a pool of
+ * {@code --threads} handlers; {@code none} removes the key so the 
platform-thread run falls back to THsHaServer's own
+ * pool, whose effective concurrency is 5), {@code --warmup W} (calls per 
client excluded from statistics).
+ *
+ * <p>For comparable RSS numbers run each mode in its own JVM ({@code --mode 
off} then {@code --mode on}).
+ */
+public final class ThriftHandlerVirtualThreadsBench {
+
+    private static final String HOST = "localhost";
+
+    private ThriftHandlerVirtualThreadsBench() {
+    }
+
+    public static void main(String[] args) throws Exception {
+        Map<String, String> opts = parseArgs(args);
+        String mode = opts.getOrDefault("mode", 
"both").toLowerCase(Locale.ROOT);
+        int clients = Integer.parseInt(opts.getOrDefault("clients", "200"));
+        int calls = Integer.parseInt(opts.getOrDefault("calls", "50"));
+        long ioMs = Long.parseLong(opts.getOrDefault("io-ms", "10"));
+        int threads = Integer.parseInt(opts.getOrDefault("threads", "64"));
+        String queueOpt = opts.get("queue-size");
+        boolean noQueue = "none".equalsIgnoreCase(queueOpt);
+        Integer queueSize = queueOpt == null || noQueue ? null : 
Integer.valueOf(queueOpt);
+        int warmup = Integer.parseInt(opts.getOrDefault("warmup", "5"));
+
+        System.out.printf(Locale.ROOT, "clients=%d calls=%d io-ms=%d 
threads=%d queue-size=%s warmup=%d%n",
+                          clients, calls, ioMs, threads, noQueue ? "none" : 
queueSize, warmup);
+        if (mode.equals("off") || mode.equals("both")) {
+            runPhase(false, clients, calls, ioMs, threads, queueSize, noQueue, 
warmup);
+        }
+        if (mode.equals("on") || mode.equals("both")) {
+            runPhase(true, clients, calls, ioMs, threads, queueSize, noQueue, 
warmup);
+        }
+    }
+
+    private static void runPhase(boolean virtual, int clients, int calls, long 
ioMs, int threads,
+                                 Integer queueSize, boolean noQueue, int 
warmup) throws Exception {
+        Map<String, Object> conf = new HashMap<>(Utils.readDefaultConfig());
+        conf.put(Config.STORM_THRIFT_TRANSPORT_PLUGIN, 
SimpleTransportPlugin.class.getName());
+        conf.put(Config.NIMBUS_THRIFT_PORT, 0);
+        conf.put(Config.NIMBUS_THRIFT_THREADS, threads);
+        conf.put(Config.STORM_VIRTUAL_THREADS_ENABLED, virtual);
+        conf.put(Config.STORM_NIMBUS_RETRY_TIMES, 0);
+        if (noQueue) {
+            conf.remove(Config.NIMBUS_QUEUE_SIZE);
+        } else if (queueSize != null) {
+            conf.put(Config.NIMBUS_QUEUE_SIZE, queueSize);
+        }
+
+        AtomicInteger inFlight = new AtomicInteger();
+        AtomicInteger peakInFlight = new AtomicInteger();
+        Nimbus.Iface handler = sleepingHandler(ioMs, inFlight, peakInFlight);
+        ThriftServer server = new ThriftServer(conf, new 
Nimbus.Processor<>(handler), ThriftConnectionType.NIMBUS);
+        Thread serveThread = new Thread(server::serve, "bench-thrift-serve");
+        serveThread.setDaemon(true);
+        serveThread.start();
+        while (!server.isServing()) {

Review Comment:
   No timeout, if `serve()` throws this spins forever.



##########
storm-client/test/jvm/org/apache/storm/daemon/worker/WorkerStateTest.java:
##########
@@ -138,4 +140,29 @@ public void 
testTransferLocalBatchDropsTuplesForUnknownTasks() throws TException
             ConfigUtils.setInstance(previousConfigUtils);
         }
     }
+
+    @Test
+    public void sharedExecutorUsesVirtualThreadsWhenEnabled() throws Exception 
{
+        Map<String, Object> conf = new HashMap<>();
+        conf.put(Config.TOPOLOGY_WORKER_SHARED_THREAD_POOL_SIZE, 2);
+        conf.put(Config.STORM_VIRTUAL_THREADS_ENABLED, true);
+        ExecutorService pool = WorkerState.makeSharedExecutor(conf);

Review Comment:
   Calls the static helper directly, so the wiring in `makeDefaultResources` 
(which conf is passed) is not tested. Reverting that line to `conf` would still 
pass.



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

Reply via email to