Spark Action认识
Spark Action在 Spark 中Action行动是触发 Spark 执行的操作。只有在遇到 Action 时Spark 才会将之前的 Transformation转换形成的 DAG 提交给执行引擎进行计算并返回结果。一、核心前提1. 区分核心概念Transformation转换操作仅定义数据处理规则、构建血缘依赖不触发任何计算如filter、select、groupBy、join等。Action行动操作触发Spark Job执行真正执行计算、读取数据、输出结果是Spark任务的“执行开关”。2. 关键原则无ActionSpark不干活只要触发Action整条血缘依赖会从头执行一次。二、数据查看类 Action核心作用调试常用用于预览数据。快速查看数据结构、内容方便调试代码不适合大数据量场景除show外。1. show()操作意义将DataFrame的数据打印到控制台可控制显示行数和是否截断字段是最常用的调试Action。案例# 默认显示前20行字符串字段自动截断df.show()# 显示前10行不截断字段适合查看长文本df.show(10,truncateFalse)# 显示所有列避免列被省略df.show(5,truncateFalse,verticalTrue)注意数据量较大时show()会自动采样显示不会拉取全量数据无OOM风险。2. collect()操作意义将集群中所有分区的DataFrame数据全部拉取到Driver端本地转换为Python List元素为Row对象。案例# 拉取全量数据到本地data_listdf.collect()# 遍历查看数据forrowindata_list:print(row[name],row[age])注意高危数据量超过Driver内存时会直接触发OOM内存溢出生产环境尽量避免使用仅用于小数据量调试。3. first()操作意义返回DataFrame的第一行数据Row对象仅拉取一行无OOM风险适合快速查看数据结构。案例# 获取第一行数据first_rowdf.first()# 提取第一行的指定字段print(first_row[id],first_row[score])4. head(n) / take(n)操作意义返回DataFrame的前n行数据转换为Python List元素为Row对象两者功能一致无本质区别。案例# 取前5行数据top5_rowsdf.head(5)# 取前3行数据等价于head(3)top3_rowsdf.take(3)注意仅拉取n行数据风险低于collect()但n过大如超过10万行仍可能OOM。5. takeAsList(n)操作意义与head(n)功能完全一致唯一区别是返回格式为标准Python Listhead(n)返回的是Spark封装的List。案例# 取前3行返回标准Listrow_listdf.takeAsList(3)三、统计计数类 Action核心作用业务常用用于数据统计。对数据进行计数、去重计数快速获取数据量相关信息是业务统计中最常用的Action。1. count()操作意义统计DataFrame的总行数触发全量数据扫描是最常用、最基础的统计Action。案例# 统计总行数total_rowsdf.count()# 统计过滤后的数据行数filtered_countdf.filter(age 18).count()2. countApprox(timeout)操作意义近似统计总行数在指定超时时间单位毫秒内返回估算结果无需扫描全量数据速度极快适合超大数据量场景。案例# 1000毫秒1秒内返回近似行数approx_totaldf.countApprox(timeout1000)注意结果是估算值精度可通过timeout调整超时时间越长精度越高。3. countApproxDistinct()操作意义近似统计某一列的去重行数无需扫描全量数据适合超大数据量的去重统计如用户数、设备数。案例# 近似统计user_id的去重数量approx_distinct_userdf.select(user_id).countApproxDistinct()三、输出写入类 Action核心作用生产必用用于数据落地。将计算后的DataFrame结果写入文件系统本地/HDFS或数据表Hive是数据处理的“最终落地步骤”。1. write 系列parquet/csv/json/orc操作意义将DataFrame以指定格式写入文件系统支持 overwrite覆盖、append追加等模式是生产环境最常用的落地方式。案例# 写入parquet格式推荐压缩比高、读取快df.write.parquet(/tmp/output/parquet,modeoverwrite)# 写入csv格式可指定分隔符df.write.csv(/tmp/output/csv,sep,,headerTrue,modeappend)# 写入json格式df.write.json(/tmp/output/json,modeignore)# 写入orc格式Hive常用df.write.orc(/tmp/output/orc,modeoverwrite)注意mode参数可选overwrite覆盖、append追加、ignore存在则跳过、errorifexists存在则报错。2. saveAsTable()操作意义将DataFrame保存为Hive数据表临时表/永久表可直接通过SQL查询适合多任务共享数据。案例# 保存为永久表需指定数据库df.write.saveAsTable(db.target_table,modeoverwrite)# 保存为临时表仅当前SparkSession有效df.write.saveAsTable(temp_table,modeoverwrite,temporaryTrue)3. insertInto()操作意义将DataFrame的数据插入到已存在的Hive表中要求DataFrame的列名、类型与目标表完全一致。案例# 插入已存在的表覆盖原有数据df.write.insertInto(db.exist_table,overwriteTrue)# 插入已存在的表追加数据df.write.insertInto(db.exist_table,overwriteFalse)4. writeTo()Spark 3.0 标准写法操作意义Spark 3.0及以上版本的标准表写入方式功能更强大支持创建表、替换表、追加数据等替代saveAsTable()的推荐写法。案例# 创建表不存在则创建存在则报错df.writeTo(db.new_table).create()# 创建或替换表存在则覆盖df.writeTo(db.target_table).createOrReplace()# 追加数据到现有表df.writeTo(db.target_table).append()四、遍历操作类 Action核心作用用于自定义数据处理。对DataFrame的每一行/每一个分区执行自定义逻辑如数据清洗、写入外部系统适合复杂业务处理。1. foreach()操作意义对DataFrame的每一行数据单独执行自定义函数lambda或普通函数一行对应一次函数调用。案例# 遍历每一行打印指定字段df.foreach(lambdarow:print(f姓名{row[name]}年龄{row[age]}))# 自定义函数处理每一行defprocess_row(row):# 自定义逻辑如写入数据库、数据清洗ifrow[age]18:print(f{row[name]}已成年)df.foreach(process_row)注意函数内的逻辑无法直接操作Driver端的变量如列表、字典需通过广播变量传递。2. foreachPartition()操作意义对DataFrame的每一个分区执行一次自定义函数一个分区对应一次函数调用函数参数为分区的迭代器性能远高于foreach()。案例# 自定义函数处理一个分区的数据defprocess_partition(partition):# partition是分区的迭代器可遍历分区内所有行forrowinpartition:print(f分区内数据{row[name]})df.foreachPartition(process_partition)# 生产常用批量写入数据库一个分区建立一次数据库连接提升性能defbatch_write(partition):connget_db_connection()# 建立数据库连接forrowinpartition:conn.execute(insert into table values(?,?),(row[id],row[name]))conn.commit()conn.close()df.foreachPartition(batch_write)注意适合大数据量场景减少函数调用次数和资源消耗如数据库连接。五、Checkpoint 类 Action核心作用防OOM专用。将DataFrame数据落盘切断血缘依赖释放内存根治长血缘导致的OOM是大数据量、长链路任务的“保命操作”。1. checkpoint(eagerTrue)操作意义将DataFrame的数据写入指定目录需提前设置setCheckpointDir彻底切断血缘依赖eagerTrue表示立即触发Action落盘数据。案例# 1. 提前设置checkpoint目录必须先执行spark.sparkContext.setCheckpointDir(/tmp/spark-checkpoints)# 2. 执行checkpoint立即落盘切断血缘dfdf.checkpoint(eagerTrue)# 3. 后续操作血缘已断不会OOMdf.groupBy(id).count()注意eagerFalse时不立即触发Action需等待后续Action才会落盘checkpoint目录不会自动清理需手动删除。2. localCheckpoint(eagerTrue)操作意义轻量级Checkpoint仅将数据暂存到本地节点不切断血缘依赖适合临时暂存数据不用于防OOM。案例# 轻量级暂存数据不切断血缘dfdf.localCheckpoint(eagerTrue)注意不能根治OOM仅适合临时缓存节点宕机后数据会丢失。六、数值聚合类 Action核心作用用于快速求值。直接计算指定列的聚合结果如求和、最大值返回单个数值无需通过agg()定义规则简化代码。1. sum() / max() / min() / avg()操作意义对指定数值列直接计算求和、最大值、最小值、平均值返回一个Row对象需提取数值。案例# 计算age列的最大值max_agedf.select(age).max()[max(age)]# 计算salary列的总和total_salarydf.select(salary).sum()[sum(salary)]# 计算score列的平均值avg_scoredf.select(score).avg()[avg(score)]2. reduce()操作意义对DataFrame转换后的RDD执行聚合计算需自定义聚合逻辑返回单个聚合结果灵活性高。案例frompyspark.sql.functionsimportcol# 将DataFrame转为RDD计算id列的最大值max_iddf.select(col(id)).rdd.reduce(lambdaa,b:aifa[id]b[id]elseb)[id]# 计算salary列的总和total_salarydf.select(col(salary)).rdd.reduce(lambdaa,b:ab)[salary]七、RDD 侧 Action核心作用DF转RDD后常用。当DataFrame转换为RDD后可使用RDD的Action操作功能与DF侧Action类似适合复杂的RDD处理场景。常用案例# 1. 拉取全量数据高危df.rdd.collect()# 2. 统计RDD行数df.rdd.count()# 3. 取前5行数据df.rdd.take(5)# 4. 取第一行数据df.rdd.first()# 5. 写入文本文件df.rdd.saveAsTextFile(/tmp/output/rdd_txt)# 6. 写入对象文件df.rdd.saveAsObjectFile(/tmp/output/rdd_obj)八、常见误区1. 以下操作不是Action不触发计算聚合函数F.countDistinct()、F.collect_set()、F.sum()、F.max()仅用于agg()中定义规则转换操作filter、select、groupBy、join、orderBy、limit、union、distinct缓存操作cache()、persist()仅标记缓存不触发计算需后续Action触发。2. 缓存与Action的关系cache()/persist()本身不触发Action当后续出现任意一个Action时Spark会顺便将数据缓存到内存/磁盘后续再触发Action时可直接复用缓存避免重复计算。3. 高危Action避坑collect()、take(n)n过大会将数据拉回Driver极易OOM生产环境尽量避免如需查看数据优先使用show()。

