Apache Beam Python SDK SqlTransform HOP窗口类型报错咨询
问题说明
你在Apache Beam SqlTransform中实现HOP滑动窗口计算时,编写SQL如下:
SELECT f_timestamp, line, COUNT(*) FROM PCOLLECTION GROUP BY line, HOP(f_timestamp, INTERVAL '30' MINUTE, INTERVAL '1' HOUR)
当前Python侧Row转换代码为:
| "Create beam Row" >> beam.Map(lambda x: beam.Row(f_timestamp= float(x["timestamp_date"]), line = unicode(x["line"])))
运行时Java侧Calcite校验抛出如下错误:
Caused by: org.apache.beam.vendor.calcite.v1_20_0.org.apache.calcite.sql.validate.SqlValidatorException: Cannot apply 'HOP' to arguments of type 'HOP(<DOUBLE>, <INTERVAL MINUTE>, <INTERVAL HOUR>)'. Supported form(s): 'HOP(<DATETIME>, <DATETIME_INTERVAL>, <DATETIME_INTERVAL>)'
已尝试两种修复方式均未解决问题:
- 将
f_timestamp设置为float类型的UNIX时间戳 - 将
f_timestamp设置为unicode类型的字符串时间戳
已知Java侧时间戳使用java.util.Date类型,需要完成类型适配解决报错。
报错原因
Calcite的HOP窗口函数第一个参数强制要求为DATETIME类型,你之前传入的float(对应SQL DOUBLE类型)、普通字符串(对应SQL VARCHAR类型)都不会被Calcite隐式转换为DATETIME类型,因此无法通过函数签名校验。
Python SDK的beam.Row不会自动将数值、普通字符串格式的时间戳转换为Java侧识别的DATETIME映射类型,必须显式做类型适配。
修复方案
方案1:Python侧构造Row时直接传入Beam原生Timestamp类型
导入Beam的Timestamp工具类,将时间戳显式转换为Beam原生时间类型,该类型会被SqlTransform自动映射为Calcite支持的DATETIME类型,修改后代码如下:
import apache_beam as beam from apache_beam.utils.timestamp import Timestamp | "Create beam Row" >> beam.Map( lambda x: beam.Row( # 若timestamp_date为毫秒级UNIX时间戳,需除以1000转为秒级再传入 f_timestamp=Timestamp(seconds=float(x["timestamp_date"])), line=str(x["line"]) ) )
该方案不需要修改原有SQL逻辑,类型匹配最稳定。
方案2:SQL逻辑中显式做时间类型转换
如果不想调整Python侧转换逻辑,可以直接在SQL中使用TIMESTAMP函数,将传入的字段显式转换为DATETIME类型,修改后SQL如下:
SELECT f_timestamp, line, COUNT(*) FROM PCOLLECTION GROUP BY line, -- 秒级数值时间戳用CAST转BIGINT再转TIMESTAMP;毫秒级需先除以1000;ISO格式字符串直接传TIMESTAMP(f_timestamp)即可 HOP(TIMESTAMP(CAST(f_timestamp AS BIGINT)), INTERVAL '30' MINUTE, INTERVAL '1' HOUR)
注意:如果传入的是毫秒级UNIX时间戳,必须先做除以1000的换算,否则会出现时间偏移问题。
内容的提问来源于stack exchange,提问作者Metehan Yıldırım
相关产品推荐
相关产品推荐

