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;
+      }
     }
   }
 

Reply via email to