This is an automated email from the ASF dual-hosted git repository.

shunping pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 07d3169e4e0 Refactor load testing infra for upcoming GCS load tests 
(#40223)
07d3169e4e0 is described below

commit 07d3169e4e01815368ae992366cba4ad190a1fc9
Author: Shunping Huang <[email protected]>
AuthorDate: Tue Sep 22 14:40:28 2026 -0400

    Refactor load testing infra for upcoming GCS load tests (#40223)
    
    * Rename FileBasedIOLT to TextIOLT
    
    * Use Awaitility to poll for Dataflow monitoring data instead of sleeping a 
fixed time
    
    Also fix a small bug where getDataProcessed was always queried with
    the legacy PCollection name.
    
    * Trigger load tests.
---
 .../beam_PostCommit_Java_IO_Performance_Tests.json |  3 +-
 it/common/build.gradle                             |  1 +
 .../beam/it/common/dataflow/LoadTestBase.java      | 50 +++++++++++++++++++---
 it/google-cloud-platform/build.gradle              |  2 +-
 .../storage/{FileBasedIOLT.java => TextIOLT.java}  | 10 ++---
 5 files changed, 52 insertions(+), 14 deletions(-)

diff --git 
a/.github/trigger_files/beam_PostCommit_Java_IO_Performance_Tests.json 
b/.github/trigger_files/beam_PostCommit_Java_IO_Performance_Tests.json
index 4f9719d7185..12baa399f31 100644
--- a/.github/trigger_files/beam_PostCommit_Java_IO_Performance_Tests.json
+++ b/.github/trigger_files/beam_PostCommit_Java_IO_Performance_Tests.json
@@ -1,5 +1,4 @@
 {
   "comment": "Modify this file in a trivial way to cause this test suite to 
run",
-  "modification": 3,
-  "https://github.com/apache/beam/pull/39990": "removing dead code from 
FnApiDoFnRunner"
+  "modification": 4,
 }
diff --git a/it/common/build.gradle b/it/common/build.gradle
index 62dd45ddacf..5d97597ca79 100644
--- a/it/common/build.gradle
+++ b/it/common/build.gradle
@@ -56,6 +56,7 @@ dependencies {
     implementation library.java.protobuf_java_util
     implementation library.java.protobuf_java
     implementation library.java.junit
+    testImplementation 'org.awaitility:awaitility:4.2.0'
     testImplementation library.java.mockito_inline
     testRuntimeOnly library.java.slf4j_simple
     // TODO: excluding Guava until Truth updates it to >32.1.x
diff --git 
a/it/common/src/test/java/org/apache/beam/it/common/dataflow/LoadTestBase.java 
b/it/common/src/test/java/org/apache/beam/it/common/dataflow/LoadTestBase.java
index cd1b71dc755..a669436731a 100644
--- 
a/it/common/src/test/java/org/apache/beam/it/common/dataflow/LoadTestBase.java
+++ 
b/it/common/src/test/java/org/apache/beam/it/common/dataflow/LoadTestBase.java
@@ -47,6 +47,8 @@ import org.apache.beam.it.common.TestProperties;
 import org.apache.beam.it.common.bigquery.BigQueryResourceManager;
 import org.apache.beam.it.common.monitoring.MonitoringClient;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
+import org.awaitility.Awaitility;
+import org.awaitility.core.ConditionTimeoutException;
 import org.checkerframework.checker.nullness.qual.Nullable;
 import org.junit.After;
 import org.junit.Before;
@@ -82,6 +84,12 @@ public abstract class LoadTestBase {
           "^All workers have finished the startup processes and began to 
receive work requests.*$");
   private static final Pattern WORKER_STOP_PATTERN = 
Pattern.compile("^Stopping worker pool.*$");
 
+  /** How often Cloud Monitoring is polled for the data of a job that just 
finished. */
+  private static final Duration METRICS_POLL_INTERVAL = Duration.ofSeconds(20);
+
+  /** How long Cloud Monitoring is polled at most before giving up on the data 
of a job. */
+  private static final Duration METRICS_POLL_TIMEOUT = Duration.ofMinutes(6);
+
   protected static final Credentials CREDENTIALS = 
TestProperties.googleCredentials();
   protected static final CredentialsProvider CREDENTIALS_PROVIDER =
       FixedCredentialsProvider.create(CREDENTIALS);
@@ -247,7 +255,12 @@ public abstract class LoadTestBase {
     metrics.put("ElapsedTime", monitoringClient.getElapsedTime(project, 
launchInfo));
 
     Double dataProcessed =
-        monitoringClient.getDataProcessed(project, launchInfo, 
config.inputPCollection());
+        monitoringClient.getDataProcessed(
+            project,
+            launchInfo,
+            RUNNER_V2.equals(launchInfo.runner())
+                ? config.inputPCollectionV2()
+                : config.inputPCollection());
     if (dataProcessed != null) {
       metrics.put("EstimatedDataProcessedGB", dataProcessed / 1e9d);
     }
@@ -331,11 +344,9 @@ public abstract class LoadTestBase {
       throws IOException, InterruptedException, ParseException {
     Map<String, Double> metrics = pipelineLauncher.getMetrics(project, region, 
launchInfo.jobId());
     if (launchInfo.runner().contains("Dataflow")) {
-      // monitoring metrics take up to 3 minutes to show up
-      // TODO(pranavbhandari): We should use a library like 
http://awaitility.org/ to poll for
-      // metrics instead of hard coding X minutes.
-      LOG.info("Sleeping for 4 minutes to query Dataflow runner metrics.");
-      Thread.sleep(Duration.ofMinutes(4).toMillis());
+      // Monitoring metrics take a few minutes to show up, so wait for them to 
be there instead of
+      // sleeping for a fixed amount of time.
+      waitUntilMonitoringDataAvailable(launchInfo);
       computeDataflowMetrics(metrics, launchInfo, config);
     } else if ("DirectRunner".equalsIgnoreCase(launchInfo.runner())) {
       computeDirectMetrics(metrics, launchInfo);
@@ -343,6 +354,33 @@ public abstract class LoadTestBase {
     return metrics;
   }
 
+  /** Waits until Cloud Monitoring has data for the given job. */
+  private void waitUntilMonitoringDataAvailable(LaunchInfo launchInfo) {
+    LOG.info("Waiting for the monitoring data of {} to be available.", 
launchInfo.jobId());
+    try {
+      Awaitility.await("monitoring data of " + launchInfo.jobId())
+          .atMost(METRICS_POLL_TIMEOUT)
+          .pollInterval(METRICS_POLL_INTERVAL)
+          .until(() -> monitoringDataAvailable(launchInfo));
+    } catch (ConditionTimeoutException e) {
+      LOG.warn(
+          "No monitoring data found for {} after {} minutes. The metrics of 
this job are"
+              + " incomplete.",
+          launchInfo.jobId(),
+          METRICS_POLL_TIMEOUT.toMinutes());
+    }
+  }
+
+  /** Returns whether Cloud Monitoring has data for the given job. */
+  private boolean monitoringDataAvailable(LaunchInfo launchInfo) {
+    try {
+      return monitoringClient.getElapsedTime(project, launchInfo) != null;
+    } catch (ParseException | RuntimeException e) {
+      LOG.warn("Error while querying the monitoring data of {}.", 
launchInfo.jobId(), e);
+      return false;
+    }
+  }
+
   /**
    * Computes CPU Utilization metrics of the given job.
    *
diff --git a/it/google-cloud-platform/build.gradle 
b/it/google-cloud-platform/build.gradle
index 164a75c06ba..f06669f31b0 100644
--- a/it/google-cloud-platform/build.gradle
+++ b/it/google-cloud-platform/build.gradle
@@ -79,7 +79,7 @@ dependencies {
 }
 
 tasks.register(
-        "GCSPerformanceTest", IoPerformanceTestUtilities.IoPerformanceTest, 
project, 'google-cloud-platform', 'FileBasedIOLT',
+        "GCSPerformanceTest", IoPerformanceTestUtilities.IoPerformanceTest, 
project, 'google-cloud-platform', 'TextIOLT',
         ['configuration':'large','project':'apache-beam-testing', 
'artifactBucket':'io-performance-temp']
         + System.properties
 )
diff --git 
a/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/FileBasedIOLT.java
 
b/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/TextIOLT.java
similarity index 96%
rename from 
it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/FileBasedIOLT.java
rename to 
it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/TextIOLT.java
index ac1a7fc103c..650149bbb9b 100644
--- 
a/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/FileBasedIOLT.java
+++ 
b/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/TextIOLT.java
@@ -54,19 +54,19 @@ import org.junit.Rule;
 import org.junit.Test;
 
 /**
- * FileBasedIO performance tests.
+ * TextIO performance tests.
  *
  * <p>Example trigger command for all tests:
  *
  * <pre>
- * mvn test -pl it/google-cloud-platform -am -Dtest="FileBasedIOLT" 
-Dproject=[gcpProject] \
+ * mvn test -pl it/google-cloud-platform -am -Dtest="TextIOLT" 
-Dproject=[gcpProject] \
  * -DartifactBucket=[temp bucket] -DfailIfNoTests=false
  * </pre>
  *
  * <p>Example trigger command for specific test running on direct runner:
  *
  * <pre>
- * mvn test -pl it/google-cloud-platform -am 
-Dtest="FileBasedIOLT#testTextIOWriteThenRead" \
+ * mvn test -pl it/google-cloud-platform -am 
-Dtest="TextIOLT#testTextIOWriteThenRead" \
  * -Dconfiguration=medium -Dproject=[gcpProject] -DartifactBucket=[temp 
bucket] -DfailIfNoTests=false
  * </pre>
  *
@@ -74,11 +74,11 @@ import org.junit.Test;
  *
  * <pre>mvn test -pl it/google-cloud-platform -am \
  * 
-Dconfiguration="{\"numRecords\":10000000,\"valueSizeBytes\":750,\"pipelineTimeout\":20,\"runner\":\"DataflowRunner\"}"
 \
- * -Dtest="FileBasedIOLT#testTextIOWriteThenRead" -Dconfiguration=local 
-Dproject=[gcpProject] \
+ * -Dtest="TextIOLT#testTextIOWriteThenRead" -Dproject=[gcpProject] \
  * -DartifactBucket=[temp bucket] -DfailIfNoTests=false
  * </pre>
  */
-public class FileBasedIOLT extends IOLoadTestBase {
+public class TextIOLT extends IOLoadTestBase {
 
   private static final String READ_ELEMENT_METRIC_NAME = "read_count";
 

Reply via email to