【发布时间】:2016-03-23 11:29:26
【问题描述】:
我是 SPARK 的新手,正在寻找更好的方法来实现以下场景。 有一个包含 3 个字段的数据库表 - 类别、金额、数量。 首先,我尝试从数据库中提取所有不同的类别。
val categories:RDD[String] = df.select(CATEGORY).distinct().rdd.map(r => r(0).toString)
现在对于每个类别,我想执行 Pipeline,它基本上从每个类别创建数据帧并应用一些机器学习。
categories.foreach(executePipeline)
def execute(category: String): Unit = {
val dfCategory = sqlCtxt.read.jdbc(JDBC_URL,"SELECT * FROM TABLE_NAME WHERE CATEGORY="+category)
dfCategory.show()
}
有可能做这样的事情吗?还是有更好的选择?
【问题讨论】:
标签: scala apache-spark dataframe apache-spark-sql rdd