【问题标题】:How to read values from two csv and do operation on b/w its column in spark java api?如何从两个csv中读取值并对spark java api中的b/w列进行操作?
【发布时间】:2017-03-22 23:43:39
【问题描述】:

我在 hadoop 中有两个 Csv,比如 csv1、csv2。两个 csv 都包含两列(时间戳和某个值),例如 csv1 的列是 t1、v1,而 csv2 的列是 t2、v2。 我想为每个 t1 = t2(对于相同的时间戳)计算 v1*v2,并使用 spark java Api 将结果作为文本文件存储在 hdfs 中。

我是新来的火花,请有人帮助我。

提前感谢。

【问题讨论】:

    标签: java hadoop spark-java bigdata


    【解决方案1】:

    我可以在 scala 中做到,也许你可以了解我正在做的事情并自己实现它:

    scala> val df1=sc.parallelize(Seq((1001,2),(1002,3),(1003,4))).toDF("t1","v1")
    df1: org.apache.spark.sql.DataFrame = [t1: int, v1: int]
    
    
    scala> val df2=sc.parallelize(Seq((1001,3),(1002,4),(1005,4))).toDF("t2","v2")
    df2: org.apache.spark.sql.DataFrame = [t2: int, v2: int]
    
    scala> df1.join(df2,df1("t1")===df2("t2"))
    res1: org.apache.spark.sql.DataFrame = [t1: int, v1: int ... 2 more fields]
    
    scala> res1.show
    +----+---+----+---+                                                             
    |  t1| v1|  t2| v2|
    +----+---+----+---+
    |1002|  3|1002|  4|
    |1001|  2|1001|  3|
    +----+---+----+---+
    
    scala> import org.apache.spark.sql.functions._
    import org.apache.spark.sql.functions._
    
    scala> val result=res1.withColumn("foo",res1("v1") * res1("v2"))
    result: org.apache.spark.sql.DataFrame = [t1: int, v1: int ... 3 more fields]
    
    scala> result.show
    +----+---+----+---+---+                                                         
    |  t1| v1|  t2| v2|foo|
    +----+---+----+---+---+
    |1002|  3|1002|  4| 12|
    |1001|  2|1001|  3|  6|
    +----+---+----+---+---+
    

    我希望这能解决你的问题。

    【讨论】:

    • 感谢解决方案,我尝试了这些概念,但没有得到确切的解决方案。时间戳列包含诸如 2016-09-01 15:31:58+00:00 之类的值。我想加载 csv 并将其拆分为列,结果应该类似于 (t1,v*v2)。
    • 然后先用convert成spark时间戳再做这些步骤,或者如果你想用简单的方式做,就用string代替。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-01-03
    • 2014-05-07
    • 2017-03-12
    • 1970-01-01
    • 1970-01-01
    • 2021-05-08
    • 1970-01-01
    相关资源
    最近更新 更多