From cc08f3c717b188008d2f007bb543928eb79f6791 Mon Sep 17 00:00:00 2001 From: aIbrahiim Date: Thu, 16 Jul 2026 13:26:32 +0300 Subject: [PATCH 1/6] Add --add-opens test JVM args for HCatalog and Dataflow worker tests --- .../trigger_files/beam_PostCommit_Java_Hadoop_Versions.json | 4 ++-- runners/google-cloud-dataflow-java/worker/build.gradle | 4 ++++ sdks/java/io/hcatalog/build.gradle | 4 ++++ 3 files changed, 10 insertions(+), 2 deletions(-) diff --git a/.github/trigger_files/beam_PostCommit_Java_Hadoop_Versions.json b/.github/trigger_files/beam_PostCommit_Java_Hadoop_Versions.json index 1bd74515152c..f1ba03a243ee 100644 --- a/.github/trigger_files/beam_PostCommit_Java_Hadoop_Versions.json +++ b/.github/trigger_files/beam_PostCommit_Java_Hadoop_Versions.json @@ -1,4 +1,4 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run", - "modification": 4 -} \ No newline at end of file + "modification": 5 +} diff --git a/runners/google-cloud-dataflow-java/worker/build.gradle b/runners/google-cloud-dataflow-java/worker/build.gradle index 21879861e9d6..0fd1dc1254ef 100644 --- a/runners/google-cloud-dataflow-java/worker/build.gradle +++ b/runners/google-cloud-dataflow-java/worker/build.gradle @@ -154,6 +154,10 @@ applyJavaNature( /******************************************************************************/ // Configure the worker root project +tasks.withType(Test).configureEach { + jvmArgs '--add-opens=java.base/java.lang=ALL-UNNAMED' +} + configurations { sourceFile diff --git a/sdks/java/io/hcatalog/build.gradle b/sdks/java/io/hcatalog/build.gradle index d3bdd8f10765..b44487ab6ec4 100644 --- a/sdks/java/io/hcatalog/build.gradle +++ b/sdks/java/io/hcatalog/build.gradle @@ -40,6 +40,10 @@ hadoopVersions.each {kv -> configurations.create("hadoopVersion$kv.key")} def hive_version = "4.0.1" +tasks.withType(Test).configureEach { + jvmArgs '--add-opens=java.base/java.net=ALL-UNNAMED' +} + dependencies { implementation library.java.vendored_guava_32_1_2_jre implementation project(path: ":sdks:java:core", configuration: "shadow") From d9ecd8aaa833817f87949a896344d5146a4aa018 Mon Sep 17 00:00:00 2001 From: aIbrahiim Date: Thu, 16 Jul 2026 15:12:59 +0300 Subject: [PATCH 2/6] add more --add-opens for Dataflow --- runners/google-cloud-dataflow-java/worker/build.gradle | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/build.gradle b/runners/google-cloud-dataflow-java/worker/build.gradle index 0fd1dc1254ef..4c1403a52a50 100644 --- a/runners/google-cloud-dataflow-java/worker/build.gradle +++ b/runners/google-cloud-dataflow-java/worker/build.gradle @@ -155,7 +155,10 @@ applyJavaNature( // Configure the worker root project tasks.withType(Test).configureEach { - jvmArgs '--add-opens=java.base/java.lang=ALL-UNNAMED' + jvmArgs '--add-opens=java.base/java.lang=ALL-UNNAMED', + '--add-opens=java.base/java.util=ALL-UNNAMED', + '--add-opens=java.base/java.util.concurrent=ALL-UNNAMED', + '--add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED' } configurations { From 2d8faee715272068c08b5c9d983e7bd331a04193 Mon Sep 17 00:00:00 2001 From: aIbrahiim Date: Fri, 17 Jul 2026 13:08:58 +0300 Subject: [PATCH 3/6] document --add-opens for WindmillStateTestUtils cache --- runners/google-cloud-dataflow-java/worker/build.gradle | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/runners/google-cloud-dataflow-java/worker/build.gradle b/runners/google-cloud-dataflow-java/worker/build.gradle index 4c1403a52a50..e6ba1ddcfe43 100644 --- a/runners/google-cloud-dataflow-java/worker/build.gradle +++ b/runners/google-cloud-dataflow-java/worker/build.gradle @@ -155,6 +155,10 @@ applyJavaNature( // Configure the worker root project tasks.withType(Test).configureEach { + // WindmillStateTestUtils.assertNoReference walks every object reachable from the Windmill + // state cache (Guava Cache / ConcurrentHashMap) to ensure no per-work-item WindmillStateReader + // leaks into the global cache. It uses reflection (Field.setAccessible) into JDK internals + // (e.g. Integer.value, AtomicReferenceArray.array), which requires --add-opens on Java 17+. jvmArgs '--add-opens=java.base/java.lang=ALL-UNNAMED', '--add-opens=java.base/java.util=ALL-UNNAMED', '--add-opens=java.base/java.util.concurrent=ALL-UNNAMED', From 76cd13f9a827a6a66c4335db6951f31910b17c73 Mon Sep 17 00:00:00 2001 From: aIbrahiim Date: Mon, 20 Jul 2026 13:19:03 +0300 Subject: [PATCH 4/6] '--add-opens=java.base/java.util.concurrent.locks=ALL-UNNAMED' --- runners/google-cloud-dataflow-java/worker/build.gradle | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/build.gradle b/runners/google-cloud-dataflow-java/worker/build.gradle index e6ba1ddcfe43..40096a0566ce 100644 --- a/runners/google-cloud-dataflow-java/worker/build.gradle +++ b/runners/google-cloud-dataflow-java/worker/build.gradle @@ -158,11 +158,13 @@ tasks.withType(Test).configureEach { // WindmillStateTestUtils.assertNoReference walks every object reachable from the Windmill // state cache (Guava Cache / ConcurrentHashMap) to ensure no per-work-item WindmillStateReader // leaks into the global cache. It uses reflection (Field.setAccessible) into JDK internals - // (e.g. Integer.value, AtomicReferenceArray.array), which requires --add-opens on Java 17+. + // (e.g. Integer.value, AtomicReferenceArray.array, ReentrantLock.sync), which requires + // --add-opens on Java 17+. jvmArgs '--add-opens=java.base/java.lang=ALL-UNNAMED', '--add-opens=java.base/java.util=ALL-UNNAMED', '--add-opens=java.base/java.util.concurrent=ALL-UNNAMED', - '--add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED' + '--add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED', + '--add-opens=java.base/java.util.concurrent.locks=ALL-UNNAMED' } configurations { From f21520093bd173aecbb042e8f2447c59c28e40d8 Mon Sep 17 00:00:00 2001 From: aIbrahiim Date: Mon, 20 Jul 2026 14:44:38 +0300 Subject: [PATCH 5/6] Add remaining --add-opens for WindmillStateTestUtils --- runners/google-cloud-dataflow-java/worker/build.gradle | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/build.gradle b/runners/google-cloud-dataflow-java/worker/build.gradle index 40096a0566ce..44a2d40f944c 100644 --- a/runners/google-cloud-dataflow-java/worker/build.gradle +++ b/runners/google-cloud-dataflow-java/worker/build.gradle @@ -158,9 +158,12 @@ tasks.withType(Test).configureEach { // WindmillStateTestUtils.assertNoReference walks every object reachable from the Windmill // state cache (Guava Cache / ConcurrentHashMap) to ensure no per-work-item WindmillStateReader // leaks into the global cache. It uses reflection (Field.setAccessible) into JDK internals - // (e.g. Integer.value, AtomicReferenceArray.array, ReentrantLock.sync), which requires - // --add-opens on Java 17+. + // (e.g. Integer.value, AtomicReferenceArray.array, ReentrantLock.sync, ReferenceQueue.head), + // which requires --add-opens on Java 17+. jvmArgs '--add-opens=java.base/java.lang=ALL-UNNAMED', + '--add-opens=java.base/java.lang.ref=ALL-UNNAMED', + '--add-opens=java.base/java.lang.reflect=ALL-UNNAMED', + '--add-opens=java.base/java.io=ALL-UNNAMED', '--add-opens=java.base/java.util=ALL-UNNAMED', '--add-opens=java.base/java.util.concurrent=ALL-UNNAMED', '--add-opens=java.base/java.util.concurrent.atomic=ALL-UNNAMED', From aec3da259f75679a33522ffab93337b9e1dc0fef Mon Sep 17 00:00:00 2001 From: aIbrahiim Date: Mon, 20 Jul 2026 15:57:11 +0300 Subject: [PATCH 6/6] fix worker tests for JDK 21 Thread.toString format --- .../worker/DataflowOperationContextTest.java | 6 ++++-- .../worker/status/ThreadzServletTest.java | 15 ++++++++++----- 2 files changed, 14 insertions(+), 7 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOperationContextTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOperationContextTest.java index 34c3b3d5373c..6692f06d75d4 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOperationContextTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/DataflowOperationContextTest.java @@ -305,11 +305,13 @@ private void verifyLullLog(boolean hasFullThreadDump) throws IOException { String infoLines = Joiner.on("\n").join(Iterables.filter(lines, line -> line.contains("\"INFO\""))); + // Match on the thread name rather than the full Thread.toString() prefix: JDK 21+ inserts + // the thread id (Thread[#51,backgroundThread,...] vs Thread[backgroundThread,...]). if (hasFullThreadDump) { assertThat( infoLines, Matchers.allOf( - Matchers.containsString("Thread[backgroundThread,"), + Matchers.containsString("backgroundThread,"), Matchers.containsString( "org.apache.beam.runners.dataflow.worker.DataflowOperationContext"), Matchers.not(Matchers.containsString(SimpleDoFnRunner.class.getName())))); @@ -318,7 +320,7 @@ private void verifyLullLog(boolean hasFullThreadDump) throws IOException { infoLines, Matchers.not( Matchers.anyOf( - Matchers.containsString("Thread[backgroundThread,"), + Matchers.containsString("backgroundThread,"), Matchers.containsString( "org.apache.beam.runners.dataflow.worker.DataflowOperationContext")))); } diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/status/ThreadzServletTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/status/ThreadzServletTest.java index 7737f7d405b5..1c2352954bd0 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/status/ThreadzServletTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/status/ThreadzServletTest.java @@ -36,13 +36,18 @@ public class ThreadzServletTest { @Test public void testDeduping() throws Exception { + // Use Thread.toString() rather than hard-coded strings: JDK 21+ includes the thread id + // (e.g. Thread[#42,Thread1,5,main] vs Thread[Thread1,5,main]). + Thread thread1 = new Thread("Thread1"); + Thread thread2 = new Thread("Thread2"); + Thread thread3 = new Thread("Thread3"); Map stacks = ImmutableMap.of( - new Thread("Thread1"), + thread1, new StackTraceElement[] {new StackTraceElement("Class", "Method1", "File", 11)}, - new Thread("Thread2"), + thread2, new StackTraceElement[] {new StackTraceElement("Class", "Method1", "File", 11)}, - new Thread("Thread3"), + thread3, new StackTraceElement[] {new StackTraceElement("Class", "Method2", "File", 17)}); Map> deduped = ThreadzServlet.deduplicateThreadStacks(stacks); @@ -54,13 +59,13 @@ public void testDeduping() throws Exception { new Stack( new StackTraceElement[] {new StackTraceElement("Class", "Method1", "File", 11)}, Thread.State.NEW), - Arrays.asList("Thread[Thread1,5,main]", "Thread[Thread2,5,main]"))); + Arrays.asList(thread1.toString(), thread2.toString()))); assertThat( deduped, Matchers.hasEntry( new Stack( new StackTraceElement[] {new StackTraceElement("Class", "Method2", "File", 17)}, Thread.State.NEW), - Arrays.asList("Thread[Thread3,5,main]"))); + Arrays.asList(thread3.toString()))); } }