能否在PySpark Structured Streaming中使用row_number()函数?
关于PySpark Structured Streaming中row_number()函数支持性的确认
PySpark SQL官方文档对row_number()函数的定义是:
返回窗口分区内从1开始的连续编号
这意味着该函数仅能在窗口场景下使用,直接执行以下代码会触发预期的异常:
df.select('*', row_number())
异常信息:
Window function row_number() requires an OVER clause
此外,row_number()的.over()方法仅兼容WindowSpec对象,如果尝试传入结构化流原生的时间窗口函数window(),会触发类型错误:
from pyspark.sql.functions import window, row_number ... df.select('*', row_number().over(window('time', '5 minutes')))
错误信息:
TypeError: window should be WindowSpec
根据ASF Jira上的相关评论:
我们所说的时间窗口是指结构化流(SS)原生支持的时间窗口类型。
WindowSpec是不被支持的。它以非时间方式定义窗口边界和行偏移,这在流场景中难以追踪。
由于Structured Streaming通常不支持WindowSpec,由此得出row_number()函数在Structured Streaming中不被支持的结论是否正确?仅需确认是否遗漏了相关信息。
内容的提问来源于stack exchange,提问作者Kai Roesner
相关产品推荐
相关产品推荐

