Skip to content

Add an opt-in virtual threads flag for blocking I/O thread pools on Java 25 - #9097

Open
GGraziadei wants to merge 1 commit into
apache:masterfrom
GGraziadei:virtual-threads-java25
Open

GGraziadei wants to merge 1 commit into
apache:masterfrom
GGraziadei:virtual-threads-java25

Conversation

@GGraziadei

Copy link
Copy Markdown
Member

What is the purpose of the change

Now that master requires Java 25, this change lets operators opt in to Java virtual threads for Storm's blocking, I/O-bound thread pools through a new cluster setting, storm.virtual.threads.enabled (default false). With the flag off nothing changes apart from thread names.

A new helper, org.apache.storm.utils.StormThreadFactory, returns a virtual-thread ThreadFactory when the flag is on and a non-daemon, normal-priority platform-thread factory otherwise. The factory is injected into the existing pools, so core sizes, queues, rejection policy and error handling are untouched and the *.threads settings keep their meaning as a concurrency bound:

  • Thrift server handlers for the SASL, TLS and simple transports (Nimbus, Supervisor and DRPC). With the simple transport an executor is now always supplied when the flag is on; without a configured queue size it uses the same unbounded queue THsHaServer builds by default.
  • AsyncLocalizer download and task executors.
  • Nimbus AssignmentDistributionService and Supervisor heartbeat pools.
  • DRPCSpout background executor.
  • The worker shared executor exposed through TopologyContext, read from the merged topology conf so a topology can override the flag.

Spout/bolt executor threads, worker transfer, JCQueue, Netty event loops and timers are deliberately left on platform threads: they are busy-poll hot loops that would lose throughput on virtual threads.

The javadoc of the new key documents the operational caveats: all virtual threads share one carrier scheduler sized to the core count, blocking file I/O (blob downloads) occupies a carrier, and with the simple transport and no configured queue size the effective handler concurrency goes from THsHaServer's fallback of 5 to *.threads.

What the benchmark shows

ThriftHandlerVirtualThreadsBench (added to examples/storm-perf) runs an in-process Nimbus Thrift server whose handler simulates blocking I/O and hammers it with N clients, each mode in its own JVM. On a 20-core machine:

clients / *.threads / io-ms flag calls/s p50 ms p99 ms server OS threads RSS after
200 / 64 / 10 off 6308 30.6 37.7 64 186 MB
200 / 64 / 10 on 6105 30.7 41.0 ~23 175 MB
500 / 512 / 20 off 21401 21.0 32.1 512 268 MB
500 / 512 / 20 on 19946 22.5 35.4 ~23 290 MB
2000 / 2000 / 50 off 28525 63.4 88.6 2000 531 MB
2000 / 2000 / 50 on 27257 61.2 121.9 ~23 424 MB

At equal *.threads the throughput is unchanged (within a few percent) because the workload is bounded by the configured concurrency either way. The gain is in resources: server-side OS threads drop from *.threads to roughly the number of cores, and RSS drops at high concurrency. The tail latency grows slightly at very high concurrency. In short, this is a resource saving on the control plane that lets operators raise nimbus.thrift.threads / supervisor.thrift.threads to absorb bursts without paying for native threads; with the shipped defaults it is neutral, and it is off by default.

How was the change tested

  • New unit tests: StormThreadFactoryTest (flag on/off/missing/unexpected type, thread naming, forced NORM_PRIORITY), WorkerStateTest (shared executor on virtual vs platform threads), AsyncLocalizerTest (download executor on virtual vs platform threads).
  • New round-trip tests in AuthTest through a real ThriftServer and NimbusClient, asserting from inside the Nimbus handler that it runs on a virtual thread when the flag is on (simple transport with and without a queue size, digest SASL transport) and on a platform thread by default; the digest case also checks ReqContext still carries the principal on a virtual handler thread.
  • Full storm-client and storm-server suites pass with the flag off (default): 650 and 487 tests respectively, 0 failures.
  • Smoke test with the flag forced on for the whole JVM via -Dstorm.options=storm.virtual.threads.enabled=true: LocalNimbusTest passes and the "platform thread by default" assertion in AuthTest fails as expected, proving the flag reaches the handler pool.
  • Benchmark runs above.

