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

Abacn 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 2b42a230a72 [Spark][#36841] Prune dead DStream helpers and add 
streaming lifecycle timeout test (#40130)
2b42a230a72 is described below

commit 2b42a230a7280ee73c327d2c7136a0d2aa2db0e7
Author: Tobias Kaymak <[email protected]>
AuthorDate: Tue Sep 15 17:33:07 2026 +0200

    [Spark][#36841] Prune dead DStream helpers and add streaming lifecycle 
timeout test (#40130)
---
 .../streaming/StreamingPipelineLifecycleTest.java  | 36 ++++++++++++++
 .../StructuredStreamingPipelineStateTest.java      | 50 ++-----------------
 .../translation/streaming/SimpleSourceTest.java    | 57 ----------------------
 3 files changed, 39 insertions(+), 104 deletions(-)

diff --git 
a/runners/spark/4/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/StreamingPipelineLifecycleTest.java
 
b/runners/spark/4/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/StreamingPipelineLifecycleTest.java
index 76462c38fe1..fd897db515a 100644
--- 
a/runners/spark/4/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/StreamingPipelineLifecycleTest.java
+++ 
b/runners/spark/4/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/StreamingPipelineLifecycleTest.java
@@ -33,6 +33,7 @@ import org.apache.beam.sdk.PipelineResult;
 import org.apache.beam.sdk.io.Read;
 import org.apache.beam.sdk.transforms.DoFn;
 import org.apache.beam.sdk.transforms.ParDo;
+import org.joda.time.Duration;
 import org.junit.ClassRule;
 import org.junit.Rule;
 import org.junit.Test;
@@ -123,6 +124,41 @@ public class StreamingPipelineLifecycleTest implements 
Serializable {
     awaitActiveQueries(count -> count == 0, "the streaming query did not stop 
after cancel");
   }
 
+  @Test
+  public void timeoutKeepsRunningState() throws Exception {
+    String tag = "lifecycle-timeout";
+    String collectorId = StreamingTestUtils.newCollectorId(tag);
+
+    SparkStructuredStreamingPipelineOptions options =
+        StreamingTestUtils.streamingOptions(checkpointDir);
+    options.setStreamingStopAfterIdleBatches(-1);
+    Pipeline pipeline = Pipeline.create(options);
+
+    pipeline
+        .apply("ReadUnbounded", Read.from(new TestUnboundedSource(tag, 1, 10)))
+        .apply("Collect", ParDo.of(new 
StreamingTestUtils.CollectDoFn<>(collectorId)));
+
+    PipelineResult result = pipeline.run();
+    try {
+      assertEquals(PipelineResult.State.RUNNING, result.getState());
+
+      awaitActiveQueries(count -> count > 0, "no streaming query started");
+
+      PipelineResult.State stateAfterTimeout = 
result.waitUntilFinish(Duration.millis(1));
+      assertEquals(PipelineResult.State.RUNNING, stateAfterTimeout);
+      assertEquals(PipelineResult.State.RUNNING, result.getState());
+
+      PipelineResult.State cancelledState = result.cancel();
+      assertEquals(PipelineResult.State.CANCELLED, cancelledState);
+      assertEquals(PipelineResult.State.CANCELLED, result.getState());
+    } finally {
+      if (result.getState() == PipelineResult.State.RUNNING) {
+        result.cancel();
+      }
+      awaitActiveQueries(count -> count == 0, "the streaming query did not 
stop after cancel");
+    }
+  }
+
   /** A failure in any leaf query surfaces through waitUntilFinish. */
   @Test
   public void failingLeafQueryFailsThePipelineAndStopsHealthySibling() throws 
Exception {
diff --git 
a/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/StructuredStreamingPipelineStateTest.java
 
b/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/StructuredStreamingPipelineStateTest.java
index 647c4334ff1..b883af5eb00 100644
--- 
a/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/StructuredStreamingPipelineStateTest.java
+++ 
b/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/StructuredStreamingPipelineStateTest.java
@@ -27,7 +27,6 @@ import static org.junit.Assert.fail;
 import java.io.Serializable;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.TimeUnit;
-import org.apache.beam.runners.spark.io.CreateStream;
 import 
org.apache.beam.runners.spark.structuredstreaming.translation.SparkSessionFactory;
 import org.apache.beam.sdk.Pipeline;
 import org.apache.beam.sdk.PipelineResult;
@@ -36,15 +35,11 @@ import org.apache.beam.sdk.options.PipelineOptionsFactory;
 import org.apache.beam.sdk.transforms.Create;
 import org.apache.beam.sdk.transforms.DoFn;
 import org.apache.beam.sdk.transforms.MapElements;
-import org.apache.beam.sdk.transforms.PTransform;
 import org.apache.beam.sdk.transforms.ParDo;
 import org.apache.beam.sdk.transforms.SimpleFunction;
-import org.apache.beam.sdk.values.PBegin;
-import org.apache.beam.sdk.values.PCollection;
 import org.apache.spark.TaskContext;
 import org.apache.spark.sql.SparkSession;
 import org.joda.time.Duration;
-import org.junit.Ignore;
 import org.junit.Rule;
 import org.junit.Test;
 import org.junit.rules.TestName;
@@ -107,21 +102,6 @@ public class StructuredStreamingPipelineStateTest 
implements Serializable {
         });
   }
 
-  private PTransform<PBegin, PCollection<String>> getValues(
-      final SparkStructuredStreamingPipelineOptions options) {
-    final boolean doNotSyncWithWatermark = false;
-    return options.isStreaming()
-        ? CreateStream.of(StringUtf8Coder.of(), Duration.millis(1), 
doNotSyncWithWatermark)
-            .nextBatch("one", "two")
-        : Create.of("one", "two");
-  }
-
-  private SparkStructuredStreamingPipelineOptions getStreamingOptions() {
-    options.setRunner(SparkStructuredStreamingRunner.class);
-    options.setStreaming(true);
-    return options;
-  }
-
   private SparkStructuredStreamingPipelineOptions getBatchOptions() {
     options.setRunner(SparkStructuredStreamingRunner.class);
     options.setStreaming(false); // explicit because options is reused 
throughout the test.
@@ -131,9 +111,9 @@ public class StructuredStreamingPipelineStateTest 
implements Serializable {
   private Pipeline getPipeline(final SparkStructuredStreamingPipelineOptions 
options) {
 
     final Pipeline pipeline = Pipeline.create(options);
-    final String name = testName.getMethodName() + "(isStreaming=" + 
options.isStreaming() + ")";
+    final String name = testName.getMethodName();
 
-    
pipeline.apply(getValues(options)).setCoder(StringUtf8Coder.of()).apply(printParDo(name));
+    pipeline.apply(Create.of("one", 
"two")).setCoder(StringUtf8Coder.of()).apply(printParDo(name));
 
     return pipeline;
   }
@@ -146,7 +126,7 @@ public class StructuredStreamingPipelineStateTest 
implements Serializable {
     try {
       final Pipeline pipeline = Pipeline.create(options);
       pipeline
-          .apply(getValues(options))
+          .apply(Create.of("one", "two"))
           .setCoder(StringUtf8Coder.of())
           .apply(
               MapElements.via(
@@ -216,45 +196,21 @@ public class StructuredStreamingPipelineStateTest 
implements Serializable {
     assertThat(result.waitUntilFinish(), is(PipelineResult.State.CANCELLED));
   }
 
-  @Ignore("TODO: Reactivate with streaming.")
-  @Test
-  public void testStreamingPipelineRunningState() throws Exception {
-    testRunningPipeline(getStreamingOptions());
-  }
-
   @Test
   public void testBatchPipelineRunningState() throws Exception {
     testRunningPipeline(getBatchOptions());
   }
 
-  @Ignore("TODO: Reactivate with streaming.")
-  @Test
-  public void testStreamingPipelineCanceledState() throws Exception {
-    testCanceledPipeline(getStreamingOptions());
-  }
-
   @Test
   public void testBatchPipelineCanceledState() throws Exception {
     testCanceledPipeline(getBatchOptions());
   }
 
-  @Ignore("TODO: Reactivate with streaming.")
-  @Test
-  public void testStreamingPipelineFailedState() throws Exception {
-    testFailedPipeline(getStreamingOptions());
-  }
-
   @Test
   public void testBatchPipelineFailedState() throws Exception {
     testFailedPipeline(getBatchOptions());
   }
 
-  @Ignore("TODO: Reactivate with streaming.")
-  @Test
-  public void testStreamingPipelineTimeoutState() throws Exception {
-    testTimeoutPipeline(getStreamingOptions());
-  }
-
   @Test
   public void testBatchPipelineTimeoutState() throws Exception {
     testTimeoutPipeline(getBatchOptions());
diff --git 
a/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/SimpleSourceTest.java
 
b/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/SimpleSourceTest.java
deleted file mode 100644
index a06d2cec1e9..00000000000
--- 
a/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/SimpleSourceTest.java
+++ /dev/null
@@ -1,57 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements.  See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership.  The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package 
org.apache.beam.runners.spark.structuredstreaming.translation.streaming;
-
-import java.io.Serializable;
-import 
org.apache.beam.runners.spark.structuredstreaming.SparkStructuredStreamingPipelineOptions;
-import 
org.apache.beam.runners.spark.structuredstreaming.SparkStructuredStreamingRunner;
-import org.apache.beam.sdk.Pipeline;
-import org.apache.beam.sdk.io.GenerateSequence;
-import org.apache.beam.sdk.options.PipelineOptionsFactory;
-import org.junit.BeforeClass;
-import org.junit.ClassRule;
-import org.junit.Ignore;
-import org.junit.Test;
-import org.junit.rules.TemporaryFolder;
-import org.junit.runner.RunWith;
-import org.junit.runners.JUnit4;
-
-/** Test class for beam to spark source translation. */
-@RunWith(JUnit4.class)
-public class SimpleSourceTest implements Serializable {
-  private static Pipeline pipeline;
-  @ClassRule public static final TemporaryFolder TEMPORARY_FOLDER = new 
TemporaryFolder();
-
-  @BeforeClass
-  public static void beforeClass() {
-    SparkStructuredStreamingPipelineOptions options =
-        
PipelineOptionsFactory.create().as(SparkStructuredStreamingPipelineOptions.class);
-    options.setRunner(SparkStructuredStreamingRunner.class);
-    options.setTestMode(true);
-    pipeline = Pipeline.create(options);
-  }
-
-  @Ignore
-  @Test
-  public void testUnboundedSource() {
-    // produces an unbounded PCollection of longs from 0 to Long.MAX_VALUE 
which elements
-    // have processing time as event timestamps.
-    pipeline.apply(GenerateSequence.from(0L));
-    pipeline.run();
-  }
-}

Reply via email to