Spark Streaming Scala程序报错:Regex$MatchIterator无法序列化求解
Spark Streaming词频统计序列化问题解决
问题描述
本人是Scala新手,正在编写一个Spark Streaming程序,用于读取文本文件内容,去除非字母数字字符后进行映射、归约并输出词频。程序启动后在val lines_map = lines.map(line => pat.findAllIn(line))行报错,错误信息为:
object not serializable (class: scala.util.matching.Regex$MatchIterator, value: empty iterator
原运行命令
spark-submit --class group.WordCount --master yarn --deploy-mode client bdpAssignmentFour.jar hdfs:///user/s3797303
原Scala代码
import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} object WordCount { def main(args: Array[String]): Unit = { val sconf = new SparkConf().setAppName("SparkWordCount") val ssc = new StreamingContext(sconf, Seconds(5)) val pat = "^[a-zA-Z0-9]*$".r val lines = ssc.textFileStream(args(0)) val lines_map = lines.map(line => pat.findAllIn(line)) lines_map.print() val wordCounts = lines_map.map((_, 1)).reduceByKey(_ + _) wordCounts.print() ssc.start() ssc.awaitTermination() } }
问题原因
- 序列化问题:
pat.findAllIn(line)返回的是Regex.MatchIterator,这是一个不可序列化的迭代器对象。Spark需要将闭包中的对象序列化后分发到Executor节点执行,迭代器无法被序列化,因此抛出错误。 - 正则逻辑错误:原正则
"^[a-zA-Z0-9]*$"是匹配整行完全由字母数字组成的内容,而不是提取每行中的字母数字单词,无法实现拆分单词的需求。
解决方法与修正代码
关键修改点
- 调整正则表达式为
"[a-zA-Z0-9]+",用于提取每行中的所有字母数字单词 - 使用
flatMap替代map,并将findAllIn(line)转换为可序列化的List[String],扁平化后得到单个单词的数据流 - 修正词频统计的逻辑,直接对单个单词进行计数
修正后的代码
import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} object WordCount { def main(args: Array[String]): Unit = { val sconf = new SparkConf().setAppName("SparkWordCount") val ssc = new StreamingContext(sconf, Seconds(5)) // 匹配每行中的所有字母数字单词 val wordPattern = "[a-zA-Z0-9]+".r val lines = ssc.textFileStream(args(0)) // flatMap将每行的单词列表扁平化,得到单个单词的流,toList将迭代器转为可序列化的集合 val words = lines.flatMap(line => wordPattern.findAllIn(line).toList) words.print() // 词频统计:每个单词映射为(单词,1),再按key归约求和 val wordCounts = words.map(word => (word, 1)).reduceByKey(_ + _) wordCounts.print() ssc.start() ssc.awaitTermination() } }
说明
flatMap会将每行拆分出的单词列表展开,让每个单词成为流中的独立元素,这是统计词频的正确姿势toList将不可序列化的MatchIterator转换为可序列化的List[String],解决了序列化问题- 修正后的正则能正确提取每行中的所有字母数字单词,符合词频统计的需求
内容的提问来源于stack exchange,提问作者Jaimee-lee Lincoln
相关产品推荐
相关产品推荐

