发布时间:2026/8/26 18:00:17
SparkStreaming 之 foreachRDD 算子详解及代码实现 摘要foreachRDD 是 Spark Streaming 里最常用、也最容易写挂的输出算子。这篇讲清它和 transform 的区别——foreachRDD 在 Driver 端拿到 RDD、真正执行在 Executor 端以及由此带来的一个高频坑连接到底该建在哪。用三种连接方式的性能对比和一段完整的写 MySQL 代码把 foreachRDD 的正确姿势讲透。关键词Spark Streaming, foreachRDD, foreachPartition, 连接管理, 事务, exactly-once一、foreachRDD 是什么先分清它和 transformDStream 的算子分两类转换操作返回新 DStream和输出操作不返回触发计算。foreachRDD属于后者是 output 操作。它和transform长得像本质区别就一条// transform拿到 RDD返回新 RDD继续流式计算valnewDsds.transform(rddrdd.map(...))// foreachRDD拿到 RDD做输出到此为止无返回值ds.foreachRDD(rdd{/* 写外部存储 */})foreachRDD是 Action 语义——一旦调用前面所有 DStream 转换才会真正执行。所以一个流里如果没有 foreachRDD或 print、saveAsTextFiles 这类 output 操作整个流是不会动的。二、最核心的坑连接建在哪这是 foreachRDD 用法的分水岭。关键要理解执行位置foreachRDD的闭包在 Driver 端定义、在 Driver 端拿到 RDD但它内部的foreachPartition/foreach真正跑在Executor端。而闭包里引用的外部变量会被序列化发送到 Executor。连接对象Connection不可序列化所以在foreachRDD这一层建连接要么报NotSerializableException要么每个 Driver 只建一个连接却要发给所有 Executor完全不对。结论一句话连接必须在 Executor 端创建也就是放进foreachPartition或foreach里面而不是foreachRDD这一层。三、三种连接方式性能差三个数量级假设要写 100 万条记录到 MySQL看三种写法的差别。方式一foreach 每条建连接反模式rdd.foreach{rowvalconnDriverManager.getConnection(url,user,pwd)save(conn,row)conn.close()}100 万条记录 100 万次 TCP 握手 认证 建连接。连接开销远大于写入本身吞吐直接崩。这是最常见的新手错误要极力避免。方式二foreachPartition 每分区建连接常用rdd.foreachPartition{itervalconnDriverManager.getConnection(url,user,pwd)try{iter.foreach(rowsave(conn,row))}finally{conn.close()}}foreachPartition的闭包在每个分区的第一条记录上执行一次所以一个分区只建一个连接分区内所有记录复用。假设每分区 1000 条连接次数就从 100 万降到 1000。这是大多数场景够用的写法。注意finally里关连接防止中途异常导致连接泄漏。方式三连接池生产最优方式二的问题在于连接数 分区数分区一多连接就多。生产上用连接池跨分区复用// Executor 端懒加载单例连接池objectConnectionPool{lazyvalds:HikariDataSource{valcfgnewHikariConfig()cfg.setJdbcUrl(url);cfg.setUsername(user);cfg.setPassword(pwd)cfg.setMaximumPoolSize(10)newHikariDataSource(cfg)}}rdd.foreachPartition{itervalconnConnectionPool.ds.getConnectiontry{iter.foreach(rowsave(conn,row))}finally{conn.close()}}连接池用lazy val单例 静态对象保证每个 Executor 只初始化一次池子里的连接跨分区复用连接数可控。四、写外部存储的事务问题foreachRDD 写外部存储默认是 at-least-once 语义处理到一半 Executor 挂了重算时会重复写入已经写过的记录。要往 exactly-once 靠需要三件事配合幂等写给记录设计唯一键用 upsert 替代 insert重复写同一行结果不变。每分区一个事务一个分区的写入放在一个事务里全部成功才 commit失败整体回滚。offset 与结果绑定只有写入成功才提交 offset这在 Direct 模式下天然支持offset 自管理。纯靠 foreachRDD 单算子做不到 exactly-once它只能做到配合外部系统幂等 事务之后的效果。这点别被网上foreachRDD 保证 exactly-once的说法误导。五、完整示例写 MySQLimportorg.apache.spark.streaming.{Seconds,StreamingContext}importcom.zaxxer.hikari.{HikariConfig,HikariDataSource}objectStreamingToMySQL{// Executor 端懒加载连接池单例objectConnectionPool{lazyvalds:HikariDataSource{valcfgnewHikariConfig()cfg.setJdbcUrl(jdbc:mysql://host:3306/db)cfg.setUsername(root)cfg.setPassword(password)cfg.setMaximumPoolSize(10)newHikariDataSource(cfg)}}defmain(args:Array[String]):Unit{valsscnewStreamingContext(local[2],stream-to-mysql,Seconds(2))vallinesssc.socketTextStream(localhost,9999)lines.foreachRDD{rddif(!rdd.isEmpty){rdd.foreachPartition{itervalconnConnectionPool.ds.getConnection conn.setAutoCommit(false)try{valpsconn.prepareStatement(INSERT INTO result(word, cnt) VALUES(?, 1) ON DUPLICATE KEY UPDATE cnt cnt 1)iter.foreach{wordps.setString(1,word)ps.addBatch()}ps.executeBatch()conn.commit()// 分区内一个事务成功才提交}catch{casee:Exceptionconn.rollback();throwe}finally{conn.close()// 归还连接池}}}}ssc.start();ssc.awaitTermination()}}这段代码把前四节的点都串起来了连接池单例、分区级连接、批处理 幂等 upsert、事务 commit/rollback。六、总结foreachRDD 是 output 操作Driver 端拿 RDD、Executor 端执行和 transform 的区别是有无返回值。连接必须建在 Executor 端foreachPartition/foreach 内foreachRDD 这层建连接会踩序列化坑。连接方式从差到好foreach 每条建连接 → foreachPartition 每分区建连接 → 连接池复用。默认 at-least-once要 exactly-once 得靠幂等写 分区事务 offset 绑定三件套。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践

