【问题标题】:Incremental update in rdd or dataframe apache sparkrdd 或数据帧 apache spark 中的增量更新
【发布时间】:2019-01-07 02:36:44
【问题描述】:

我有一个用例,其中我有一组数据(例如:一个包含大约 1000 万行和大约 25 列的 csv 文件)。 我有一组规则(大约 1000 条规则),我需要更新记录,这些规则必须按顺序执行。

我写了一个代码,我在其中循环每个规则,并为每个规则更新数据。

假设规则是这样的

col1=5 和 col2=10 然后 col25=updatedValue

rulesList.foreach(rule=> {
    var data = data.map(line(col1, col2, .., col25) => if(rule){
        line(col1, col2, .., updatedValue)
    } else {line(col1, col2, .., col25)})
})

这些规则将按顺序执行,最后 a 将获得更新的记录。

但问题是,如果规则和数据少于它正确执行但如果数据大于我得到 ​​StackOverflow 错误,原因可能是因为它正在映射所有规则并像 map-reduce 一样最后执行它。

有什么方法可以让我逐步更新这些数据。

【问题讨论】:

    标签: scala apache-spark dataframe apache-spark-sql


    【解决方案1】:

    尝试在 RDD 上进行一次映射,然后循环遍历映射内的规则,从而减少数据移动。所有规则都将在数据本地应用,从而导致更新记录 - 而不是创建 1000 个 RDD

    【讨论】:

      【解决方案2】:

      给定 RDD 中的一条记录,如果您可以增量地应用所有更新但独立于其他记录,我建议您先执行映射,然后遍历映射内的 rulesList:

      val result = data.map { case line(col1, col2, ..., col25) => 
          var col25_mutable = col25
          rulesList.foreach{ rule => 
              col25_mutable = if(rule) updatedValue else col25_mutable
          }
          line(col1, col2, ..., col25_mutable)
      }
      

      如果 rulesList 是一个简单的可迭代对象,例如 Array 或 List,这种方法应该是线程安全的。

      我希望它对您有用,或者至少可以帮助您实现目标。

      干杯

      【讨论】:

        猜你喜欢
        • 2023-03-17
        • 2017-02-24
        • 1970-01-01
        • 1970-01-01
        • 2023-03-26
        • 2021-12-06
        • 2016-11-04
        • 1970-01-01
        • 2015-06-14
        相关资源
        最近更新 更多