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