ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Beam 2.36.0 版本深度解析:Kafka 停止读取时间、cloudpickle 序列化与破坏性变更全指南

Apache Beam 2.36.0 版本深度解析:Kafka 停止读取时间、cloudpickle 序列化与破坏性变更全指南 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 2.36.02022-02-07 发布是一次聚焦于 I/O 能力增强与 Python SDK 运行时体验改进的版本Java 版 KafkaIO 引入了 SDFSplittable DoFn场景下的stopReadTime停止读取时间控制Python SDK 新增 cloudpickle 序列化库、BigQuery 流式写入触发频率与 Dataflow 工件缓存等实用选项同时带来 RedisIO 的 jedis 3.x→4.x 升级等一系列破坏性变更。读完本文你将掌握 2.36.0 的核心新 API 用法、关键命令行参数、破坏性变更的迁移要点以及如何规避本版本已知问题。版本概览与获取方式2.36.0 是 Apache Beam 在 2022 年初发布的正式版本包含功能改进与全新能力官方发布的下载入口为项目网站的下载页面对应 2.36.0 / 2022-02-07 条目。详细的逐项变更可查阅官方 JIRA 的 Release Notes。本仓库的版本发布公告位于 website/www/site/content/en/blog/beam-2.36.0.md其内容划分为 I/Os、New Features / Improvements、Breaking Changes、Known Issues 与贡献者致谢几个部分下文逐节展开。I/O 增强KafkaIO SDF 新增 stopReadTime 停止读取时间2.36.0 在 Java SDK 的 KafkaIO 上落地了 [BEAM-13171]为基于 SDFSplittable DoFn的读取路径新增stopReadTime支持允许用户指定一个绝对时间戳让读取在到达该时间点后停止。源码中的 API 形态从 KafkaIO.java 可以看到该能力的 Builder 方法public ReadK, V withStopReadTime(Instant stopReadTime) { return toBuilder().setStopReadTime(stopReadTime).build(); }stopReadTime内部存储为Longepoch 毫秒在翻译translation阶段被转回Instant.ofEpochMilli(...)见 KafkaIO.java在 SDF 的ReadFromKafkaDoFn中stopReadTime会与startReadTime一同被传递给 Kafka 的offsetForTimes查询用于把时间戳换算为每个分区的起始/结束 offset见 KafkaIO.java。使用前提与失败语义该方法与已有的withStartReadTime(Instant)KafkaIO.java配套使用其 Javadoc 明确了两点硬性约束仅支持 Kafka Client 0.10.1.0 及以上版本且消息格式版本需在 0.10.0 之后即消息必须携带时间戳两种情况下会硬失败hard failure某个分区内不存在时间戳大于等于目标时间戳的消息分区消息格式版本早于 0.10.0消息没有时间戳。典型用法示例pipeline.apply( KafkaIO.String, Stringread() .withBootstrapServers(broker:9092) .withTopic(events) .withStartReadTime(Instant.parse(2022-02-07T00:00:00Z)) .withStopReadTime(Instant.parse(2022-02-07T12:00:00Z)) .withKeyDeserializer(StringDeserializer.class) .withValueDeserializer(StringDeserializer.class));配套的测试覆盖可见于 KafkaIOTest.java 与 ReadFromKafkaDoFnTest.java它们验证了 stopReadTime 在 SDF 读取路径上的行为。新特性与新功能开箱即用的 ARM64 / Apple M1 支持[BEAM-11703] 为 Beam 补齐了 ARM64 架构支持在 Apple M1以及各类 ARM64 环境上无需额外配置即可直接运行 Beam 相关组件。对于需要在本地 M1 上开发或测试 Dataflow 管道的用户这是一个重要的体验改进。Python SDKcloudpickle 序列化库支持[BEAM-8123] 为 Python SDK 引入了 cloudpickle 作为新的序列化后端。cloudpickle 擅长序列化在__main__等交互式环境中定义的函数与类即主会话状态能显著降低无法 pickle 本地定义的函数这类常见报错。启用方式设置管道选项--pickle_librarycloudpickle。从 pipeline_options.py 的选项定义可以看到完整的取值集合--pickle_library defaultdefault choices[cloudpickle, default, dill, dill_unsafe]default由 Beam 自行选择默认序列化库cloudpickle使用内置于apache_beam.internal.cloudpickle的 cloudpickle 实现dill/dill_unsafe使用 dill 系列需要额外安装apache-beam[dill]extras否则提交作业时会得到明确的错误提示。从 pickler.py 的实现看Beam 内部通过USE_CLOUDPICKLE/USE_DILL/USE_DILL_UNSAFE三个标记管理序列化库的切换并在set_pickler中根据选项完成cloudpickle_pickler与dill_pickler的挂钩覆盖pickler.py。值得注意的联动行为--save_main_session选项的默认行为会随序列化库变化——在 Dataflow runner 上当选用 cloudpickle 作为 pickle 库时save_main_session默认开启见 pipeline_options.py 的帮助文本。这意味着迁移到 cloudpickle 后主会话中定义的函数与变量会被自动保存并同步到 worker交互式开发场景更顺畅。Python BigQuery流式写入触发频率triggering_frequency[BEAM-12865] 为 Python SDK 的 BigQuery I/O 增加了触发频率选项用于控制流式写入批次提交的节奏。该参数在 bigquery.py 的WriteToBigQuery中通过triggering_frequency参数暴露其语义如下bigquery.py类型为float会被转换为int每隔triggering_frequency秒触发一次批次提交当有数据等待写入时最多每triggering_frequency秒提交一批流中的行每隔triggering_frequency秒提交一次。示例from apache_beam.io.gcp.bigquery import WriteToBigQuery rows | WriteToBigQuery WriteToBigQuery( tableproject:dataset.table, write_dispositionWriteToBigQuery.WriteDisposition.WRITE_APPEND, triggering_frequency60, # 每 60 秒触发一次批次提交 )需要特别注意的是参数约束当triggering_frequency与STREAMING_INSERTS结合使用时必须同时开启with_auto_sharding否则会抛出校验错误见 bigquery.py。这一约束源于流式插入路径对负载分片的依赖迁移时务必检查现有管道是否满足。Python Dataflow工件缓存enable_artifact_caching[BEAM-13459] 为 Python Dataflow 作业增加了上传工件缓存开关。启用后已上传的工件artifact会在 GCS staging bucket 中跨作业提交复用减少重复上传的开销。启用方式设置管道选项--enable_artifact_caching。从 pipeline_options.py 的定义看该选项默认关闭defaultFalse帮助文本说明开启后工件将在 GCS staging bucket 中被跨作业缓存官方明确表示该行为将在未来版本默认开启。因此当前版本属于主动开启、提前体验阶段升级时可提前验证缓存与既有 staging 策略的兼容性。破坏性变更与迁移指南RedisIOjedis 从 3.x 升级到 4.x[BEAM-12092] 将 Java RedisIO 底层客户端从 jedis 3.x 升级到 4.x。当前仓库中 RedisIO 的 build.gradle 锁定的版本为redis.clients:jedis:4.0.1。影响范围如果你在管道中直接使用 jedis而非仅通过 RedisIO 封装需要参照 jedis 官方的 3-to-4 迁移文档更新 API 调用——jedis 4.x 在Jedis连接管理、命令返回值类型与事务 API 上有大量调整。从源码可见RedisConnectionConfiguration.java 与 RedisIO.java 已适配 4.x 的StreamEntryID、ScanParams、XAddParams等新 API。只使用 RedisIO 自身 API如RedisIO.read()/RedisIO.write()的用户通常无需改动但若项目中直接依赖了 jedis 旧版本需同步升级并处理 API 差异。AWS SQSSqsMessage 时间戳字段类型变更[BEAM-13638] 调整了 AWS IOsSDK v2中SqsMessage的结构时间戳字段数据类型从String改为long所有字段的可见性从package private修复为public。这意味着直接构造或读取SqsMessage时间戳字段的代码需要改为long语义epoch 时间同时跨包访问字段不再受限。Java SDKDoFn 输出时间戳严格校验[BEAM-12931] 为 Java SDK 增加了对 DoFn 输出元素、定时器timers以及onWindowExpiration回调中输出时间戳的校验。此前这些时间戳缺少统一检查违反时间戳约束的输出现在会被更严格地拦截并报错相关管道需要确保输出时间戳合法例如不超过 watermark 前进规则允许的范围。Python DataFrameDeferredDataFrame.xs 非 tuple 键缺陷修复[BEAM-13421] 修复了DeferredDataFrame.xs在使用非 tuple 键时的 bug。此前以单个标量键调用xs可能产生错误结果2.36.0 起该场景行为正确。Python SDKgoogle-cloud-pubsub 最低版本要求Python SDK 现在要求google-cloud-pubsub2.1.0。apache_beam.io.gcp.pubsub的 API 面没有变化但直接使用 PubSub 客户端库的代码可能需要随之更新例如依赖 2.x 系列 API 的方法签名调整。升级依赖时请确认环境中google-cloud-pubsub满足该版本下限。已知问题Known Issues2.36.0 存在以下官方记录的已知问题使用时需留意规避ArithmeticExceptionJava当输出元素的时间戳与 DoFn 允许的 allowedSkew 之差超出Integer.MAX_VALUE即设置的 allowedSkew 大于该阈值时可能抛出意外的java.lang.ArithmeticException。规避方式是将 DoFn 的allowedSkew控制在Integer.MAX_VALUE以内。S3 对象元数据检索失效PythonPython SDK 中 S3 对象元数据检索存在回归对应 [BEAM-13980]依赖 S3 元数据如对象大小、最后修改时间的管道在 2.36.0 上可能受影响。完整的受影响问题清单可在官方 JIRA 中按affectedVersion 2.36.0过滤查看。升级建议综合上述变更从 2.36.0 之前的版本升级时建议按以下清单自查使用 RedisIO 且直接依赖 jedis 的项目先完成 jedis 3→4 的 API 迁移再升级 BeamAWS SQS 相关代码检查SqsMessage时间戳字段的String→long变更Java 管道检查 DoFn / timer /onWindowExpiration输出时间戳的合法性适配新的严格校验Python 环境升级google-cloud-pubsub到2.1.0若计划使用 cloudpickle可直接设置--pickle_librarycloudpickle并利用其默认save_main_session行为若使用 dill 系列请安装apache-beam[dill]BigQuery 流式写入如需设置triggering_frequency在STREAMING_INSERTS模式下务必同时开启with_auto_shardingDataflow 用户可提前开启--enable_artifact_caching验证工件缓存为后续版本默认开启做准备关注已知问题清单尤其是 Java 的 allowedSkew 溢出异常与 Python S3 元数据问题。致谢与社区根据官方公告2.36.0 由大量社区贡献者共同完成公告中列出的贡献者名单覆盖来自 Google、AWS、Wizeline 等多家公司的开发者。向所有参与测试、修复与功能开发的贡献者致谢。参考资源版本发布公告website/www/site/content/en/blog/beam-2.36.0.mdKafkaIO 停止读取时间实现sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.javaPython pickle 库选项定义sdks/python/apache_beam/options/pipeline_options.pyPython 序列化库切换实现sdks/python/apache_beam/internal/pickler.pyBigQuery 触发频率参数sdks/python/apache_beam/io/gcp/bigquery.pyRedisIO jedis 版本sdks/java/io/redis/build.gradle赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Nomad 停止支持版本变更日志深度解读0.1.0 至 1.9.x 的安全修复、破坏性变更与演进脉络Nomad 停止支持版本变更日志深度解读0.1.0 至 1.9.x 的安全修复、破坏性变更与演进脉络 导读 本篇文章围绕仓库根目录的 CHANGELOG un任务调度云原生运维后端UMAP 逆变换inverse_transform实战指南从低维嵌入还原高维数据样本UMAP 逆变换inverse_transform实战指南从低维嵌入还原高维数据样本 UMAPUniform Manifold Approximatio大数据批处理流处理数据工程Apache DataFusion 43.0.0 版本深度解读破坏性变更、SQL 能力增强与性能优化全景Apache DataFusion 43.0.0 版本深度解读破坏性变更、SQL 能力增强与性能优化全景 导读 本文以官方 43.0.0 版本变更日志为核心大数据数据分析后端上一篇OpenEBS Replicated PV Mayastor 磁盘池静态加密At-Rest Encryption设计解析与实战指南下一篇开源剧本软件Trelby让创作回归内容本质的专业编剧工具创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进