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 261ed247388 [Spark 4] Reuse the legacy maxRecordsPerBatch option in
the Structured Streaming runner (#39952)
261ed247388 is described below
commit 261ed2473888a2716f1bed8d5e564bbe5f511b45
Author: Tobias Kaymak <[email protected]>
AuthorDate: Tue Sep 1 16:35:30 2026 +0200
[Spark 4] Reuse the legacy maxRecordsPerBatch option in the Structured
Streaming runner (#39952)
---
.../SparkStructuredStreamingPipelineOptions.java | 15 +++---
...parkStructuredStreamingPipelineOptionsTest.java | 57 ++++++++++++++++++++++
2 files changed, 63 insertions(+), 9 deletions(-)
diff --git
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java
index 29cc4cb99cf..391350fd348 100644
---
a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java
+++
b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java
@@ -48,16 +48,13 @@ public interface SparkStructuredStreamingPipelineOptions
extends SparkCommonPipe
void setWatermarkDelayMillis(long value);
- // Note: deliberately NOT named getMaxRecordsPerBatch. The legacy Spark
runner's
- // SparkPipelineOptions already declares Long getMaxRecordsPerBatch(); a
same-name getter with a
- // different return type breaks proxy generation for every registered
PipelineOptions interface.
- @Description(
- "Maximum number of records to read per micro-batch from a streaming
source "
- + "(streaming mode only).")
- @Default.Integer(1000)
- int getMaxRecordsPerMicroBatch();
+ // Mirrors the legacy SparkPipelineOptions declaration exactly, so users
migrating from the
+ // legacy runner keep the same flag.
+ @Description("Max records per micro-batch. For streaming sources only.")
+ @Default.Long(-1)
+ Long getMaxRecordsPerBatch();
- void setMaxRecordsPerMicroBatch(int value);
+ void setMaxRecordsPerBatch(Long maxRecordsPerBatch);
@Description(
"Maximum duration in milliseconds of a micro-batch trigger interval
(streaming mode only).")
diff --git
a/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptionsTest.java
b/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptionsTest.java
new file mode 100644
index 00000000000..48c13f0ad57
--- /dev/null
+++
b/runners/spark/src/test/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptionsTest.java
@@ -0,0 +1,57 @@
+/*
+ * 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;
+
+import static org.junit.Assert.assertEquals;
+
+import org.apache.beam.runners.spark.SparkPipelineOptions;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.PipelineOptionsFactory;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/** Tests for {@link SparkStructuredStreamingPipelineOptions}. */
+@RunWith(JUnit4.class)
+public class SparkStructuredStreamingPipelineOptionsTest {
+
+ /**
+ * {@code maxRecordsPerBatch} is declared by both {@link
SparkPipelineOptions} and {@link
+ * SparkStructuredStreamingPipelineOptions} with identical signatures, so
one options proxy serves
+ * both and the flag carries over when migrating between the runners.
+ */
+ @Test
+ public void maxRecordsPerBatchIsSharedWithLegacyOptions() {
+ PipelineOptions options = PipelineOptionsFactory.create();
+
+ SparkStructuredStreamingPipelineOptions streamingOptions =
+ options.as(SparkStructuredStreamingPipelineOptions.class);
+ assertEquals(Long.valueOf(-1), streamingOptions.getMaxRecordsPerBatch());
+
+ streamingOptions.setMaxRecordsPerBatch(500L);
+ assertEquals(Long.valueOf(500),
options.as(SparkPipelineOptions.class).getMaxRecordsPerBatch());
+ }
+
+ @Test
+ public void maxRecordsPerBatchIsParsedFromArgs() {
+ SparkStructuredStreamingPipelineOptions options =
+ PipelineOptionsFactory.fromArgs("--maxRecordsPerBatch=42")
+ .as(SparkStructuredStreamingPipelineOptions.class);
+ assertEquals(Long.valueOf(42), options.getMaxRecordsPerBatch());
+ }
+}