KSQLDB中能否在CREATE STREAM语句内拼接列作为时间戳?
需要将RlcDate与RlcTime两列拼接后作为时间戳字段,用于计算滚动窗口(TUMBLING WINDOW),尝试在CREATE STREAM语句中直接实现拼接,编写的SQL如下:
create STREAM testStream (Price INT,concat(cast(RlcDate as varchar),LPAD(CAST(RlcTime AS VARCHAR),8,'0')) as timestamp, Id VARCHAR) WITH (kafka_topic='testTopic', partitions=1,value_format='JSON', timestamp = 'timestamp', timestamp_format = 'yyyyMMddHHmmssSS');
执行后返回错误:
Caused by: line 1:42: mismatched input '(' expecting {'STRING', 'EMIT',
'CHANGES', 'FINAL', 'ESCAPE', 'INTEGER', 'DATE', 'TIME', 'TIMESTAMP',
'INTERVAL', 'YEAR', 'MONTH', 'DAY', 'HOUR', 'MINUTE', 'SECOND', 'ZONE',
'PARTITION', 'STRUCT', 'EXPLAIN', 'ANALYZE', 'TYPE', 'TYPES', 'SHOW', 'TABLES',
'COLUMNS', 'COLUMN', 'PARTITIONS', 'FUNCTIONS', 'FUNCTION', 'ARRAY', 'MAP',
'SET', 'RESET', 'SESSION', 'DECIMAL', 'KEY', 'SINK', 'SOURCE', 'PRIMARY',
'REPLACE', 'ASSERT', 'ADD', 'ALTER', 'IF', IDENTIFIER, DIGIT_IDENTIFIER,
QUOTED_IDENTIFIER, BACKQUOTED_IDENTIFIER, VARIABLE}
Caused by: org.antlr.v4.runtime.InputMismatchException
咨询:是否可以在CREATE STREAM语句中使用CONCAT和CAST函数?若不可行,有何解决方案?
1. 核心结论
不能直接在CREATE STREAM的字段定义中使用CONCAT、CAST这类函数。流处理SQL(如KSQL)的CREATE STREAM语法规则中,字段定义仅支持字段名 数据类型的格式,不允许嵌入函数表达式,这也是你遇到语法解析错误的根本原因。
2. 可行实现方案
推荐通过两步创建流的方式完成需求:
步骤1:创建包含原始字段的基础流
先定义包含RlcDate、RlcTime等原始字段的基础流,暂不处理时间戳逻辑:
CREATE STREAM testStream_raw ( Price INT, RlcDate DATE, RlcTime TIME, Id VARCHAR ) WITH ( kafka_topic='testTopic', partitions=1, value_format='JSON' );
步骤2:创建派生流,生成拼接后的时间戳并绑定为流时间戳
基于基础流,通过SELECT语句计算拼接后的时间戳字段,同时在WITH子句中指定该字段为流的时间戳,用于后续滚动窗口计算:
CREATE STREAM testStream WITH ( timestamp = 'timestamp_col', timestamp_format = 'yyyyMMddHHmmssSS' ) AS SELECT Price, CONCAT(CAST(RlcDate AS VARCHAR), LPAD(CAST(RlcTime AS VARCHAR), 8, '0')) AS timestamp_col, Id FROM testStream_raw;
补充说明
如果不想创建中间基础流,也可以在后续的窗口查询语句中直接使用拼接表达式作为时间戳(例如在WINDOW子句中指定),但这种方式无法将时间戳绑定到流本身。若需要流级别的全局时间戳用于窗口计算,上述两步法是最简洁可靠的实现方式。
内容的提问来源于stack exchange,提问作者Fatemeh

