【问题标题】:Spark : subtractByKey issue (pyspark)Spark:subtractByKey 问题 (pyspark)
【发布时间】:2016-11-22 12:20:42
【问题描述】:

我有一些关于subtractByKey 的问题。

我有 2 个文件: 第一个是这样的:(客户 ID + 客户邮件)

client_id   emails

4A85FD8E-197D-2AE3-B939-A527AFF16A04    imperdiet.non.vestibulum@mon***tur.com 
D48D530C-CF68-DAF1-18F0-E0A0A03F3E06    rutrum.urna@estm***ncus.net:facilisis@i****m.ca
40815230-25DC-9EA0-01D1-2706B4B56958    iaculis.nec.eleifend@gr****nc.net 
...

第二个:(仅限邮件)

pharetra@P****s.com  
ut.aliquam@o****m.org  
erat@a****e.edu 
....

第一个文件中的某些行可以包含 2 封(或更多)这种格式的邮件:

mail:mail

我做了什么:

*test1=sc.textFile("file1")
*test2=sc.textFile("file2")
*test3=test1.subtractByKey(test2)

结果是……:

[(u'A', u'B'), (u'A', u'D'), (u'A', u'1'), (u'A', u'D'), (u'A', u'D'), (u'A', u'B'), (u'A', u'F'), (u'A', u'E'), (u'A', u'9'), (u'A', u'5'), (u'A', u'9'), (u'A', u'6'), (u'c', u'l'), (u'E', u'8'), (u'E', u'4'), (u'E', u'6'), (u'E', u'6'), (u'E', u'7'), (u'E', u'5'), (u'E', u'5'), (u'E', u'5'), (u'E', u'2'), (u'E', u'8'), (u'C', u'2'), (u'C', u'5'), (u'C', u'6'), (u'C', u'C'), (u'C', u'E'), (u'C', u'3'), (u'C', u'F'), (u'C', u'4'), (u'C', u'B'), (u'C', u'F'), (u'C', u'F'), (u'C', u'8'), (u'C', u'0'), (u'1', u'D'), (u'1', u'2'), (u'1', u'3'), (u'1', u'8'), (u'1', u'0'), (u'1', u'F'), ... ]

我想删除第一个文件中的客户,他们的邮件在第二个文件中,但它不起作用。

【问题讨论】:

  • 你最好用code格式编辑你的问题,因为我看到的很混乱
  • 里面没有代码...除了3行,是代码格式。
  • 这是第一个可以包含多封电子邮件的文件,对吧?
  • 是的,对不起。正在编辑!

标签: pyspark


【解决方案1】:

注意:我对 pyspark 不是很熟悉,但 spark api 应该是 一样。

首先你应该把电子邮件作为关键

rdd1=sc.textFile("file1").map(lambda line: (line.split(" ")[0], line.split(" ")[1]))

这会给你一个 rdd

[(4A85FD8E-197D-2AE3-B939-A527AFF16A04,imperdiet.non.vestibulum@mon***tur.com)]

那么可能有多个电子邮件,你应该做一个flatMapValues()

rdd2 = rdd1.flatMapValues(lambda email: email.split(":"))

这将为您提供一对 rdd,每个仅包含一封电子邮件

现在您可以切换键和值

rdd3=rdd2.map(lambda kv: (kv[1], kv[0]))

现在你得到一个使用用户电子邮件作为键和 UUID 作为值的 rdd 比如

[(imperdiet.non.vestibulum@mon***tur.com, 4A85FD8E-197D-2AE3-B939-A527AFF16A04)]

现在您应该找到哪个 UUID 的电子邮件包含在 file2 中,为此您应该将第二个文件加载为 rdd:

secondRdd = sc.textFile("file2").map(lambda line: (line, 1))

然后执行join 并调整连接结果 rdd。

rdd4 = rdd3.join(secondRdd).map(lambda kv: (kv[1][0], kv[0]))

如果现在一切正常,您应该得到一个格式为(UUID, email) 的rdd,它代表电子邮件出现在file2 中的所有用户,

然后你可以用我们最初得到的rdd1 做一个subtractByKey()

【讨论】:

  • 好的,非常感谢。如果我可以添加一个问题:如果电子邮件在 file2 中,它将删除客户端。好的,但是如果客户有 2 封邮件,它会删除其中一个还是两个都删除? (我猜你的平面图复制了客户端?)
  • 你应该注意到,rdd1 还不是 flatMap,是的,我认为如果有重复的 'UUID',subtractByKey()ease all
猜你喜欢
  • 2017-12-06
  • 1970-01-01
  • 1970-01-01
  • 2018-04-13
  • 2021-06-16
  • 2019-04-19
  • 2017-10-27
  • 2021-01-15
  • 1970-01-01
相关资源
最近更新 更多