我试图编写一个程序,使我能够在Scala中对Kafka主题运行预定义的KSQL操作,但我不想每次都打开KSQL Cli。因此,我想从Scala程序中启动KSQL“服务器”。如果我正确理解KSQL源代码,我必须构建并启动KsqlRestApplication:
def restServer = KsqlRestApplication.buildApplication(new
KsqlRestConfig(defaultServerProperties), true, new VersionCheckerAgent
{override def start(ksqlModuleType: KsqlModuleType, properties:
Properties): Unit = ???})
但是,当我尝试这样做时,会出现以下错误:
Exception in thread "main" java.lang.NoSuchMethodError: org.apache.kafka.streams.StreamsConfig.getConsumerConfigs(Ljava/lang/String;Ljava/lang/String;)Ljava/util/Map;
at io.confluent.ksql.rest.server.BrokerCompatibilityCheck.create(BrokerCompatibilityCheck.java:62)
at io.confluent.ksql.rest.server.KsqlRestApplication.buildApplication(KsqlRestApplication.java:241)
我查看了BrokerCompatibilityCheck中的函数调用,以及它调用StreamsConfig的create函数。getConsumerConfigs(),使用2个字符串作为参数,而不是中定义的参数
https://kafka.apache.org/0102/javadoc/org/apache/kafka/streams/StreamsConfig.html#getConsumerConfigs(StreamThread,%20java.lang.String,%20java.lang.String)
。
我的KSQL和卡夫卡版本是不兼容还是我做错了什么?
我使用的是KSQL版本4.1.0-SNAPSHOT和Kafka版本1.0.0。