Spark Streaming文本文件流重复执行问题及单词统计优化咨询
文件流处理问题解决方案与Scala下划线详解
一、Spark文件流处理的循环/单次执行问题
1. 无限循环的原因与解决
用findAllIn正则提取时出现无限循环,大概率是处理过程中生成的文件被写入了监控目录,导致Spark重复读取新生成的文件;另外,若手动修改源目录下的文件(比如重命名),也会触发Spark Streaming的重复检测。
解决步骤:
- 分离源目录与输出目录:确保处理后的结果写入独立的HDFS路径,比如源目录设为
hdfs://input/texts,输出目录设为hdfs://output/counts,绝对不要让输出文件回流到源目录。 - 使用正确的流API:必须用Spark的流处理API而非批处理API,示例代码如下:
// Spark Streaming (DStream) 示例 import org.apache.spark.streaming.{StreamingContext, Seconds} import org.apache.spark.SparkConf val conf = new SparkConf().setAppName("WordCountStream") val ssc = new StreamingContext(conf, Seconds(10)) // 每10秒扫描一次目录 // 监控HDFS目录下的新增文件 val lines = ssc.textFileStream("hdfs://input/texts") // 用正则提取不含特殊字符的单词(字母+数字) val validWords = lines.flatMap(line => """[a-zA-Z0-9]+""".r.findAllIn(line)) val wordCounts = validWords.map(word => (word, 1)).reduceByKey(_ + _) wordCounts.print() ssc.start() ssc.awaitTermination() - 避免修改源文件:源目录只负责接收新增文件,不要在处理过程中对其进行移动、重命名等操作。
2. split仅处理一次的原因与解决
split按空格分割只执行一次,是因为你用了批处理API(比如sc.textFile("path")),它只会读取一次目录下的现有文件,不会监控新增内容。
解决:切换到流处理API(textFileStream或Structured Streaming的readStream),就能持续监控目录中的新增文件并处理。另外,split(" ")无法过滤特殊字符,比如hello!会被当成一个完整单词,不符合你的需求,所以还是推荐用正则提取。
如果用Structured Streaming(更推荐的新版本API),示例代码:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.{explode, regexp_extract_all} val spark = SparkSession.builder.appName("StructuredWordCount").getOrCreate() import spark.implicits._ // 监控新增文本文件 val df = spark.readStream.text("hdfs://input/texts") // 提取所有符合规则的单词并统计 val wordCounts = df.select( explode(regexp_extract_all($"value", """[a-zA-Z0-9]+""", 0)).as("word") ).groupBy("word").count() // 输出到控制台,也可以写入HDFS val query = wordCounts.writeStream .outputMode("complete") // 每次输出完整的统计结果 .format("console") .start() query.awaitTermination()
二、Scala中下划线的常见用法
你在Spark代码里遇到的下划线,主要有这几种场景:
- 匿名函数参数占位符:比如
reduceByKey(_ + _),等价于(a: Int, b: Int) => a + b,两个下划线分别代表匿名函数的第一个和第二个参数。 - 忽略无关参数:模式匹配中
case (_, count) => println(count),这里下划线表示忽略第一个参数,只关注第二个。 - 导入包下所有成员:
import org.apache.spark._,表示导入spark包下的所有类、对象和方法。 - 变量默认值初始化:
var total: Int = _,给变量赋对应类型的默认值(数值类型为0,字符串为null,布尔值为false)。 - 函数部分应用:比如
val add = (a: Int, b: Int) => a + b; val add5 = add(5, _),下划线表示第二个参数暂未指定,生成一个新的函数add5(b: Int) = 5 + b。 - 匹配集合剩余元素:比如
case head :: tail => ...里的tail是列表除第一个元素外的剩余部分,或者case Array(_, second, _*)表示忽略第一个元素,取第二个,忽略后面所有元素。
内容的提问来源于stack exchange,提问作者lordtaekwon
相关产品推荐
相关产品推荐

