【问题标题】:Inner join with Pyspark for a Cohort Study内部加入 Pyspark 进行队列研究
【发布时间】:2017-03-13 12:44:52
【问题描述】:

我正在尝试建立一个队列研究来跟踪应用内用户的行为,我想问你是否知道我在使用 .join() 时如何在 pyspark 中指定条件 给定:

rdd1 = sc.parallelize ([(u'6df99638e4584a618f92a9cfdf318cf8',
    ((u'service1',
      u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A',
      u'2016-02-08',
      u'2016-39',
      u'2016-6',
      u'2016-2',
      '2016-10-19'),
     (u'service2',
      u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A',
      u'1',
      u'67.0',
      u'2016-293',
      u'2016-42',
      u'2016-10',
      '2016-10-19')))])


rdd2 = sc.parallelize ([(u'6df99638e4584a618f92a9cfdf318cf8',
    ((u'serice1',
      u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A',
      u'2016-02-08',
      u'2016-39',
      u'2016-6',
      u'2016-2',
      '2016-10-20'),
     (u'service2',
      u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A',
      u'10',
      u'3346.0',
      u'2016-294',
      u'2016-42',
      u'2016-10',
      '2016-10-20')))])

这两个rdds代表一个用户的信息,ID为'6df99638e4584a618f92a9cfdf318cf8',在2016-10-19和2016-10-20登录了服务1和服务2。我的目标是加入我的两个 rdds,每个 rdds 至少包含 20 000 行。所以它必须是一个内部连接。真正的目标是获取所有已经在 2016-10-19' 登录并且也在 2016-10-20 登录的用户。所以更具体地说,我的最终目标是在内部连接之后得到 rxemple 的结果,只有 rdd2 的内容。

预期输出:

    [(u'6df99638e4584a618f92a9cfdf318cf8',
((u'serice1', u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A', u'2016-02-08', u'2016-39', u'2016-6', u'2016-2', '2016-10-20'), 
(u'service2', u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A', u'10', u'3346.0', u'2016-294', u'2016-42', u'2016-10', '2016-10-20'))
) ] 

一个简单的连接 rdd1.join(rdd2) 从逻辑上给我一个 RDD,其中包含与两个 rdd 匹配的所有元素对。 leftOuterJoin 或 rightOuterJoin 也不适合我的土地,因为我想要一个内部连接(只是 rdd1 和 rdd2 中已经存在的 ID)..

预期输出:假设我们有两个字典:dict1 = {'a': 'man', 'b': woman, 'c': 'baby'} 和 dict2 = {'a': 'Zara', 'x':芒果,'y':'Celio'}。预期的输出必须是:output_dict = {'a': 'Zara'}。 'a'(键)已经存在于 dict 1 中,我想要的是键,来自 dict2 的值!

它试图这样做:

rdd1.map(lambda (k, v) : k).join(rdd2)

这段代码给了我一个空的rdd。

怎么办? PS:我必须处理 rdds,而不是数据帧!所以我不想将我的 rdds 转换为 DataFrames :D 任何帮助表示赞赏。谢谢!

【问题讨论】:

  • 预期输出是什么?
  • @Yaron: [(u'6df99638e4584a618f92a9cfdf318cf8', ((u'serice1', u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A', u'2016-02-08', u'2016 -39', u'2016-6', u'2016-2', '2016-10-20'), (u'service2', u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A', u'10' , u'3346.0', u'2016-294', u'2016-42', u'2016-10', '2016-10-20')))]
  • @Yaron : rdd2 的内容。我正在寻找 rdd1 (2016-10-19) 和 rdd2 (2016-10-20) 中存在的用户。
  • 对不起,请您简单解释一下,应该对输入执行什么算法,预期输出是什么?
  • 预期输出: rdd2 的所有行都存在于 rdd1 中。假设我们有两个字典:dict1 = {'a': 'man', 'b': woman, 'c': 'baby'} and dict2 = {'a': 'Zara', 'x': Mango, 'y':'西里奥'}。预期的输出必须是:output_dict = {'a': 'Zara'}。 'a'(键)已经存在于 dict 1 中,我想要的是键,来自 dict2 的值!谢谢!

标签: python-2.7 join apache-spark pyspark rdd


【解决方案1】:

因此,您正在寻找 rdd1 和 rdd2 的连接,它将仅从 rdd2 获取键和值:

rdd_output = rdd1.join(rdd2).map(lambda (k,(v1,v2)):(k,v2))

结果是:

print rdd_output.take(1)

[(u'6df99638e4584a618f92a9cfdf318cf8', (
(u'serice1', u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A', u'2016-02-08', u'2016-39', u'2016-6', u'2016-2', '2016-10-20'), 
(u'service2', u'D8B75AA2-7408-49A7-A70D-6442C12E2B6A', u'10', u'3346.0', u'2016-294', u'2016-42', u'2016-10', '2016-10-20')
))]

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-01-28
    • 2022-09-29
    • 2015-12-20
    • 1970-01-01
    相关资源
    最近更新 更多