【问题标题】:How can I update one column value in an RDD[Row]?如何更新 RDD[Row] 中的一列值?
【发布时间】:2019-04-30 19:59:38
【问题描述】:

我使用scala for spark,我想更新RDD中的一列值,我的数据格式是这样的:

[510116,8042,1,8298,20170907181326,1,3,lineno805]
[510116,8042,1,8152,20170907182101,1,3,lineno805]
[510116,8042,1,8154,20170907164311,1,3,lineno805]
[510116,8042,1,8069,20170907165031,1,3,lineno805]
[510116,8042,1,8061,20170907170254,1,3,lineno805]
[510116,8042,1,9906,20170907171417,1,3,lineno805]
[510116,8042,1,8295,20170907174734,1,3,lineno805]

我的scala代码是这样的:

 val getSerialRdd: RDD[Row]=……

我想更新包含数据20170907181326的列,我希望数据如下格式:

[510116,8042,1,8298,2017090718,1,3,lineno805]
[510116,8042,1,8152,2017090718,1,3,lineno805]
[510116,8042,1,8154,2017090716,1,3,lineno805]
[510116,8042,1,8069,2017090716,1,3,lineno805]
[510116,8042,1,8061,2017090717,1,3,lineno805]
[510116,8042,1,9906,2017090717,1,3,lineno805]
[510116,8042,1,8295,2017090717,1,3,lineno805]

并输出RDD类型,如RDD[Row]。

我该怎么做?

【问题讨论】:

  • 1) 你已经有了 RDD[Row],为什么没有数据框呢(可选问题)? 2) 要更新的列的 Row 或数据类型的架构是什么?可以发rdd.take(1)(0).schema吗?
  • 是的,我已经有 RDD[Row],数据列包含 20170907181326 是 String ,是时间列,我想从 20170907181326 列中得到 2017090718 。
  • 您是否有机会使用数据框或必须使用 RDD?并且列中的所有字符串也是相同长度的。
  • 所有相同长度的字符串。对于这个 getSerialRdd: RDD[Row],我进行了一些转换和操作,所以 getSerialRdd: RDD[Row] 是转换结果。我只想在 RDD[row] 中转换一列值,然后得到新的 RDD[row]。

标签: scala apache-spark rdd


【解决方案1】:

您可以像这样定义update 方法来更新行中的字段:

import org.apache.spark.sql.Row

def update(r: Row): Row = {
    val s = r.toSeq
    Row.fromSeq((s.take(4) :+ s(4).asInstanceOf[String].take(10)) ++ s.drop(5))
}

rdd.map(update(_)).collect

//res13: Array[org.apache.spark.sql.Row] = 
//       Array([510116,8042,1,8298,2017090718,1,3,lineno805], 
//             [510116,8042,1,8152,2017090718,1,3,lineno805], 
//             [510116,8042,1,8154,2017090716,1,3,lineno805], 
//             [510116,8042,1,8069,2017090716,1,3,lineno805], 
//             [510116,8042,1,8061,2017090717,1,3,lineno805], 
//             [510116,8042,1,9906,2017090717,1,3,lineno805], 
//             [510116,8042,1,8295,2017090717,1,3,lineno805])

更简单的方法是使用 DataFrame API 和 substring 函数:

1) 从 rdd 创建一个数据框:

val df = spark.createDataFrame(rdd, rdd.take(1)(0).schema)
// df: org.apache.spark.sql.DataFrame = [_c0: string, _c1: string ... 6 more fields]

2) 使用substring 转换列:

df.withColumn("_c4", substring($"_c4", 0, 10)).show
+------+----+---+----+----------+---+---+---------+
|   _c0| _c1|_c2| _c3|       _c4|_c5|_c6|      _c7|
+------+----+---+----+----------+---+---+---------+
|510116|8042|  1|8298|2017090718|  1|  3|lineno805|
|510116|8042|  1|8152|2017090718|  1|  3|lineno805|
|510116|8042|  1|8154|2017090716|  1|  3|lineno805|
|510116|8042|  1|8069|2017090716|  1|  3|lineno805|
|510116|8042|  1|8061|2017090717|  1|  3|lineno805|
|510116|8042|  1|9906|2017090717|  1|  3|lineno805|
|510116|8042|  1|8295|2017090717|  1|  3|lineno805|
+------+----+---+----+----------+---+---+---------+

3) 将数据帧转换为rdd很简单:

val getSerialRdd = df.withColumn("_c4", substring($"_c4", 0, 10)).rdd

【讨论】:

  • 谢谢,Psidom。我对我的程序使用第一种方法,没关系。我是火花学习的新手。我觉得有很多东西要学。
【解决方案2】:

在某些情况下,您可能希望使用架构更新行

import org.apache.spark.sql.Row
import org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema

def update(r: Row, i: Int, a: Any): Row = {

    val s: Array[Any] = r
      .toSeq
      .toArray
      .updated(i, a)

    new GenericRowWithSchema(s, r.schema)
}

rdd.map(update(_)).show(false)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-12-16
    • 1970-01-01
    • 1970-01-01
    • 2017-11-15
    • 1970-01-01
    • 2019-09-09
    • 1970-01-01
    相关资源
    最近更新 更多