Nexmark Q4改造:如何从Interval Join传播时间属性做Window Aggregation?
修改Nexmark Query 4生成追加流时的时间属性问题
需求
修改Nexmark基准测试中的Query 4(品类平均价格查询),尽可能仅使用Flink SQL生成品类平均价格的追加流(原查询输出为更新流)。
当前实现步骤
- 对auction与bid流执行Interval Join,按
auction.id和category分组取最高bid.price得到商品售价,此分组操作会产生更新和删除流; - 尝试将上述结果通过Windowed Aggregation转换为追加流,但无法从Interval Join的结果中传递时间属性,执行SQL时触发以下错误:
org.apache.flink.table.api.ValidationException: The window function requires the timecol is a time attribute type, but is TIMESTAMP(3)
已排除的方案
- Window Join:不符合业务重叠区间逻辑
- Upsert Sink:并非最终目标
内容的提问来源于stack exchange,提问作者Vignesh Chandramohan
相关产品推荐
相关产品推荐

