十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spark Action认识

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()。
返回列表