@GGraziadei
GGraziadei requested review from reiabreu and rzo1 September 20, 2026 12:02
@reiabreu

Copy link
Copy Markdown
Contributor

Currently traveling and unable to properly review this. Happy to do it once I'm able to

@reiabreu

Copy link
Copy Markdown
Contributor

This comment was drafted with the help of an LLM (Claude) and reviewed by me before posting.

Reviewed this in detail — two real behavioral side effects worth addressing before merge, plus some smaller items. Ordered by importance:

1. DRPC_INVOCATIONS concurrency silently jumps 5→64 when the flag is enabled

DRPC_INVOCATIONS has no queue-size config (hardcoded null in ThriftConnectionType), so the new guard if (queueSize != null || StormThreadFactory.isVirtualEnabled(topoConf)) is always true for this type once the flag is on. That forces an explicit ThreadPoolExecutor(64, 64, ...) where before it fell back to THsHaServer's own pool, capped at 5 by corePoolSize (unbounded queue means maxWorkerThreads never actually kicks in). For DRPC_INVOCATIONS this isn't an edge case — it's a permanent, unconditional 13x concurrency jump, not just a thread-implementation swap. Worth documenting explicitly, or capping the new executor at the old effective concurrency for connection types with no queue-size config.

2. Worker shared-executor pool size becomes topology-overridable

// before: bare `conf` in this instance method resolves to this.conf (daemon-level config)
int threadPoolSize = ObjectReader.getInt(conf.get(Config.TOPOLOGY_WORKER_SHARED_THREAD_POOL_SIZE));
// after: explicitly passes this.topologyConf (merged: daemon conf + topology's own overrides)
return ImmutableMap.of(WorkerTopologyContext.SHARED_EXECUTOR, makeSharedExecutor(topologyConf));

Pre-PR, topology.worker.shared.thread.pool.size only ever read from the daemon/cluster conf. Post-PR it reads from the merged topology conf, silently letting a topology submitter override a pool size an operator may be relying on as a cluster-wide cap (TestUtilsForWorkerState.java needed a new entry added to its topologyConf map, or the test throws — confirms the read source really moved). Probably fine as a side effect of making the flag itself topology-overridable (which is documented), but worth confirming that's deliberate — it's an undocumented trust-boundary change on multi-tenant clusters.

3. Virtual threads are always daemon threads, with no override

StormThreadFactory's platform branch explicitly sets .daemon(false); the virtual branch can't — Thread.Builder.OfVirtual has no .daemon(...) method, so virtual threads are always daemon. This is a real asymmetry across all 7 pools this PR touches, but since Storm daemons shut down via explicit Utils.exitProcess()/System.exit() calls rather than relying on natural-exit semantics, the practical impact may be minimal — System.exit() doesn't wait for non-daemon threads either, unless a shutdown hook explicitly joins them. Worth a quick check for any addShutdownHook usage around these pools before dismissing it, and a one-line doc-comment note either way.

4. Benchmark's proxy handler misroutes Object methods (ThriftHandlerVirtualThreadsBench.sleepingHandler)

if (method.getDeclaringClass() == Object.class) {
    return method.invoke(inFlight, margs);   // should be `proxy`, not `inFlight`
}

equals/hashCode/toString get invoked on the captured AtomicInteger counter instead of the proxy. Harmless today (nothing calls these on the proxy), confined to the new example/benchmark file — just a copy-paste slip worth a one-line fix.

5. AsyncLocalizer thread names changed format

"AsyncLocalizer Download Executor - 0" → "AsyncLocalizer-Download-Executor-0". No functional effect found in-repo, but worth a heads-up for anyone with external log/metrics tooling matching the old literal name.

Smaller cleanup, non-blocking:

  • StormThreadFactory.isVirtualEnabled() writes its own boolean-parsing logic instead of using ObjectReader.getBoolean, which is already used for every other @IsBoolean config key in this codebase.
  • type.name().toLowerCase(Locale.ROOT) + "-handler" is copy-pasted identically across SimpleTransportPlugin/SaslTransportPlugin/TlsTransportPlugin — could be one method on ThriftConnectionType, consistent with how every other per-type value is already sourced.
  • AssignmentDistributionService's new @SuppressWarnings("unchecked") cast could be avoided by having StormThreadFactory accept a raw Map.
  • The four new virtual-thread tests in AuthTest.java are near-identical copy-paste — could collapse into one @ParameterizedTest.

