【问题标题】:SPARK - Use RDD.foreach to Create a Dataframe and execute actions on the DataframeSPARK - 使用 RDD.foreach 创建数据框并在数据框上执行操作
【发布时间】: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


    【解决方案1】:
    // You could get all your data with a single query and convert it to an rdd
    val data = sqlCtxt.read.jdbc(JDBC_URL,"SELECT * FROM TABLE_NAME).rdd
    
    // then group the data by category
    val groupedData = data.groupBy(row => row.getAs[String]("category"))
    
    // then you get an RDD[(String, Iterable[org.apache.spark.sql.Row])]
    // and you can iterate over it and execute your pipeline
    groupedData.map { case (categoryName, items) =>
      //executePipeline(categoryName, items)
    }
    

    【讨论】:

    • groupBy 在这种情况下是有风险的 - 假设每个类别都有很多记录,这可能会导致 OutOfMemory 错误,因为单个不可分发的记录会包含太多数据。另外,这里真正的问题是不能在工作端调用executePipeline,因为它使用SQLContext 来加载DataFrame - 你不能序列化它。
    • 没有。我没有在工人端加载任何东西,这就是我的解决方案的重点:)
    • 关于其他评论,是的 groupBy 可能是一个潜在的昂贵操作,但在这种情况下,我们真的想要分组,如果数据集真的那么大,我会开始考虑规避它的方法.
    • 我注意到很多人不明白这个概念。
    【解决方案2】:

    您的代码将因TaskNotSerializable 异常而失败,因为您尝试在execute 方法中使用SQLContext(不可序列化),该方法应被序列化并发送给工作人员以在其上执行categories RDD 中的每条记录。

    假设您知道类别的数量是有限,这意味着类别列表不会太大而无法容纳您的驱动程序记忆,您应该将类别收集到驱动程序,并使用foreach 迭代该本地集合:

    val categoriesRdd: RDD[String] = df.select(CATEGORY).distinct().rdd.map(r => r(0).toString)
    val categories: Seq[String] = categoriesRdd.collect()
    categories.foreach(executePipeline)
    

    另一个改进是重用您加载的数据框,而不是执行另一个查询,对每个类别使用 过滤器

    def executePipeline(singleCategoryDf: DataFrame) { /* ... */ }
    
    categories.foreach(cat => {
      val filtered = df.filter(col(CATEGORY) === cat)
      executePipeline(filtered)
    })
    

    注意:为确保 df 的重复使用不会在每次执行时重新加载它,请确保在收集类别之前cache() 它。

    【讨论】:

    • 如果您收集驱动程序上的类别并对其进行迭代,它会按顺序运行,这有点违背了 spark 的目的。
    • 我假设#Categories records 好得多(这是基于groupBy解决方案可以 - 每个类别都会创建一个巨大的Iterable,只能由单个任务处理)。如果该假设是错误的 - 实际上这不是最佳解决方案。似乎将 executePipeline 更改为在 Iterable[Row] 上工作要么很难要么不可能,因为 OP 暗示它应该“应用一些机器学习”——如果它们的意思是 MLLib,则输入必须是 RDD/DataFrame。
    • @TzachZohar ,您的意思是说 Iterable[Row] ,这些数据是由 Row 迭代的吗?它不会在 spark 客户端上消耗更多的内存吗?
    • @user3252097 这就是为什么我说你不应该使用Iterable[Row](这就是为什么我仍然认为不推荐这里建议的其他解决方案)
    • @DanielB “如果你收集驱动程序上的类别并遍历它们,它会按顺序运行,这有点违背了 spark 的目的。”那么我应该如何处理这种情况,我也有同样的情况问题..stackoverflow.com/questions/54416623/…
    猜你喜欢
    • 2018-03-31
    • 2017-11-30
    • 2022-01-18
    • 1970-01-01
    • 1970-01-01
    • 2017-10-10
    • 1970-01-01
    • 2020-09-22
    • 1970-01-01
    相关资源
    最近更新 更多