相关新闻

2026/8/26 18:00:17

AI大模型与数学·第52课傅里叶级数基础2:任意周期T展开、复数形式傅里叶级数——扩散模型复频域噪声分解核心工具

本课定位 第51课我们完整学习了周期为2\pi的标准傅里叶级数,仅适配周期恰好等于2\pi的信号。 但现实AI工程里,时序、音频、图像、振动数据周期千差万别:音频采样周期、24小时时序、设备振动周期都不是2\pi,因此本节课首先推广任意…

2026/8/26 18:50:41

阿里巴巴国际站运营科普:平台逻辑、关键要素与增长路径

1. 阿里国际站是什么 阿里巴巴国际站是阿里巴巴集团旗下面向全球市场的B2B跨境电商平台,致力于连接中国供应商与海外采购商。自上线以来,平台逐步从早期的信息展示与商机撮合,演进为覆盖交易、支付、物流、通关等环节的数字化外贸基础设施。…

2026/8/26 18:50:41

一文学完linux必要点

本文仅作为了解使用linux。(个人笔记) 目录 一、linux认知与命令格式 1、 命令格式 1)选项的两种风格: 2)帮助系统 二、文件与目录操作 1、浏览与定位 1)ls 2)cd 2、创建操作 3、删除…

2026/8/26 18:45:40

Deep Learning for Computer Vision——Recurrent Neural Networks

第一部分:为什么需要 RNN?(序列建模问题)传统 CNN 和全连接网络(FC)的输入和输出大小是固定的(比如 224x224 的图片输入,输出 1000 个类别)。 但现实中有很多任务输入或输…

2026/8/26 9:13:28

[光学原理与应用-521]:对光的错误理解与纠偏

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

2026/8/25 11:48:27

SIP通话转接原理与REFER方法实战解析

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

2026/8/25 16:56:43

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

2026/8/26 0:04:32

Python random 模块常用函数详解:从入门到实战

目录 1. 引言2. 准备工作3. 基础随机函数4. 序列相关函数5. 随机种子与复现6. 实战案例7. 注意事项8. 常见问题与排查9. 总结 1. 引言 摘要: 本文系统介绍 Python 标准库 random 模块中最常用的随机数生成函数。内容涵盖基础随机函数(random()、unifor…

2026/8/26 1:19:35

JSON总结

JSON概念 JSON(JavaScript Object Notation) 是一种轻量级的数据交换格式,主要用于跟服务器进行交换数据。它基于ECMAScript的一个子集。 JSON采用完全独立于语言的文本格式,但是也使用了类似于C语言家族的习惯(包括C、C、C#、Java、JavaScr…

2026/8/26 1:19:35

保存连接sse 是什么原理,为什么不会一直请求

“保持连接”用的是 SSE(Server-Sent Events),本质是一个没有马上结束的 HTTP 请求。 过程是: 拷贝机发送一次请求: GET /api/code-sync/events服务器返回: Content-Type: text/event-stream但不关闭响应&…

2026/8/24 13:42:17

实测才敢推 AI论文网站 2026最新测评与推荐

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。一、综…

2026/8/24 18:13:48

2026必备!AI论文网站测评:最新推荐与深度对比

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

2026/8/25 1:08:14

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…