【发布时间】:2020-03-08 15:24:58
【问题描述】:
我有一个格式为 [(ID, Date), (ID, Date)...] 的元组列表,其中日期为日期时间格式。作为我正在使用的 RDD 的一个示例:
[('1', datetime.date(2012, 1, 01)),
('2', datetime.date(2012, 1, 01)),
('3', datetime.date(2012, 1, 01)),
('4', datetime.date(2012, 1, 01)),
('5', datetime.date(2012, 1, 01)),
('1', datetime.date(2011, 1, 01)),
('2', datetime.date(2013, 1, 01)),
('3', datetime.date(2015, 1, 01)),
('4', datetime.date(2010, 1, 01)),
('5', datetime.date(2018, 1, 01))]
我需要收集 ID 以及与每个 ID 关联的最短日期。据推测,这是一个reduceByKey 操作,但我一直无法理清关联的功能。我猜我只是把事情复杂化了,但是在确定适当的 lambda 时会得到帮助(或者如果 reduceByKey 在这种情况下不是最有效的方法)。
我搜索了 StackOverflow 并找到了类似的答案 here、here 和 here,但同样,我无法成功修改这些答案以适应我的特定情况。通常,日期时间格式似乎会出错(日期时间格式本身是由于我解析 xml 的方式造成的,因此如果有帮助,我可以返回并将其解析为字符串)。
我尝试了以下方法,但都收到了错误:
.reduceByKey(min) - IndexError: 元组索引超出范围
reduceByKey(lambda x, y: (x, min(y))) - IndexError: 元组索引超出范围(如果日期时间被转换为字符串,或者如果是日期时间格式则出现以下错误)
.reduceByKey(lambda x, y: (x[0], min(y))) - TypeError: 'datetime.date' 对象不可下标
我希望最终结果如下:
[('1', datetime.date(2011, 1, 01)),
('2', datetime.date(2012, 1, 01)),
('3', datetime.date(2012, 1, 01)),
('4', datetime.date(2010, 1, 01)),
('5', datetime.date(2012, 1, 01))]
【问题讨论】:
-
试试
rdd.reduceByKey(lambda x,y: x if x<y else y) -
@jxc -
IndexError: tuple index out of range当我尝试.take()时。 -
你的RDD元素是(id, date)键值对吗?
-
@jxc - 我理解它们是键值对。但是,在 Spark 中究竟是什么指定了它们呢?我认为键始终是第一个元素,值是 2 元素(1 元组/行)RDD 中的第二个...
-
你是对的,在 Pair RDD 中,元组的第一项始终是键,第二项是值。在 reduceByKey(lambda x, y:) 中,x 和 y 都是值,所以你不应该使用 x[0] 或 x[1] 等。x 和 y 都是 datetime.date 对象。
标签: python-3.x apache-spark pyspark rdd