pandas可正常运行的时间滚动窗口计算在koalas中报错
Koalas分组时间滚动窗口计算报类型错误排查
问题复现
一段基于pandas的分组时间滚动窗口统计逻辑运行正常,迁移到Koalas执行时抛出TypeError: '<' not supported between instances of 'str' and 'int'错误。
正常运行的Pandas版本
首先导入依赖、构造测试数据集:
import pandas as pd import databricks.koalas as ks Timestamp = pd.Timestamp df = pd.DataFrame([[Timestamp('2022-05-18 18:10:50.021831300'), '65.78.97'], [Timestamp('2022-05-24 09:48:39.787426700'), '65.78.97'], [Timestamp('2022-05-24 16:06:18.765405500'), '65.78.97'], [Timestamp('2022-05-25 03:04:01.860841300'), '65.78.97'], [Timestamp('2022-05-26 05:01:08.335874700'), '47.41.196'], [Timestamp('2022-05-31 03:57:15.060167500'), '47.41.196'], [Timestamp('2022-05-31 06:32:37.177199300'), '47.41.196']], columns=['time', 'ip'])
pandas侧按ip字段分组,基于时间索引执行窗口大小为24h、min_periods=1的滚动求和,统计每个IP过去24小时内的会话数量,代码如下:
print(pd.DataFrame(df).groupby('ip').apply(lambda x: pd.Series([1]*len(x), index=x['time']).rolling(window='24h', min_periods=1).sum()))
运行输出符合预期:
ip time 47.41.196 2022-05-26 05:01:08.335874700 1.0 2022-05-31 03:57:15.060167500 1.0 2022-05-31 06:32:37.177199300 2.0 65.78.97 2022-05-18 18:10:50.021831300 1.0 2022-05-24 09:48:39.787426700 1.0 2022-05-24 16:06:18.765405500 2.0 2022-05-25 03:04:01.860841300 3.0
报错的Koalas版本
将逻辑改写为Koalas接口,时间窗口参数替换为等价的1d,代码如下:
print(ks.DataFrame(df).groupby('ip').apply(lambda x: ks.Series([1]*len(x), index=x['time']).rolling(window='1d', min_periods=1).sum()))
运行抛出错误:
TypeError: '<' not supported between instances of 'str' and 'int'
排查思路
- 根因定位:Koalas的
groupby.apply对自定义函数内部的对象类型推断存在缺陷,直接用分组子集的时间列作为ks.Series索引时,时间戳类型没有被正确识别,传入时间滚动窗口逻辑后,框架会错误地将时间偏移字符串和默认生成的整数索引做大小比较,最终触发类型错误。 - 参数校验:Koalas时间滚动窗口确实不支持
24h这类小时级偏移写法,替换为等价的1d是正确操作,不是报错诱因。 - 逻辑适配:在
groupby.apply内部手动构造Series再调用rolling是pandas单机场景的写法,没有适配Koalas的分布式执行模型,本身就不属于Koalas推荐的实现方式。
解决方案
放弃在groupby.apply内部构造Series做滚动计算的写法,直接使用Koalas原生支持的分组滚动API实现,既可以避开自定义函数内的类型推断问题,执行性能也远高于自定义apply写法:
# 先将时间列转为标准时间类型、设置为索引 ks_df = ks.DataFrame(df) ks_df['time'] = ks.to_datetime(ks_df['time']) ks_df = ks_df.set_index('time') # 直接调用分组滚动接口统计窗口内记录数 result = ks_df.groupby('ip').rolling(window='1d', min_periods=1).size() print(result)
上述代码运行结果和pandas版本完全一致,不会触发类型错误。
如果业务逻辑必须在apply内实现复杂自定义计算,可以在函数内部将传入的Koalas分组子集转为pandas对象完成计算后再返回,但该方式会触发跨节点的数据拉取,性能很差,仅适合小批量数据场景。
内容的提问来源于stack exchange,提问作者Lei
相关产品推荐
相关产品推荐

