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();
- }
-}