【问题标题】:PySpark - use datetime object with a PandasUDFType.GROUPED_MAPPySpark - 使用带有 PandasUDFType.GROUPED_MAP 的日期时间对象
【发布时间】:2020-01-16 18:51:57
【问题描述】:

我创建了一个 PandasUDF 来返回每个 ID 的最新“计数”。 spark DF 中的“日期”列是字符串类型(YYYY-mm-dd)。在下面的函数中,我使用 pd.to_datetime 将字符串转换为 datetype 以获取每个 ID 的 max(date)。该函数(下)在应用于熊猫数据框时工作得很好。但是当我尝试在 spark 中使用它时,出现以下错误。

AttributeError("Can only use .dt accessor with datetimelike " "values")

我尝试先将日期列转换为 DateType(),但错误仍然存​​在。

@pandas_udf("id string, count int", PandasUDFType.GROUPED_MAP)
def recent_date(pdf):
    pdf['date'] = pd.to_datetime(pdf.date)
    latest_data = (pdf[pdf['date'] == max(pdf['date'])]).copy()
    return latest_data[['id', 'count']]

我正在使用以下调用调用该函数:

df.groupby('id').apply(recent_date)

任何帮助将不胜感激。谢谢。

【问题讨论】:

  • 你为什么要为此应用pandas_udf?这可以在原生 pyspark 本身中轻松完成。
  • 这只是我的问题的精简版,我使用的 pandas_udf 要复杂得多
  • 作为pandas_udf 的输入和输出的数据帧的架构需要具有相同的架构。
  • 我不明白,你是说@pandas_udf("id string, count int", PandasUDFType.GROUPED_MAP) 需要包含日期列吗?
  • 是的。请关注此documentation link 并阅读GROUPED_MAP 聚合

标签: python pandas pyspark


【解决方案1】:

根据answer 和检查supported types,当前的 pandas_udf 不支持带有分组映射 UDF 的 date 类型(但奇怪的是我可以以某种方式使用带有 date 类型的分组聚合 UDF,不确定它是否因为在我的情况下它没有遇到任何类型检查逻辑)。

我所做的只是将 date 类型列(在您的情况下为“日期”)转换为 timestamp 类型,然后它对我有用。

df.withColumn('date', unix_timestamp(col('date'), "yyyy-MM-dd").cast("timestamp")) \
    .groupBy('id').apply(recent_date)

希望这会有所帮助。

【讨论】:

  • 我最终在此过程的早期将 datetype 更改为 string 类型,并且它起作用了。感谢您的回答。
猜你喜欢
  • 1970-01-01
  • 2021-09-19
  • 2017-08-20
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-08-27
  • 1970-01-01
相关资源
最近更新 更多