Apache Spark Scala用Regex处理RDD提取每行首个逗号前内容问题
错误原因分析
你的实现存在3个核心问题:
- 对RDD调用了
collect()方法,直接把分布式RDD转换成了本地Scala数组,失去了Spark分布式计算的能力,数据量较大时会直接触发本地内存溢出 String.matches()方法的返回值本身就是Boolean类型,仅用于判断整个字符串是否符合正则规则,不会返回正则匹配到的具体内容- 即使正则逻辑正确,你也没有提取正则捕获组里的内容,自然拿不到目标子串
正确实现方案
方案1:用split方法(更简单,性能更好)
对于按第一个逗号分割的简单场景,直接用split指定分割次数即可,避免正则的额外开销:
// 不要提前collect,直接对RDD做map操作 val originalRdd = spark.sparkContext.parallelize(List( "Et: NT=grouptoClassify,hadoop-exec,sparkConnection,Ready", "group: NT=app_1,hadoop-exec,sparkConnection,Ready", "group: NT=app_exmpl_2,DB-exec,MDBConnection,NR", "group: NT=apprexec,hadoop-exec,sparkConnection,Ready", "group: NT=nt_prblm_app,hadoop-exec,sparkConnection,NR" )) // split第二个参数设为2,代表最多分割成2段,取第一段就是第一个逗号之前的内容 val resultRdd = originalRdd.map(line => line.split(",", 2)(0)) // 测试环境可以调用collect打印结果,生产环境不要随意使用collect resultRdd.collect().foreach(println)
方案2:正则提取(适合更复杂的匹配场景)
如果一定要用正则实现,需要调用find方法匹配内容,再提取捕获组:
val regex1 = """^(.+?),""".r val resultRdd = originalRdd.map { line => regex1.findFirstMatchIn(line).map(_.group(1)).getOrElse(line) }
输出结果
两种方案都可以得到预期输出:
Et: NT=grouptoClassify group: NT=app_1 group: NT=app_exmpl_2 group: NT=apprexec group: NT=nt_prblm_app
内容的提问来源于stack exchange,提问作者CloudSparkie
相关产品推荐
相关产品推荐

