Apache Beam 2.76.0 正式发布
Apache Beam 2.76.0 正式发布。本次更新包含多项功能改进与新特性,可前往下载页面获取。
如需了解 2.76.0 的详细变更,请查阅完整发行说明。
亮点
- 新增完整的 Iceberg 批处理与流式 changelog 数据源(CDC)(#38831)
- (Java)Dataflow Streaming Runner 支持逐元素的 OpenTelemetry 链路追踪跨阶段传播,可通过
--experiments=enable_otel_defaults,element_metadata_supported,disable_portable_worker启用。注意 Cloud Trace 会产生额外费用。(#33176) - (Java)KafkaIO 和 PubSubIO 的读写操作均支持 OpenTelemetry 头部传播。(#33176)
- (Java)SpannerIO 变更流新增 OpenTelemetry 链路追踪支持(#33176)
- (Python)JmsIO(IBM MQ、ActiveMQ 及其他提供商)现可通过跨语言调用在 Python 中使用(#30716)。
I/O
- Iceberg 依赖升级至 1.11.0(Java)(#38925)。
- 新增 ArrowFlight IO(Java)(#20116)。
- 新增 Delta Lake 批处理 changelog 数据源(CDC)(#39492)
新特性 / 改进
- Go SDK 新增
GroupIntoBatches转换算子及标准beam:coder:sharded_key:v1编码器,并引入beam.Coder.IsDeterministic、beam.PCollection.WindowingStrategy和coder.RegisterDeterministicCoder,支持按需注册确定性自定义编码器(Go)(#19868)。 - TriggerStateMachineRunner 中 BitSetCoder 替换为 SentinelBitSetCoder 用于编码已完成的 bitset。两者状态兼容,均可解码对方编码的字节(#38139)。
envoy-data-plane(及其传递依赖 betterproto);EnvoyRateLimiter 现在改用内置的一小份 protobuf 定义,解决了下游项目的依赖冲突问题(#37854)。sys.path。这样,通过 ‘–files_to_stage’ 流水线选项提供的 Python 文件即可在流水线代码中直接导入,也方便通过 --beam_plugins 流水线选项在启动时初始化 Python SDK harness。详情见依赖管理文档中的 Staging Individual Files 一节。传入 ‘–experiments=no_staged_dir_in_sys_path’ 流水线选项可禁用此行为(#39431)。equal_to_approx,这是一个 assert_that 匹配器,可用可配置的容差比较流水线的数值输出(#18028)。Timestamp 现已支持可变的亚秒精度,最高到纳秒级。便携式 beam:logical_type:timestamp:v1 逻辑类型现在映射到 Python 的 Timestamp(#39344)。UnboundedSource,这是一个读取无限数据流的接口,支持 checkpoint、watermark 上报和 bundle 收尾。可通过 beam.io.Read 使用(#19137)。Watch 转换,它会为每个输入元素轮询不断增长的输出集合,跨轮询轮次对输出去重,并按用户指定的终止条件停止(#21521)。Breaking Changes
(Python) 已从 SDK 容器镜像中移除 `google-perftools`。如需使用 `--profiler_agent=tcmalloc`,请在自定义容器镜像中单独安装 google-perftools APT 包(#39323)。
[IcebergIO] 读取 `timestamptz` 列现在将返回 `Timestamp.MICROS` Beam 逻辑类型以保留微秒精度(旧的 `Schema.FieldType#DATETIME` 原始类型会截断毫秒之后的部分)。当存在 `timestamptz` 列时,以下场景可能受到影响:
- 已有的流式读取管道。
- 从旧版 SDK 升级后的托管 Iceberg 批处理读取。
- Python 读取。
使用管道选项 `--updateCompatibilityVersion=2.75.0`(或更早版本)可保持旧行为(#39344)。
`DoFn.process` 返回 `str`、`bytes` 或 `dict`(而非包装在可迭代对象中)时,现在会抛出 `TypeError`,不再像之前那样静默地按字符/字节/键迭代(Python)(#18712)。
(Java) `PipelineResult` 新增 `DRAINING` 和 `DRAINED` 状态,包括 runner 状态映射和 Dataflow 更新处理(#39020)。
(Java) 由于 Iceberg 1.11.0 升级,IcebergIO 及使用它的项目现在必须使用 Java 17 或更高版本构建(#38925)。
Bugfixes
- 修复 Python Dataflow Flex Templates 中未解析的运行时 `ValueProvider` 选项被字符串化的问题(#39499)。
- 修复可拆分 DoFn 在 portable Flink runner(Java)上自 checkpoint 时 checkpoint 状态无限增长的问题(#27648)。
- 通过避免在创建缓存 invoker 时反复解析 DoFn 类型描述符,提升了 Java pipeline 的性能(#39309)。
- (Python)修复了 Python SDK 中的内存泄漏:此前将可能包含大 stack frame 的异常对象存入缓存会导致泄漏(#39406)。
已知问题
- (Java)使用 Flink 2.1 及以上版本的 Flink runner,且项目中依赖了需要
org.lz4:lz4-java的库(如 Kafka 客户端)时,可能遇到 Gradle capability 冲突,因为 Flink 2.1+ 自带at.yawk.lz4:lz4-java并声明了相同的 capability。解决方法:在build.gradle中添加一条capabilitiesResolution规则,指定选用at.yawk.lz4:lz4-java(#38947)。
根据 git shortlog 记录,以下贡献者参与了 2.76.0 版本的开发,感谢所有贡献者!
ADITYA RAJ, Abdelrahman Ibrahim, Ahmed Abualsaud, Alexander Kolb, Amar3tto, Andrew Crites, Arun Pandian, Aryankn29, Atharva Moroney, Avi Kondareddy, Bruno Volpato, Chamikara Jayalath, Chris Qiu, Claire McGinty, Danny McCormick, Derrick Williams, Elia Liu, Florian TREHAUT, Guflly, HansMarcus01, Ian Liao, Ivy Xu, Jack McCluskey, KRITI MITTAL, Kenneth Knowles, Lalit Yadav, Manvith Panyam, Minh Vu, Nikita Grover, PRADDZY, Peter Tran, Radosław Stankiewicz, Ryan Wigglesworth, Shahar Epstein, Shunping Huang, SreeramaYeshwanthGowd, Steven van Rossum, Tarun Annapareddy, Tejas Iyer, Tobias Kaymak, Tomasz Wojdat, Utkarsh Parekh, Venkata Bharath Malapati, Vitaly Terentyev, Yi Hu, ZIHAN DAI, aibrahiim, akshayjadiyanv, atognolas, claudevdm, janaom, jayjayakumar, raman118, shunping, tvalentyn