Flink SQL DDL定义Watermark采用何种策略?能否自主选择?
Flink 1.12 DDL水位线策略问题解答
你提供的DDL示例采用的是PeriodicWatermarks(周期性生成水位线),默认每隔200ms生成一次。通过
WATERMARK FOR user_action_time AS user_action_time - INTERVAL '5' SECOND这种表达式定义的水位线逻辑,属于周期性水位线生成机制,Flink会按固定时间间隔计算并发出水位线。在Flink 1.12的Table API中,仅靠DDL语法无法直接切换到PunctuatedWatermarks(每条新记录立即生成水位线)。如果需要使用这种策略,得通过以下方式实现:
- 先在DataStream API层面自定义
WatermarkGenerator,实现每条记录触发水位线的逻辑; - 再将该DataStream转换为Table,或是通过自定义TableSource来集成自定义的水位线生成逻辑。
- 先在DataStream API层面自定义
内容的提问来源于stack exchange,提问作者Tom
相关产品推荐
相关产品推荐