相关新闻

Linux IIO子系统

Linux IIO子系统

因实际项目资料涉密,在此将项目中Linux common的部分和学习心得总结于此,以备时习之。(部分示例代码为脱敏经过删减和简化处理,仅供参考,请见谅!)参考图来源于网络,侵删。keyword&am…

2026/8/22 11:47:59 阅读更多 →
ROS机器人自主导航实战:从定位、建图到路径规划全链路解析

ROS机器人自主导航实战:从定位、建图到路径规划全链路解析

在实际机器人开发中,自主导航是衡量一个机器人系统是否“智能”的核心能力。它要求机器人能够回答“我在哪?”(定位)、“周围环境什么样?”(建图)以及“我该怎么去那里?”&#xff0…

2026/8/22 12:43:19 阅读更多 →
《我的世界》多人跑酷地图【跑酷惊魂+】核心玩法与部署指南

《我的世界》多人跑酷地图【跑酷惊魂+】核心玩法与部署指南

这次我们来看一个《我的世界》(Minecraft)社区地图项目:【跑酷惊魂】。这不是一个模组或插件,而是一张精心设计的多人小游戏地图,核心玩法是“随机生成的跑酷大比拼”。对于喜欢在《我的世界》里和朋友联机、寻求快节奏…

2026/8/23 17:20:51 阅读更多 →

