【问题标题】:How do I find the minimum date for each unique key in an RDD with PySpark?如何使用 PySpark 找到 RDD 中每个唯一键的最短日期?
【发布时间】: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 并找到了类似的答案 hereherehere,但同样,我无法成功修改这些答案以适应我的特定情况。通常,日期时间格式似乎会出错(日期时间格式本身是由于我解析 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


【解决方案1】:

我想通了。有几个问题。对于初学者,这是适用的语法。首先(当然是在创建 SparkSession 之后),我将 RDD 转换为数据帧:

df = spark.createDataFrame(df, ['col1', 'col2'])

然后进行了groupBy和聚合功能。这些你会在其他 SO 答案中看到,但我想我会在这里发布这个,因为它是在我的特定场景的上下文中。

from pyspark.sql import functions as F
df= df.groupBy('col1').agg(F.min('col2'))

然后将数据返回到 RDD 格式,我使用了

result = df.rdd.map(lambda x: (x[0], x[4]))

在这种特殊情况下,我还将数据帧的第 0 列和第 4 列的元素映射回分配给 result 的元组。

在这个过程中,我还发现了一些可能对其他人有帮助的有趣点:

  1. 我的数据框中有一些我不知道的 Null 值 不断导致,主要是“NoneType 不可下标” 错误。虽然这是有道理的,但我花了一段时间才弄清楚 NoneType 所在的位置。
  2. 我的一些 XML 被错误地解析,因此它返回的是 (None) 元组,而不是上述数据格式所要求的 (None, None) 元组。

这些更正使我能够.show() 数据帧(而不仅仅是.printSchema().groupBy 及其关联对象从来都不是问题。

【讨论】:

  • @jxc 在 cmets 中的回答可能在我解决了此答案中描述的“无”问题后有效,但我没有尝试尝试。
猜你喜欢
  • 2018-10-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-01-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多