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