KSQL无头模式执行查询报错:STREAM_NAME不存在,交互模式正常
问题背景
你遇到的这个问题我之前也碰到过:明明STREAM_NAME已经在交互式模式中确认存在,但在KSQL 4.1.0版本的Non-Interactive (Headless)模式下,用--queries-file启动服务器执行建表查询时,却被提示流不存在。
你的启动命令如下:
$path-to-ksql/bin/ksql-server-start \ $path-to-ksql/etc/ksql/ksql-server.properties \ --queries-file /tmp/ksql-queries/queries.sql \ >path-to-logdirectory/ksql-server-1_`date '+%Y%m%d_%H_%M_%S'`.log 2>&1 &
查询文件中的语句为:
create table TABLE_NAME as select a, min(b) from STREAM_NAME WINDOW TUMBLING (size 1 minute) group by a;
日志抛出的异常是:
Exception in thread "main" io.confluent.ksql.parser.exception.ParseFailedException: Parsing failed on KsqlEngine msg: STREAM_NAME does not exist. at io.confluent.ksql.KsqlEngine.parseQueries(KsqlEngine.java:278) at io.confluent.ksql.KsqlEngine.createQueries(KsqlEngine.java:593)
而且相同的查询在交互式模式下完全可以正常执行。
问题根源
在KSQL 4.1.0的无交互模式中,服务器启动时会立刻解析并执行--queries-file里的语句,但此时KSQL还没完成已有元数据(比如你在交互式模式创建的流)的加载。另外,无交互模式默认会初始化一个独立的会话上下文,和交互式模式的元数据存储是隔离的,除非你主动配置让两者共享。
可行的解决方案
1. 在查询文件中预先声明流
最简单的处理方式是在queries.sql开头添加DECLARE STREAM语句,明确告知KSQL这个流的结构,这样解析查询时就不会报错:
-- 替换成你STREAM_NAME的实际字段名和对应数据类型 DECLARE STREAM STREAM_NAME (a VARCHAR, b INT); create table TABLE_NAME as select a, min(b) from STREAM_NAME WINDOW TUMBLING (size 1 minute) group by a;
2. 配置KSQL加载已有元数据
修改ksql-server.properties配置文件,添加以下参数,让无交互模式启动时加载之前的元数据:
# 确保这个ID和你交互式模式使用的一致,默认值是default_ ksql.service.id=default_ # 替换成你的Schema Registry实际地址 ksql.schema.registry.url=http://your-schema-registry:8081 # 指向KSQL存储元数据的目录 ksql.streams.state.dir=/path/to/ksql-state-directory
重启服务器后,它就能识别到之前创建的STREAM_NAME流了。
3. 先启动服务器再执行查询
如果你的查询数量不多,可以拆分操作步骤:先启动KSQL服务器,等它加载完元数据后,再用CLI的--execute参数执行建表语句:
首先启动服务器:
$path-to-ksql/bin/ksql-server-start $path-to-ksql/etc/ksql/ksql-server.properties >path-to-logdirectory/ksql-server-1_`date '+%Y%m%d_%H_%M_%S'`.log 2>&1 &
等待服务器启动完成(可以通过日志确认),然后执行查询:
$path-to-ksql/bin/ksql http://localhost:8088 --execute "create table TABLE_NAME as select a, min(b) from STREAM_NAME WINDOW TUMBLING (size 1 minute) group by a;"
这种方式和交互式模式共享上下文,自然能识别到已存在的流。
验证方法
完成操作后,你可以通过KSQL CLI连接到服务器,执行SHOW STREAMS;确认STREAM_NAME存在,再执行SHOW TABLES;检查TABLE_NAME是否成功创建,同时查看日志有没有新的异常信息。
内容的提问来源于stack exchange,提问作者Rahul Bharti

