Spark MLlib Rating功能无法运行,Zeppelin无报错终止执行求助
嗨,我看到你作为Scala和Spark新手,虽然有Python的丰富经验,但在Zeppelin里用org.apache.spark.mllib.recommendation.Rating时遇到了代码执行突然终止却无报错的问题,结合你的代码片段,我来帮你分析下可能的原因和解决办法:
核心问题:类型匹配不兼容
先看你代码里的这部分关键逻辑:
case sql.Row(user: Int, product: Int, rating: Int) => Rating(user.toInt, product.toInt, rating.toInt)
这里有两个明显的问题:
- CSV默认读取类型是String:Spark读取CSV文件时,默认会把所有字段识别为
String类型,你直接在模式匹配里指定Int类型,会导致匹配失败,触发Scala的MatchError。但Zeppelin的前端有时候不会显示这类异常,只会直接终止执行。 - 冗余的类型转换:就算真的匹配到了
Int类型,再调用.toInt完全是多余的,反而可能引发不必要的问题。
修正后的代码
把模式匹配改成针对String类型,再转换为Int,就能解决这个问题:
val ratings = spark.sql("SELECT _c0 as user,_c1 as product,_c2 as rating FROM datasetframe").rdd.map { case sql.Row(user: String, product: String, rating: String) => Rating(user.toInt, product.toInt, rating.toInt) }
另外,建议你先执行data.printSchema()来确认CSV读取后的字段类型,这能帮你快速验证类型是否符合预期。
额外的排查与优化建议
- 查看Zeppelin后端日志:Zeppelin前端没显示报错,不代表没有错误。你可以去Zeppelin安装目录下的
logs文件夹查看后端日志,那里会记录更详细的异常信息,帮你精准定位问题。 - 添加异常捕获逻辑:在map操作里加上异常处理,能打印出出错的具体行数据,方便排查:
val ratings = spark.sql("SELECT _c0 as user,_c1 as product,_c2 as rating FROM datasetframe").rdd.map { case sql.Row(user: String, product: String, rating: String) => try { Rating(user.toInt, product.toInt, rating.toInt) } catch { case e: Exception => println(s"Failed to process row: user=$user, product=$product, rating=$rating") throw e } }
- 优先使用DataFrame API:Spark的DataFrame/Dataset API比RDD更简洁、类型安全,推荐用这种方式实现相同逻辑,出错时也会给出更明确的提示:
import org.apache.spark.sql.types.IntegerType import org.apache.spark.mllib.recommendation.Rating // 先把字段转换为Int类型 val ratingsDF = spark.sql("SELECT _c0 as user,_c1 as product,_c2 as rating FROM datasetframe") .withColumn("user", $"user".cast(IntegerType)) .withColumn("product", $"product".cast(IntegerType)) .withColumn("rating", $"rating".cast(IntegerType)) // 转换为RDD[Rating] val ratings = ratingsDF.as[(Int, Int, Int)].map { case (u, p, r) => Rating(u, p, r) }.rdd
内容的提问来源于stack exchange,提问作者Wired604
相关产品推荐
相关产品推荐

