This is an automated email from the ASF dual-hosted git repository.
kennknowles 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 b544f40326e Cache and reuse shared metrics pipeline execution in
MetricsTest (#40322)
b544f40326e is described below
commit b544f40326e35ae2761bd86e0e59ec07a2f1bc44
Author: Reuven Lax <[email protected]>
AuthorDate: Wed Sep 30 14:03:33 2026 -0700
Cache and reuse shared metrics pipeline execution in MetricsTest (#40322)
---
.../org/apache/beam/sdk/metrics/MetricsTest.java | 157 +++++++++++----------
1 file changed, 86 insertions(+), 71 deletions(-)
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/metrics/MetricsTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/metrics/MetricsTest.java
index 2fc800adca5..fbf0ab17282 100644
--- a/sdks/java/core/src/test/java/org/apache/beam/sdk/metrics/MetricsTest.java
+++ b/sdks/java/core/src/test/java/org/apache/beam/sdk/metrics/MetricsTest.java
@@ -34,6 +34,7 @@ import java.io.IOException;
import java.io.Serializable;
import java.util.List;
import java.util.NoSuchElementException;
+import java.util.concurrent.ConcurrentHashMap;
import org.apache.beam.sdk.PipelineResult;
import org.apache.beam.sdk.coders.Coder;
import org.apache.beam.sdk.coders.VarIntCoder;
@@ -91,6 +92,9 @@ public class MetricsTest implements Serializable {
/** Shared test helpers and setup/teardown. */
public abstract static class SharedTestBase implements Serializable {
+ private static final ConcurrentHashMap<String, PipelineResult>
CACHED_METRIC_PIPELINE_RESULTS =
+ new ConcurrentHashMap<>();
+
@Rule public final transient ExpectedException thrown =
ExpectedException.none();
@Rule public final transient TestPipeline pipeline = TestPipeline.create();
@@ -101,77 +105,88 @@ public class MetricsTest implements Serializable {
}
protected PipelineResult runPipelineWithMetrics() {
- final Counter count = Metrics.counter(MetricsTest.class, "count");
- StringSet sideinputs = Metrics.stringSet(MetricsTest.class,
"sideinputs");
- final TupleTag<Integer> output1 = new TupleTag<Integer>() {};
- final TupleTag<Integer> output2 = new TupleTag<Integer>() {};
- pipeline
- .apply(Create.of(5, 8, 13))
- .apply(
- "MyStep1",
- ParDo.of(
- new DoFn<Integer, Integer>() {
- Distribution bundleDist =
Metrics.distribution(MetricsTest.class, "bundle");
-
- @StartBundle
- public void startBundle() {
- bundleDist.update(10L);
- }
-
- @SuppressWarnings("unused")
- @ProcessElement
- public void processElement(ProcessContext c) {
- Distribution values =
Metrics.distribution(MetricsTest.class, "input");
- StringSet sources = Metrics.stringSet(MetricsTest.class,
"sources");
- BoundedTrie boundedTrieSources =
- Metrics.boundedTrie(MetricsTest.class,
"boundedTrieSources");
- count.inc();
- values.update(c.element());
-
- c.output(c.element());
- c.output(c.element());
- sources.add("gcs");
- sources.add("gcs"); // repeated should appear once
- sources.add("gcs", "gcs"); // repeated should appear once
- sideinputs.add("bigtable", "spanner");
- boundedTrieSources.add(ImmutableList.of("ab_source",
"cd_source"));
- boundedTrieSources.add(ImmutableList.of("ef_source"));
- }
-
- @DoFn.FinishBundle
- public void finishBundle() {
- bundleDist.update(40L);
- }
- }))
- .apply(
- "MyStep2",
- ParDo.of(
- new DoFn<Integer, Integer>() {
- @SuppressWarnings("unused")
- @ProcessElement
- public void processElement(ProcessContext c) {
- Distribution values =
Metrics.distribution(MetricsTest.class, "input");
- Gauge gauge = Metrics.gauge(MetricsTest.class,
"my-gauge");
- StringSet sinks =
Metrics.stringSet(MetricsTest.class, "sinks");
- BoundedTrie boundedTrieSinks =
- Metrics.boundedTrie(MetricsTest.class,
"boundedTrieSinks");
- Integer element = c.element();
- count.inc();
- values.update(element);
- gauge.set(12L);
- c.output(element);
- sinks.add("bq", "kafka", "kafka"); // repeated
should appear once
- sideinputs.add("bigtable", "sql");
- boundedTrieSinks.add(ImmutableList.of("ab_sink",
"cd_sink"));
- boundedTrieSinks.add(ImmutableList.of("ef_sink"));
- c.output(output2, element);
- }
- })
- .withOutputTags(output1, TupleTagList.of(output2)));
- PipelineResult result = pipeline.run();
-
- result.waitUntilFinish();
- return result;
+ String cacheKey =
+ pipeline.getOptions().getRunner().getName()
+ + ":"
+ +
System.getProperty(TestPipeline.PROPERTY_BEAM_TEST_PIPELINE_OPTIONS, "");
+ synchronized (CACHED_METRIC_PIPELINE_RESULTS) {
+ PipelineResult cachedResult =
CACHED_METRIC_PIPELINE_RESULTS.get(cacheKey);
+ if (cachedResult != null) {
+ return cachedResult;
+ }
+ final Counter count = Metrics.counter(MetricsTest.class, "count");
+ StringSet sideinputs = Metrics.stringSet(MetricsTest.class,
"sideinputs");
+ final TupleTag<Integer> output1 = new TupleTag<Integer>() {};
+ final TupleTag<Integer> output2 = new TupleTag<Integer>() {};
+ pipeline
+ .apply(Create.of(5, 8, 13))
+ .apply(
+ "MyStep1",
+ ParDo.of(
+ new DoFn<Integer, Integer>() {
+ Distribution bundleDist =
Metrics.distribution(MetricsTest.class, "bundle");
+
+ @StartBundle
+ public void startBundle() {
+ bundleDist.update(10L);
+ }
+
+ @SuppressWarnings("unused")
+ @ProcessElement
+ public void processElement(ProcessContext c) {
+ Distribution values =
Metrics.distribution(MetricsTest.class, "input");
+ StringSet sources =
Metrics.stringSet(MetricsTest.class, "sources");
+ BoundedTrie boundedTrieSources =
+ Metrics.boundedTrie(MetricsTest.class,
"boundedTrieSources");
+ count.inc();
+ values.update(c.element());
+
+ c.output(c.element());
+ c.output(c.element());
+ sources.add("gcs");
+ sources.add("gcs"); // repeated should appear once
+ sources.add("gcs", "gcs"); // repeated should appear
once
+ sideinputs.add("bigtable", "spanner");
+ boundedTrieSources.add(ImmutableList.of("ab_source",
"cd_source"));
+ boundedTrieSources.add(ImmutableList.of("ef_source"));
+ }
+
+ @DoFn.FinishBundle
+ public void finishBundle() {
+ bundleDist.update(40L);
+ }
+ }))
+ .apply(
+ "MyStep2",
+ ParDo.of(
+ new DoFn<Integer, Integer>() {
+ @SuppressWarnings("unused")
+ @ProcessElement
+ public void processElement(ProcessContext c) {
+ Distribution values =
Metrics.distribution(MetricsTest.class, "input");
+ Gauge gauge = Metrics.gauge(MetricsTest.class,
"my-gauge");
+ StringSet sinks =
Metrics.stringSet(MetricsTest.class, "sinks");
+ BoundedTrie boundedTrieSinks =
+ Metrics.boundedTrie(MetricsTest.class,
"boundedTrieSinks");
+ Integer element = c.element();
+ count.inc();
+ values.update(element);
+ gauge.set(12L);
+ c.output(element);
+ sinks.add("bq", "kafka", "kafka"); // repeated
should appear once
+ sideinputs.add("bigtable", "sql");
+ boundedTrieSinks.add(ImmutableList.of("ab_sink",
"cd_sink"));
+ boundedTrieSinks.add(ImmutableList.of("ef_sink"));
+ c.output(output2, element);
+ }
+ })
+ .withOutputTags(output1, TupleTagList.of(output2)));
+ PipelineResult result = pipeline.run();
+
+ result.waitUntilFinish();
+ CACHED_METRIC_PIPELINE_RESULTS.put(cacheKey, result);
+ return result;
+ }
}
}