Introduce storm.virtual.threads.enabled (default false) and
StormThreadFactory, which returns a virtual-thread ThreadFactory when the
flag is on and a non-daemon, normal-priority platform-thread factory
otherwise. Thread names are <prefix>-<n> in both modes.

The factory is injected into the existing pools, so core sizes, queues,
rejection policy and error handling are unchanged and the *.threads
settings keep their meaning as a concurrency bound:

- Thrift server handlers for the SASL, TLS and simple transports
  (Nimbus, Supervisor and DRPC). With the simple transport an executor is
  now always supplied when the flag is on; without a configured queue
  size it uses the same unbounded queue THsHaServer builds by default.
- AsyncLocalizer download and task executors.
- Nimbus AssignmentDistributionService and Supervisor heartbeat pools.
- DRPCSpout background executor.
- Worker shared executor exposed through TopologyContext, read from the
  merged topology conf so a topology can override the flag.

Spout/bolt executor threads, worker transfer, JCQueue, Netty and timer
threads are deliberately left on platform threads.

Round-trip tests through NimbusClient assert that Thrift handlers run on
virtual threads when the flag is on and on platform threads by default;
unit tests cover the factory, the localizer executor and the worker
shared executor.

ThriftHandlerVirtualThreadsBench in storm-perf measures the flag against
an in-process Nimbus Thrift server whose handler simulates blocking I/O.
At equal *.threads the throughput is unchanged, while server-side OS
threads drop from *.threads to roughly the number of cores.
@rzo1
rzo1 force-pushed the virtual-threads-java25 branch from f6547f1 to 94a4adc Compare October 1, 2026 11:16
@GGraziadei

Copy link
Copy Markdown
Member Author

Hi @reiabreu, thanks for the review. I will fix asap.

@rzo1 rzo1 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for this. With the flag off it is indeed only thread names. I agree with @reiabreu's points, all of them are still open on the current head, so I only add what is not covered there. Two corrections to that list though:

  • Reading topology.worker.shared.thread.pool.size from the topology conf is not a trust boundary change. The pool lives in the submitter's own worker, which runs their code anyway. It is rather a fix for a topology.* key that was silently ignored before. Still worth a line in the description.
  • ObjectReader.getBoolean throws on a String, so switching to it would break the -Dstorm.options=...=true case. See inline.

On the design: one flag covers the Thrift handlers, the localizer and pools a topology controls. Given the thread dump issue below I would prefer separate switches, or at least start with the Thrift pools only. Open for discussion.

The benchmark only covers the simple transport (THsHaServer). SASL/TLS use TThreadPoolServer with a thread per connection, which is the setup most secure clusters run, and that is not measured. Also, rows 2 and 3 are well below min(clients, threads) * 1000 / io-ms (86% / 71% with the flag off), so they are not bounded by concurrency as stated, and +37% p99 at 2000 is more than "slightly". Single runs without variance, so I would not draw conclusions either way. The server OS thread column is not printed by the tool; I assume it is peak minus client threads, please say so.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

* 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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

return (Boolean) value;
}
if (value instanceof String) {
return Boolean.parseBoolean((String) value);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

} No newline at end of file

@Test
public void simpleTransportRunsHandlerOnVirtualThreadWhenEnabled() throws Exception {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Thread serveThread = new Thread(server::serve, "bench-thrift-serve");
serveThread.setDaemon(true);
serveThread.start();
while (!server.isServing()) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No timeout, if serve() throws this spins forever.

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;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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


public static void main(String[] args) throws Exception {
Map<String, String> opts = parseArgs(args);
String mode = opts.getOrDefault("mode", "both").toLowerCase(Locale.ROOT);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

@reiabreu

reiabreu commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

This comment was generated with the help of an LLM (Claude) and reviewed by me before posting.

Two points in my earlier review were wrong; @rzo1 is right on both:

  • Please keep isVirtualEnabled() as is. ObjectReader.getBoolean throws on a String, which would break the -Dstorm.options=...=true case.
  • The worker shared pool change is not a security concern, since the pool lives in the submitter's own worker, which already runs their code. It's a topology.* key that used to be ignored and is now honored, which is worth a line in the description.

The rest of my review stands.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants