本地Spark任务组合使用转换与行动算子后卡住无法完成如何排查?
排查步骤与解决方案
1. 调整日志级别定位报错
当前你将Spark日志级别设置为ERROR,会屏蔽任务运行时的异常信息,多数“卡住”场景实际是任务反复报错重试。先将日志级别调整为WARN或INFO,查看控制台输出的错误栈:
修改代码行:
sc.setLogLevel("WARN")
优先排查是否有数组越界、空指针类的运行时异常。
2. 修复解析方法的边界问题
你当前的parse方法没有做格式兼容:如果某一行数据不符合预期的XML格式,line.indexOf(subs)会返回-1,后续的substring操作会直接抛出StringIndexOutOfBoundsException。
之前的wikiRdd.count()、wikiRdd.take(2)能快速返回,是因为这类操作仅处理前N行数据,不会遍历全量数据集;而加了filter之后需要遍历所有行执行解析、过滤逻辑,一旦碰到格式异常的行就会报错,日志被屏蔽后就会误以为任务卡住。
可以先给parse方法增加异常兼容验证问题:
def parse(line: String): Option[WikipediaArticle2] = { val subs = "</title><text>" val i = line.indexOf(subs) if (i < 14) return None try { val title = line.substring(14, i) val text = line.substring(i + subs.length, line.length-16) Some(WikipediaArticle2(title, text)) } catch { case e: Exception => None } }
调整RDD生成逻辑过滤无效数据:
val wikiRdd: RDD[WikipediaArticle2] = sc.parallelize(WikipediaData2.lines) .map(WikipediaData2.parse) .filter(_.isDefined) .map(_.get)
重新运行过滤逻辑验证是否恢复正常。
3. 新增缓存避免重复计算
你当前的WikipediaData2.lines用def定义,每次访问都会重新读取文件流生成List,且wikiRdd未做缓存,每次行动算子触发都会重新执行读文件、解析的全流程,数据量较大时会有明显性能损耗,可做两点优化:
- 将
lines改为lazy val,避免重复读取文件:
lazy val lines: List[String] = { // 原有逻辑保持不变 }
- 对
wikiRdd增加缓存:
val wikiRdd: RDD[WikipediaArticle2] = sc.parallelize(WikipediaData2.lines) .map(WikipediaData2.parse) .cache()
4. 优化语言匹配逻辑性能
当前mentionsLanguage中text.split(' ').contains(lang)的写法,会将整段文本全部切割为数组后再匹配,大文本场景下性能极差,可直接用子串查找结合边界判断实现相同逻辑,性能提升非常明显:
def mentionsLanguage(lang: String): Boolean = { val idx = text.indexOf(lang) if (idx == -1) return false // 匹配前后边界,避免匹配到单词片段 val beforeValid = idx == 0 || !text.charAt(idx - 1).isLetter val afterValid = idx + lang.length == text.length || !text.charAt(idx + lang.length).isLetter beforeValid && afterValid }
内容的提问来源于stack exchange,提问作者sqoor
相关产品推荐
相关产品推荐