最新新闻

半导体AI岗位面试准备与核心技术解析

半导体AI岗位面试准备与核心技术解析

1. 项目概述:半导体行业AI岗位的独特挑战 陕西华码半导体作为西北地区领先的集成电路设计企业,其AI开发岗位要求候选人同时具备算法工程化能力和半导体行业知识。我在参与他们2023年秋招时发现,技术笔试中30%的题目涉及芯片设计数据预处理&am…

2026/8/23 18:03:36 阅读更多 →
如何在 .NET MAUI 里加相机:Camera.MAUI 拍照与扫码完整指南

如何在 .NET MAUI 里加相机:Camera.MAUI 拍照与扫码完整指南

如何在 .NET MAUI 里加相机:Camera.MAUI 拍照与扫码完整指南 【免费下载链接】Camera.MAUI A CameraView Control for preview, take photos and control the camera options 项目地址: https://gitcode.com/gh_mirrors/ca/Camera.MAUI 想在 .NET MAUI 应用里…

2026/8/23 18:03:36 阅读更多 →
从竞争到合作:技术架构思维转变与开放生态构建

从竞争到合作:技术架构思维转变与开放生态构建

1. 从“为竞争而合作”到“为合作而竞争”:一个技术架构思维的转变这个话题乍一看有点宏大,甚至带点哲学意味,但它背后指向的,是我们在设计系统、构建平台、乃至组织技术团队时,一个非常根本的思维模式差异。如果你正在…

2026/8/23 18:03:36 阅读更多 →
PoeCharm:把配装试错变成一次模拟

PoeCharm:把配装试错变成一次模拟

PoeCharm:把配装试错变成一次模拟 【免费下载链接】PoeCharm Path of Building Chinese version 项目地址: https://gitcode.com/gh_mirrors/po/PoeCharm 做流放之路的Build规划时,每次换装都要翻浏览器查英文词缀,窗口切来切去太慢了…

2026/8/23 18:03:35 阅读更多 →
开源双臂机器人Hei-rebot-lift:从零构建具身智能硬件平台的工程实践

开源双臂机器人Hei-rebot-lift:从零构建具身智能硬件平台的工程实践

你有没有过这样的经历:看着网上那些炫酷的双臂机器人视频,心里痒痒的,也想自己动手造一个,但一查资料,要么是动辄几十上百万的工业级产品,要么就是只有论文和概念,连个螺丝钉都找不到&#xff1…

2026/8/23 18:03:35 阅读更多 →
46 个现成 Dify 工作流:导入、实战到排错的保姆级上手指南

46 个现成 Dify 工作流:导入、实战到排错的保姆级上手指南

46 个现成 Dify 工作流:导入、实战到排错的保姆级上手指南 【免费下载链接】Awesome-Dify-Workflow 分享一些好用的 Dify DSL 工作流程,自用、学习两相宜。 Sharing some Dify workflows. 项目地址: https://gitcode.com/GitHub_Trending/aw/Awesome-D…

2026/8/23 18:02:35 阅读更多 →

日新闻

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

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

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

2026/8/23 0:00:50 阅读更多 →
SIP通话转接原理与REFER方法实战解析

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

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

2026/8/23 0:00:50 阅读更多 →
Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

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

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

2026/8/23 0:00:50 阅读更多 →

周新闻

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

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

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

2026/8/23 0:00:50 阅读更多 →
SIP通话转接原理与REFER方法实战解析

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

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

2026/8/23 0:00:50 阅读更多 →
Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

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

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

2026/8/23 0:00:50 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/22 18:08:39 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/23 12:10:44 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/22 3:22:48 阅读更多 →