【问题标题】:List (or iterator) of tuples returned by MAP (PySpark)MAP (PySpark) 返回的元组列表(或迭代器)
【发布时间】:2014-02-01 12:03:31
【问题描述】:

我有一个映射器方法:

def mapper(value):
    ...
    for key, value in some_list:
        yield key, value

实际上,我需要的与普通的字数示例相差不远。我已经有了工作脚本,但前提是映射器方法看起来像这样:

def mapper(value):
    ...
    return key, value

它的调用是这样的:

sc.textFile(sys.argv[2], 1).map(mapper).reduceByKey(reducer).collect()

我花了 2 个小时尝试编写支持 mapper 中的生成器的代码。但不能那样做。我什至同意只返回一个列表:

def mapper(value):
    ...
    result_list = []
    for key, value in some_list:
        result_list.append( key, value )
    return result_list

这里:https://groups.google.com/forum/#!searchin/spark-users/flatmap$20multiple/spark-users/1WqVhRBaJsU/-D5QRbenlUgJ 我发现我应该使用 flatMap,但它没有解决问题 - 然后我的减速器开始获取 (key1, value1, key2, value2, value3, ...) 之类的输入 - 但它应该是 [(key1, value1 ), (key2, value2, value3)...]。也就是说,reducer 开始只取单件,不知道是 value 还是 key,如果 value - 它属于哪个 key。

那么如何使用返回迭代器或列表的映射器呢?

谢谢!

【问题讨论】:

  • P.S.与传统的 map-reduce 框架(我们都知道我在说什么)相比,Spark 看起来快速且有前途,但确实很奇怪……而且它确实缺少更多 Python 示例,因为我在那里没有找到任何关于我的问题的信息。也许在有人帮助我之后我会提交一些东西=)

标签: python apache-spark


【解决方案1】:

如果你想要一个返回多个输出的映射函数,你可以使用flatMap

传递给flatMap的函数可以返回一个可迭代对象:

>>> words = sc.textFile("README.md")
>>> def mapper(line):
...     return ((word, 1) for word in line.split())
...
>>> words.flatMap(mapper).take(4)
[(u'#', 1), (u'Apache', 1), (u'Spark', 1), (u'Lightning-Fast', 1)]
>>> counts = words.flatMap(mapper).reduceByKey(lambda x, y: x + y)
>>> counts.take(5)
[(u'all', 1), (u'help', 1), (u'webpage', 1), (u'when', 1), (u'Hadoop', 12)]

也可以是生成器函数:

>>> words = sc.textFile("README.md")
>>> def mapper(line):
...     for word in line.split():
...         yield (word, 1)
...
>>> words.flatMap(mapper).take(4)
[(u'#', 1), (u'Apache', 1), (u'Spark', 1), (u'Lightning-Fast', 1)]
>>> counts = words.flatMap(mapper).reduceByKey(lambda x, y: x + y)
>>> counts.take(5)
[(u'all', 1), (u'help', 1), (u'webpage', 1), (u'when', 1), (u'Hadoop', 12)]

您提到您尝试了flatMap,但它将所有内容都压缩为一个列表[key, value, key, value, ...],而不是一个列表[(key, value), (key, value)...]的键值对。我怀疑这是您的地图功能中的问题。如果您仍然遇到此问题,能否发布更完整的地图功能版本?

【讨论】:

  • 谢谢乔希!虽然现在我遇到了 java 的“内存不足”错误 =) 好吧会解决的
猜你喜欢
  • 1970-01-01
  • 2015-09-23
  • 2019-09-26
  • 1970-01-01
  • 1970-01-01
  • 2013-01-30
  • 2011-04-02
  • 2022-06-29
相关资源
最近更新 更多