From 7e0f07bd8791276bbdda5c2d1556915f503b33c3 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Thu, 4 Jun 2026 18:54:52 +0000 Subject: [PATCH 1/2] Fix thread safety of HotKey logger --- .../beam/runners/dataflow/worker/HotKeyLogger.java | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/HotKeyLogger.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/HotKeyLogger.java index 00d93890d9a2..b13ccce2eb90 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/HotKeyLogger.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/HotKeyLogger.java @@ -18,6 +18,7 @@ package org.apache.beam.runners.dataflow.worker; import com.google.api.client.util.Clock; +import javax.annotation.concurrent.GuardedBy; import org.apache.beam.runners.dataflow.util.TimeUtil; import org.joda.time.Duration; import org.slf4j.Logger; @@ -33,6 +34,7 @@ public class HotKeyLogger { * The previous time the HotKeyDetection was logged. This is used to throttle logging to every 5 * minutes. */ + @GuardedBy("this") private long prevHotKeyDetectionLogMs = 0; /** Throttles logging the detection to every loggingPeriod */ @@ -83,10 +85,12 @@ public void logHotKeyDetection(String userStepName, Duration hotKeyAge, Object h protected boolean isThrottled() { // Throttle logging the HotKeyDetection to every 5 minutes. long nowMs = clock.currentTimeMillis(); - if (nowMs - prevHotKeyDetectionLogMs < loggingPeriod.getMillis()) { - return true; + synchronized(this) { + if (nowMs - prevHotKeyDetectionLogMs < loggingPeriod.getMillis()) { + return true; + } + prevHotKeyDetectionLogMs = nowMs; } - prevHotKeyDetectionLogMs = nowMs; return false; } } From d53696f548381c2d06afba3f6071f18d4ffe6b86 Mon Sep 17 00:00:00 2001 From: Arun Pandian Date: Fri, 5 Jun 2026 04:45:27 +0000 Subject: [PATCH 2/2] spotless --- .../org/apache/beam/runners/dataflow/worker/HotKeyLogger.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/HotKeyLogger.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/HotKeyLogger.java index b13ccce2eb90..c449f9d6c825 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/HotKeyLogger.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/HotKeyLogger.java @@ -85,7 +85,7 @@ public void logHotKeyDetection(String userStepName, Duration hotKeyAge, Object h protected boolean isThrottled() { // Throttle logging the HotKeyDetection to every 5 minutes. long nowMs = clock.currentTimeMillis(); - synchronized(this) { + synchronized (this) { if (nowMs - prevHotKeyDetectionLogMs < loggingPeriod.getMillis()) { return true; }