PySpark使用元组作为多字段key调用reduceByKey结果异常如何解决
问题根因
你的代码问题出在两个地方:
- PySpark中RDD是不可变对象,
reduceByKey是转换算子,会返回新的RDD,你没有将计算结果赋值回变量,后续collect()打印的还是执行map操作之后未聚合的原始RDD - 代码缺少
datetime模块导入,运行会直接报错
另外元组本身就是PySpark支持的合法Key类型,不需要额外做序列化或格式转换处理,不存在多字段元组无法作为Key分组的问题。
修正后代码
from pyspark import SparkContext, SparkConf import sys import datetime # 补全缺失的依赖导入 conf = SparkConf().setAppName("test") sc = SparkContext(conf=conf) from operator import add def convertion(num): return datetime.datetime.fromtimestamp(num).strftime('%Y-%m-%d') def compute(strs, num): if strs == 'apple': return -num return num rdd = sc.parallelize([ {'user':'user','tpe':'apple','timstamp':1500000000,'amount':1}, {'user':'user','tpe':'pear','timstamp':1500000001,'amount':2}, {'user':'user2','tpe':'apple','timstamp':1505000002,'amount':3} ]) rdd = rdd.map(lambda x: ((x['user'],convertion(x['timstamp'])),compute(x['tpe'],x['amount']))) # 将reduceByKey的结果赋值回rdd变量,替换之前未聚合的RDD rdd = rdd.reduceByKey(lambda x, y: x+y) print(rdd.collect())
运行后即可输出你期望的结果:[(('user', '2017-07-13'), 1), (('user2', '2017-09-09'), -3)]
内容的提问来源于stack exchange,提问作者Sailormoon
相关产品推荐
相关产品推荐

