【发布时间】: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