We are happy to present the new 2.77.0 release of Beam.
This release includes both improvements and new functionality.
See the download page for this release.
For more information on changes in 2.77.0, check out the detailed release notes.
Highlights
I/Os
- Added
schema_update_optionstoWriteToBigQueryfile loads, allowing BigQuery load jobs to add nullable fields or relax required fields when appending data (Python) (#21141). - BigQueryIO now supports reading BigQuery Lakehouse runtime catalog (BigLake metastore) Iceberg tables with the Storage Read API, using 4-part
project.catalog.namespace.tableidentifiers (or aTableReferencewith a compositecatalog.namespacedataset id). Previously such references were silently mis-parsed (Java) (#39597) . - SolaceIO now supports reading and writing binary and text content data payload (Java) (#39875).
- ClickHouseIO: support writing
Decimal(P, S)/Decimal32/64/128/256columns (Java) (#39840). - SolaceIO now supports reading and writing user properties (message metadata) (Java) (#40099).
- [IcebergIO] AddFiles (
IcebergAddFilesin YAML) can evolve the table schema before registering files, withschema_evolution_options,required_columns,incompatible_schema_handlingandunverifiable_file_handling(Java/YAML, batch only) (#40144).
New Features / Improvements
- (Java/Python)
Watchcan bound its deduplication state by event time, retiring an output key once the greatest emitted timestamp has moved more than the allowed lateness past it. Java addsWatch.growthOf(...).withTimestampCursor(). Python addsallowed_latenessfor the existingtimestamp_cursoroption (#18459). - (Java) Spark Structured Streaming runner: stateful ParDo with state, timers,
@RequiresTimeSortedInputand tagged outputs is now supported in batch mode (#39779). - (Python) Added support for Vertex AI Model Monitoring V2 in RunInference (#39738).
- [Flink Runner] Added opt-in static round-robin split assignment for small bounded sources via the new
sourceStaticSplitThresholdMbpipeline option. The default of 0 keeps the existing lazy pull-based assignment (#39873). - Added automatic caching of bounded, single-pane side-input views for classic Java Flink DataStream execution (#39866).
- (Python) Added
Sample.Any, the Python equivalent of Java'sSample.any, which returns up to n arbitrary elements from a PCollection (#18552).
Breaking Changes
- Portable Java SDK now encodes SchemaCoders in a portable way (#34672).
- Original custom Java coder encoding can still be obtained using StreamingOptions.setUpdateCompatibilityVersion("2.76") (#34672).
- Fixes (#36496), (#30276), (#29245).
- (Python)
TensorRTEngineHandlerNumPynow requires TensorRT 10 or later. TensorRT 8.x is no longer supported, since TensorRT 10 removed the engine binding API the handler was written against (#36306).- Engines serialized by TensorRT 8.x must be rebuilt, as an engine can only be deserialized by the major version that built it.
- TensorRT 10 and later require a GPU with compute capability 7.5 or higher, which excludes NVIDIA Pascal and Volta GPUs.
- If dropping TensorRT 8.x support is a hard blocker for you, please comment on (#36306).
Bugfixes
- (Java) Fixed the Spark runner firing processing-time timers in reverse timestamp order (#39824).
- (Java) Fixed the Spark runner dropping the stored watermark of a streaming source with no update in a batch (#39822).
- (Python) Fixed incorrect profiler options handling on portable runners (#39613).
- (Java) KafkaIO dynamic reads no longer require the obsolete
beam_fn_apiexperiment (#29998). - (Prism) Self-checkpointing splittable DoFns now resume after their requested delay instead of immediately, so polling SDFs no longer busy-spin (#39848).
- (Java) MongoDbIO read splitting now preserves non-ObjectId
_idtypes (e.g. string ids) instead of failing to parse the generated range filters (#39900). - (Go) Fixed GCS glob matching silently dropping objects when the glob pattern contains multi-byte characters (#39969).
- (Python) Fixed
TensorRTEngineHandlerNumPyfailing withCUDA_ERROR_INVALID_VALUEon models with a single-element input or output tensor (#36306). - (Python) Fixed
PickleCoder/_MemoizingPickleCoder.as_deterministic_coder()raisingTypeErrorinstead of returning a working deterministic coder (#28558).
According to git shortlog, the following people contributed to the 2.77.0 release. Thank you to all contributors!
Abdelrahman Ibrahim, Aditya Narayan, Ahmed Abualsaud, Alex Bevilacqua, Alexander Pochill, Ali Ebrahim, Andrew Crites, Arun Pandian, Ashwin S, Bruno Volpato, Chamikara Jayalath, Chris Gavin, Claire McGinty, Danny McCormick, Derrick Williams, Eiji Ogiwara, Elia Liu, Fabian Loris, Goutam Adwant, HansMarcus01, Israel Herraiz, Jack McCluskey, Jan Lukavský, Jeremy Schoemaker, Kenneth Knowles, Lalit Yadav, Lawrence Qiu, M Junaid Shaukat, Makoto Nagai, Maksym Tymoshyk, Manvith Panyam, Mattie Fu, Michael Gruschke, Mukesh Bhandarkar, Nicolas Gibanel, Paulius Kuzmickas, Radosław Stankiewicz, Ryan Wigglesworth, Sam Whittle, Sharan Teja M, Shizuma5, Shunping Huang, SreeramaYeshwanthGowd, Tobias Kaymak, Tom Newton, Udit Jain, Vitaly Terentyev, Yi Hu, ZIHAN DAI, akshayjadiyanv, claudevdm, darshan-sj, feefs, junaiddshaukat, kellen, nitinware, parveensania, tvalentyn