【问题标题】:Joining two HDFS files in in Spark在 Spark 中加入两个 HDFS 文件
【发布时间】:2023-03-16 20:50:02
【问题描述】:

我想使用 spark shell 连接来自 HDFS 的两个文件。 这两个文件都是制表符分隔的,我想加入第二列

尝试过的代码 但没有给出任何输出

val ny_daily= sc.parallelize(List("hdfs://localhost:8020/user/user/NYstock  /NYSE_daily"))

val ny_daily_split = ny_daily.map(line =>line.split('\t'))

val enKeyValuePair = ny_daily_split.map(line => (line(0).substring(0, 5), line(3).toInt))


val ny_dividend= sc.parallelize(List("hdfs://localhost:8020/user/user/NYstock/NYSE_dividends"))

val ny_dividend_split = ny_dividend.map(line =>line.split('\t'))

val enKeyValuePair1 = ny_dividend_split.map(line => (line(0).substring(0, 4),     line(3).toInt))

enKeyValuePair1.join(enKeyValuePair)

但我没有得到任何关于如何在特定列上加入文件的信息 请推荐

【问题讨论】:

    标签: scala hadoop apache-spark


    【解决方案1】:

    我没有得到任何关于如何在特定列上加入文件的信息

    RDD 在它们的键上连接,因此您在编写时决定要连接的列:

    val enKeyValuePair = ny_daily_split.map(line => (line(0).substring(0, 5), line(3).toInt))
    ...
    val enKeyValuePair1 = ny_daily_split.map(line => (line(0).substring(0, 4), line(3).toInt))
    

    您的 RDD 将加入来自 line(0).substring(0, 5)line(0).substring(0, 4) 的值。

    您可以找到 join 函数(以及许多其他有用的函数)hereSpark Programming Guide 是了解 Spark 工作原理的绝佳参考。

    尝试过代码但没有给出任何输出

    为了看到输出,你必须让 Spark 打印出来:

    enKeyValuePair1.join(enKeyValuePair).foreach(println)
    

    注意:要从文件中加载数据,您应该使用sc.textFile()sc.parallelize() 仅用于从 Scala 集合中生成 RDD。

    以下代码应该可以完成这项工作:

    val ny_daily_split = sc.textFile("hdfs://localhost:8020/user/user/NYstock/NYSE_daily").map(line =>line.split('\t'))
    val ny_dividend_split = sc.textFile("hdfs://localhost:8020/user/user/NYstock/NYSE_dividends").map(line =>line.split('\t'))
    
    val enKeyValuePair = ny_daily_split.map(line => line(0).substring(0, 5) -> line(3).toInt)
    val enKeyValuePair1 = ny_dividend_split.map(line => line(0).substring(0, 4) -> line(3).toInt)
    
    enKeyValuePair1.join(enKeyValuePair).foreach(println)
    

    对了,你提到你想加入第二列,但你实际上使用的是line(0),这是故意的吗?

    希望这会有所帮助!

    【讨论】:

    • 我应该在 JOIN 的键和值中加入什么,因为我想加入列,作为输出,我应该能够看到整个加入的数据集
    • 然后将您的 map 函数分别更改为 ny_daily_split.map(line => line(1) -> line.mkString("\t"))ny_dividend_split.map(line => line(1) -> line.mkString("\t"))
    猜你喜欢
    • 2018-11-18
    • 1970-01-01
    • 2016-12-05
    • 2016-01-24
    • 1970-01-01
    • 2016-01-11
    • 2017-11-08
    • 2015-08-18
    • 2017-06-26
    相关资源
    最近更新 更多