大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 2.11.0 是该项目于 2019 年发布的一个重要里程碑版本聚焦于改进与新增功能双线推进。本文以官方发布博客 beam-2.11.0.md 为核心骨架结合当前仓库源码逐项拆解该版本的依赖升级清单、I/O 能力增强Kafka 偏移量消费者、Parquet 压缩编解码器、BigQuery KMS 加密、GCS KMS 拷贝、新特性Python 3 实验性支持、ZStandard 压缩、CombineFn.compact、Spark/Flink Runner 优化及弃用项帮助读者快速评估升级收益与迁移影响。版本概览Apache Beam 2.11.0 于 2019-02-26 发布官方公告见仓库中的 beam-2.11.0.md。该版本同时包含改进与新增功能主要亮点集中在I/O 能力增强跨语言变换的 Portable Flink Runner 支持、GCS 拷贝的 Cloud KMS 支持、KafkaIO 偏移量消费者参数、ParquetIO 写入压缩编解码器、BigQuery KMS 密钥传递新特性Python 3实验性DirectRunner/DataflowRunner 支持、Java SDK 的 ZStandard 压缩、Python CombineFn.compact、Spark Runner GroupByKey 非合并窗口优化与 bundleSize 参数、Flink Runner 可移植 Runner savepoint/升级支持依赖大规模升级Java 侧 grpc、netty、google 生态组件批量升级Python 侧约束收紧弃用MongoDbwithKeepAlive因 Mongo 驱动中已弃用。依赖升级与变更Java 依赖升级2.11.0 对 Java 生态依赖进行了系统性升级直接影响到使用这些组件的 I/O 连接器如 gRPC 相关连接器、Bigtable、Spanner、Pub/Sub、GCS 等。核心升级包括类别组件新版本解析antlr / antlr_runtime4.7gRPC 全家桶grpc_all / grpc_auth / grpc_core / grpc_pubsub_v1 / grpc_protobuf / grpc_protobuf_lite / grpc_netty / grpc_stub1.17.1Nettynetty_handler / netty_transport_native_epoll4.1.30.Finalnetty_tcnative_boringssl_static2.0.17.FinalGoogle Cloudgoogle_api_common / google_auth_library_credentials / google_auth_library_oauth2_http1.7.0 / 0.12.0 / 0.12.0google_cloud_core / google_cloud_core_grpc1.61.0google_api_services_dataflowv1b3-rev20190126-1.27.0google_cloud_bigquery_storage / google_cloud_bigquery_storage_proto0.79.0-alpha / 0.44.0google_cloud_spanner / proto_google_cloud_spanner_admin_database_v11.6.0 / 1.6.0bigtable_client_core1.8.0gax_grpc1.38.0数据访问bigdataoss_gcsio / bigdataoss_util1.9.16cassandra-driver-core / cassandra-driver-mapping3.6.0压缩commons-compress1.18zstd_jni1.3.8-3提示表格中各组件的新版本号与官方发布说明一致。对于使用 Google Cloud 连接器的任务升级到 2.11.0 时建议核对自身依赖树避免与上表版本产生冲突。Python 依赖变更Python SDK 在 2.11.0 中收紧了部分依赖版本约束其中带有python_version 3.0条件的项仅作用于 Python 2 环境futures3.2.0,4.0.0; python_version 3.0Python 2 下并发原语依赖pyvcf0.6.8,0.7.0; python_version 3.0google-apitools0.5.26,0.5.27google-cloud-core0.28.1google-cloud-bigtable0.31.1这与该版本同步引入的 Python 3 实验性支持相呼应Python 2 环境继续通过条件依赖保持兼容而 Python 3 环境则开始获得独立运行能力。I/O 能力增强Portable Flink Runner 跨语言变换支持2.11.0 中 Portable Flink Runner 开始支持运行跨语言cross-language变换即在同一 Pipeline 中混合使用不同 SDK如 Java 与 Python编写的变换。该能力建立在 Beam 的可移植性架构之上使 Flink 作为执行后端时可以调度来自多语言 SDK 的扩展服务。这意味着用户可以利用 Python/Java 各自的生态例如使用 Python 编写的数据处理逻辑配合 Java 的成熟 I/O 连接器由 Flink Runner 统一调度执行。GCS 拷贝的 Cloud KMS 支持该版本为 GCS 拷贝操作增加了 Cloud KMS 支持允许在跨存储桶复制对象时使用客户管理的加密密钥CMEK。这为数据在 GCS 间的迁移场景提供了加密密钥的可控性使企业可以在合规要求下统一管理对象加密密钥。KafkaIO 偏移量消费者配置2.11.0 为KafkaIO.read()新增了偏移量消费者offset consumer相关参数。从当前仓库的 KafkaIO.java 源码看Kafka 读取后端ReadFromKafkaDoFn实际运行两个消费者主消费者main consumer真正从 Kafka 读取数据次级偏移量消费者secondary offset consumer通过拉取每个分区的latest offset来估算 backlog积压量用于推进 watermark 与流量控制。默认情况下偏移量消费者继承主消费者的配置并使用自动生成的group.id。但在安全加固的 Kafka 集群如启用 SASL/SSL中这一默认行为可能失败运行时会出现如下 WARN 日志exception while fetching latest offset for partition {}. will be retried此时可通过新增的配置 API 为偏移量消费者单独注入配置。仓库中相关方法为 withOffsetConsumerConfigOverrides用法示例KafkaIO.String, Stringread() .withBootstrapServers(broker:9092) .withTopic(my-topic) .withOffsetConsumerConfigOverrides(ImmutableMap.of( sasl.mechanism, PLAIN, security.protocol, SASL_PLAINTEXT)) ...同时仓库还提供了对称的withConsumerConfigUpdates在默认消费者属性基础上合并更新与withConsumerConfigOverrides整体替换主消费者配置等方法见 KafkaIO.java#L2841-L2865 与 KafkaIO.java#L3016-L3019两者配合即可分别治理主消费者与偏移量消费者的配置。ParquetIO 写入压缩编解码器2.11.0 允许在ParquetIO.write()中显式设置压缩编解码器。仓库中 ParquetIO.java 的Sink.withCompressionCodec(CompressionCodecName)即对应此能力默认值为CompressionCodecName.SNAPPY见 ParquetIO.java#L1069-L1072pipeline.apply(...) .apply(ParquetIO.write(FileSystems.matchNewResource(/tmp/out.parquet, false)) .withSchema(schema) .withCompressionCodec(CompressionCodecName.GZIP));该设置会经由open()中的AvroParquetWriter.Builder传递到底层 Parquet 写入器ParquetIO.java#L1191-L1200。写入 Sink 上还提供了系列配套参数便于在开启压缩的同时精细控制文件布局与内存占用方法作用默认值/约束withCompressionCodec设置压缩编解码器SNAPPYwithConfiguration指定 Hadoop ConfigurationMap 或 Configuration—withRowGroupSize设置 row-group 大小必须为正整数否则使用底层默认值withPageSize设置 page 大小1 MBwithDictionaryEncoding开关字典编码默认开启withBloomFilterEnabled开关 bloom filter默认关闭withMinRowCountForPageSizeCheckpage 大小检查前最少缓冲行数大行场景可调低如 1100withMaxRowCountForPageSizeCheck强制 page 大小检查的最大缓冲行数防止行大小差异大时缓冲区溢出默认由 Parquet 估算上限 10000BigQuery 变换的 KMS 密钥传递该版本为 BigQuery 变换增加了kms_key支持并传递给 Dataflow 执行。仓库中 BigQueryIO.java 的读TypedRead、DynamicRead与写Write侧均提供withKmsKey(String)方法见 BigQueryIO.java#L3568-L3569BigQueryIO.writeTableRows() .to(project:dataset.table) .withKmsKey(projects/my-project/locations/global/keyRings/my-ring/cryptoKeys/my-key)在批量加载实现 BatchLoads.java 中kmsKey作为加载作业参数被传递从而在 Dataflow 执行 BigQuery 导入/导出时使用指定的客户管理密钥加密目标表满足数据静态加密的合规诉求。新特性与改进Python 3 实验性支持2.11.0 为 DirectRunner 与 DataflowRunner 引入了 Python 3 的实验性支持。这是 Beam Python SDK 向 Python 3 迁移进程中的重要一步意味着上述两个 Runner 已可运行基于 Python 3 编写的 Pipeline但该能力在当时仍处于实验阶段官方尚未将其标注为生产级稳定支持。与之配套的 Python 依赖约束如futures仅在 Python 2 下引入也体现了双版本共存的过渡策略。Java SDK 的 ZStandard 压缩支持该版本为 Java SDK 增加了 ZStandardzstd压缩支持。仓库中 ZstdCoder.java 提供了若干工厂方法可用于组合出带 zstd 压缩的 Coder// 包裹内层 coder使用默认压缩级别 ZstdCoderT coder ZstdCoder.of(innerCoder); // 指定压缩级别 ZstdCoderT coder ZstdCoder.of(innerCoder, level); // 同时指定压缩字典与级别 ZstdCoderT coder ZstdCoder.of(innerCoder, dict, level);zstd 以高压缩比与较快的压缩/解压速度著称适合对 PCollection 中间数据或落盘数据进行压缩以节省存储与网络带宽。其依赖zstd_jni也正是上文依赖升级表中被提升至 1.3.8-3 的组件。Python CombineFn.compactPython SDK 新增CombineFn.compact与 Java SDK 的CombineFn.compact行为对齐。仓库中 core.py 的CombineFn基类及其派生实现均声明了compact(accumulator, *args, **kwargs)如 core.py#L1165其作用是在合并累加器之前对累加器进行压缩/规约以减少内存占用和后续合并的开销。典型场景是累加器体积较大、且存在冗余信息可在合并前去重的组合逻辑如超大规模去重、直方图/摘要统计等实现该方法可显著降低分布式执行中的传输与存储成本。Spark Runner 优化2.11.0 中 Spark Runner 有两项针对性优化GroupByKey 非合并窗口优化当窗口无需合并non-merging windows如 FixedWindows/SlidingWindows 之外的固定窗口场景时GroupByKey 的调度与数据组织得到优化减少了不必要的窗口合并开销bundleSize 参数新增bundleSize参数用于控制 Spark source 的分裂splitting粒度。仓库中 SparkPipelineOptions.java 定义了该选项其语义为若设置了 bundleSize将用它来分裂 BoundedSource否则使用默认值在 SourceRDD.java 中bundleSize 0时按该值决定每个分片的目标字节数否则回退到DEFAULT_BUNDLE_SIZE。合理调大 bundleSize 可减少任务切分数、降低调度开销调小则可提升并行度需要结合数据规模与集群资源权衡。Flink Runner可移植 Runner savepoint / 升级支持Flink Runner 的便携式portable执行路径新增了 savepoint 与作业升级支持允许用户在升级作业拓扑或 Beam 版本时基于 Flink 的 savepoint 机制保留状态、实现无状态丢失的平滑升级。该能力对生产环境中的流式作业尤为重要。Bugfixes 与弃用Bugfixes官方发布说明以Various bug fixes and performance improvements概括了该版本的缺陷修复与性能改进清单。具体条目可参考 Apache JIRA 上该版本的 Release Notes详见原文档中指向的 JIRA 链接升级时建议结合自身使用到的连接器与 Runner 核对其修复项。弃用MongoDb withKeepAlive2.11.0 弃用了 MongoDB 连接器的withKeepAlive配置原因是该参数对应的功能已在底层 Mongo 驱动中废弃。使用旧版 API 的用户应迁移到驱动推荐的心跳/保活机制避免依赖已弃用行为。贡献者根据git shortlog统计2.11.0 版本共有 60 余位贡献者参与完整名单见官方发布博客 beam-2.11.0.md 的 Contributors 小节其中包括 Ahmet Altay、Kenneth Knowles、Robert Bradshaw、Tyler Akidau、Reuven Lax 等 Apache Beam PMC 成员与活跃社区开发者。升级建议小结针对 2.11.0 的升级评估可遵循以下要点核对依赖Java 侧 gRPC 1.17.1、Netty 4.1.30、Google Cloud 组件版本均有提升确认与自身项目无版本冲突利用新 I/O 能力安全 Kafka 集群务必配置withOffsetConsumerConfigOverridesParquet 输出按存储与查询成本选择压缩编解码器BigQuery/GCS 场景通过 KMS 密钥满足加密合规Python 3 迁移可开始用 DirectRunner/DataflowRunner 验证 Python 3 实验性支持但生产环境需评估实验状态的风险Runner 调优Spark 场景可使用bundleSize控制分裂粒度Flink 流式作业可利用 savepoint 升级能力规划滚动升级方案清理弃用 API移除 MongoDBwithKeepAlive调用。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam 2.18.0 版本全解析Spark Structured Streaming Runner、SQS/RabbitMQ I/O 与 SQL 能力升级Apache Beam 2.18.0 版本全解析Spark Structured Streaming Runner、SQS/RabbitMQ I/O 与 SQ大数据批处理流处理数据工程WSABuilds 完整指南3 步在 Windows 上跑起安卓应用WSABuilds 完整指南3 步在 Windows 上跑起安卓应用 想玩的某款安卓游戏Windows 上翻来覆去只有个只带亚马逊应用商店的 WSAGoo大数据批处理流处理数据工程Apache Beam 2.15.0 发布解读I/O 增强、SQL ParquetTable 与 Runner 改进全解析Apache Beam 2.15.0 发布解读I/O 增强、SQL ParquetTable 与 Runner 改进全解析 Apache Beam 2.15.大数据批处理流处理数据工程上一篇如何快速诊断nvm-windows问题自动化日志分析工具终极指南下一篇ObjectivePGP高级技巧如何优化加密性能与安全性创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
阅读完成 · 觉得有帮助?