【问题标题】:how to join two datasets by key in scala spark如何在scala spark中按键连接两个数据集
【发布时间】:2016-10-03 14:43:07
【问题描述】:

我有两个数据集,每个数据集都有两个元素。 以下是示例。

Data1:(名称,动物)

('abc,def', 'monkey(1)')
('df,gh', 'zebra')
...

Data2:(名称,水果)

('a,efg', 'apple')
('abc,def', 'banana(1)')
...

预期结果:(名称、动物、水果)

('abc,def', 'monkey(1)', 'banana(1)')
... 

我想通过使用第一列“名称”来连接这两个数据集。我已经尝试这样做了几个小时,但我无法弄清楚。谁能帮帮我?

val sparkConf = new SparkConf().setAppName("abc").setMaster("local[2]")
val sc = new SparkContext(sparkConf)
val text1 = sc.textFile(args(0))
val text2 = sc.textFile(args(1))

val joined = text1.join(text2)

以上代码不工作!

【问题讨论】:

  • 您在哪里将输入文本拆分为(key, value) 元组?
  • 你得到什么样的错误?它告诉你什么?
  • @maasg 它说''无法解析符号连接。'
  • 但是,它不是已经有 (key, value) 的格式了吗?我真的很困惑..
  • 这也可以为数据集更新吗?解决方案仅包含 RDD

标签: scala apache-spark


【解决方案1】:

join 定义在 Pair 的 RDD 上,即 RDD[(K,V)] 类型的 RDD。 所需的第一步是将输入数据转换为正确的类型。

我们首先需要将String类型的原始数据转换成(Key, Value)的对:

val parse:String => (String, String) = s => {
  val regex = "^\\('([^']+)',[\\W]*'([^']+)'\\)$".r
  s match {
    case regex(k,v) => (k,v)
    case _ => ("","")
  }
}

(请注意,我们不能使用简单的split(",") 表达式,因为键包含逗号)

然后我们使用该函数来解析文本输入数据:

val s1 = Seq("('abc,def', 'monkey(1)')","('df,gh', 'zebra')")
val s2 = Seq("('a,efg', 'apple')","('abc,def', 'banana(1)')")

val rdd1 = sparkContext.parallelize(s1)
val rdd2 = sparkContext.parallelize(s2)

val kvRdd1 = rdd1.map(parse)
val kvRdd2 = rdd2.map(parse)

最后,我们使用join方法来加入两个RDD

val joined = kvRdd1.join(kvRdd2)

//我们来看看结果

joined.collect

// res31: Array[(String, (String, String))] = Array((abc,def,(monkey(1),banana(1))))

【讨论】:

  • 非常感谢!
  • 我还有一个问题。如何在数据中保留单引号?
  • @tobby 更改正则表达式以保留引号。
  • 我试过 "^\('([^']+)',[\\W]*'([^']+)'\)$".. 但它不起作用?
  • 如果你不想修改正则表达式,你总是可以在解析后放回引号。给猫剥皮的方法有很多,有些方法比其他方法更快。
【解决方案2】:

您必须首先为您的数据集创建 pairRDD,然后您必须应用连接转换。您的数据集看起来不准确。

请考虑以下示例。

**Dataset1**

a 1
b 2
c 3

**Dataset2**

a 8
b 4

你的代码在 Scala 中应该如下所示

    val pairRDD1 = sc.textFile("/path_to_yourfile/first.txt").map(line => (line.split(" ")(0),line.split(" ")(1)))

    val pairRDD2 = sc.textFile("/path_to_yourfile/second.txt").map(line => (line.split(" ")(0),line.split(" ")(1)))

    val joinRDD = pairRDD1.join(pairRDD2)

    joinRDD.collect

这是 scala shell 的结果

res10: Array[(String, (String, String))] = Array((a,(1,8)), (b,(2,4)))

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-07-16
    • 1970-01-01
    • 2018-09-13
    • 1970-01-01
    • 1970-01-01
    • 2019-01-27
    • 2019-11-01
    相关资源
    最近更新 更多