Apache Beam is an open source unified programming model for batch and streaming data processing pipelines. Recent activity on the master branch brings initial Spark 4 stateful streaming support, key row format alignment across SDKs, and important Python runtime cleanup. These updates address runner capabilities, engine serialization fixes, and connector modernizations for data engineering teams.
Spark 4 Structured Streaming gains stateful processing ¶
Engine capabilities for Apache Spark continue to expand as the runner updates for Spark 4. In commit dfc4d667c, support landed for stateful ParDo execution in Spark Structured Streaming using the native transformWithState operator. The implementation maps Beam user state into Spark state store primitives while keeping timers synchronized per key.
The translator delegates user state to SparkStateInternals backed by Spark MapState. Timers use SparkTimerInternals stored inside a single ValueState entry, triggering one Spark wakeup per key. File updates to SparkStateInternals.java wire this storage layout directly into the state processor lifecycle.
public class BeamStatefulProcessor extends StatefulProcessor<K, V, O> {
@Override
public void handleInputRows(K key, Iterator<V> inputRows, TimerValues timerValues) {
// Process input elements and evaluate timers per key
}
}
This translation layer comes with explicit boundary constraints. Pipelines using unsupported timer domains, window expiration, time sorted inputs, or merging windows will fail early during graph translation. Related runner validation tests in commit 32d54e821 now run serially due to JVM wide metric collection, while commit 6fe144b37 fixes isEmpty evaluations on state maps after item removals.
Python SDK drops version 3.10 and fixes async execution isolation ¶
Python runtime requirements are moving forward across the codebase. In commit debfa0e60, the project dropped Python 3.10 support. Python 3.11 is now the minimum supported runtime version for Python SDK workers and pipeline submitters.
Transform isolation for asynchronous operations also received critical fixes. In commit 483f14b6d, async DoFn execution was updated to isolate asynchronous tasks by key and window. Previously, state decoding could allow asynchronous work units to share context across window boundaries under high load.
In text processing, commit 0f106b321 fixed an ingestion bug in ReadFromText. When an escape character directly preceded a record delimiter, the reader previously skipped the next delimiter and merged adjacent records. In transform utilities, commit c5d866c50 updated util.py so SortAndBatchElements allows batch weights to reach max_batch_weight instead of truncating batches prematurely.
Go SDK aligns INT16 row encoding and unblocks channel recovery ¶
Data serialization format consistency across SDK boundaries is vital for multi language pipelines. In commit 22abffca3, the Go SDK row encoder was changed to encode int16 and uint16 struct fields as 2 byte big endian values. As documented in CHANGES.md, this aligns Go row wire format with Java and Python, but it represents an update incompatible breaking change for active streaming pipelines decoding Go rows containing 16 bit integers.
func EncodeInt16(value int16, w io.Writer) error {
var data [2]byte
binary.BigEndian.PutUint16(data[:], uint16(value))
_, err := ioutilx.WriteUnsafe(w, data[:])
return err
}
Worker harness resilience also improved. In commit be4211d1d, data channel recovery logic was updated so the channel manager lock is released before executing channel close instructions. This prevents deadlock scenarios when a worker attempts to recreate failed state or data channels during network disruption.
For cross language pipelines, commit 7a5f89199 introduced propagation of pipeline level resource hints to cross language transform expansions. Resource hints configured on Go pipelines now correctly pass through expansion service calls to external worker pools.
ActiveMQ migration to Jakarta and parallel GCS file rewrites ¶
Connector dependencies and object storage utilities received significant modernizations. In commit 40852a71e, JmsIO updated to ActiveMQ 6.2.5 and migrated all package imports from javax.jms to jakarta.jms using API version 3.1.0.
Object store operations in GcsUtilV2.java were overhauled in commit e26d6bacb. Copy and move operations now execute per file rewrite requests in parallel using worker thread pools.
- Added
gcsMaxConcurrentRewritesconfiguration option inGcsOptions, defaulting to 32 concurrent rewrite threads. - Deprecated
gcsRewriteDataOpBatchLimit, which applied only to legacy utility implementations. - Added a 5 minute wait loop on failure to clean up in flight rewrite tasks safely.
What to watch ¶
Operators and pipeline developers should keep the following operational details in mind when upgrading:
- Go row encoding breaking change: Streaming Go pipelines using
int16oruint16schema fields must be drained before upgrading to avoid binary decoding errors. - Spark 4 streaming limitations: Stateful ParDo pipelines on Spark 4 Structured Streaming must avoid window merging or time sorted inputs until translator support expands.
- Python worker runtime drop: Build images and container environments running Python 3.10 must be updated to Python 3.11 before taking new SDK dependencies.