Apache Spark Scala拆分RDD提取每行首段 仅返回首条问题求助
问题原因
- 你调用
parallelize时传入的是List(value),这里的value是完整的多行文本字符串,所以生成的RDD仅包含1个元素(全部文本),而非每行对应一个独立的RDD元素,拆分后自然只能得到第一行逗号前的内容。 - 提前调用
collect()会把全量RDD数据拉取到Driver端转成本地数组,数据量较大时会触发内存溢出,非必要不要在处理流程中提前调用。
正确实现代码
如果你的原始输入是多行字符串,先按换行拆分后再生成RDD:
// 示例原始多行字符串 val rawValue = """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 """ // 先按换行拆分得到每行的集合,过滤空行后生成RDD val lineRdd = spark.sparkContext.parallelize(rawValue.split("\n").filter(_.trim.nonEmpty)) // 提取每行第一个逗号前的内容 val resRdd = lineRdd.map(line => line.split(",")(0)) // 需要查看结果时再调用collect打印 resRdd.collect().foreach(println)
如果你的数据来源于文件,直接用textFile读取即可,默认按行拆分生成RDD:
val lineRdd = spark.sparkContext.textFile("你的数据文件路径") val resRdd = lineRdd.map(_.split(",")(0))
输出结果
Et: NT=grouptoClassify group: NT=app_1 group: NT=app_exmpl_2 group: NT=apprexec group: NT=nt_prblm_app
内容的提问来源于stack exchange,提问作者CloudSparkie
相关产品推荐
相关产品推荐

