Quarkus应用中Java版kSqlDB客户端问题:BatchedQueryResult的completableFuture get()方法挂起
看起来你遇到了一个典型的阻塞调用与事件循环线程冲突的问题,结合Quarkus的运行环境特性,我来帮你梳理可能的原因和解决办法。
先还原下你的核心问题:
你用kSqlDB Java客户端执行pull查询时,代码卡在List<Row> resultRows = batchedQueryResult.get();这一行彻底挂起。虽然kSqlDB服务器日志显示返回了200响应,但客户端这边一直处于等待状态,且没有任何异常抛出。你已经确认了:
- kSqlDB服务器正常运行(
/info接口返回状态正常) - 目标物化表
QUERYABLE_CLIENTS_VEHICLES存在,直接在Confluent控制台执行查询无问题 - 手动调用
POST /query-streamREST接口也能正常拿到结果
你的代码参考
import io.confluent.ksql.api.client.BatchedQueryResult; import io.confluent.ksql.api.client.Client; import io.confluent.ksql.api.client.ClientOptions; import io.confluent.ksql.api.client.Row; import org.junit.jupiter.api.Test; import java.util.List; import java.util.Map; import java.util.Properties; import java.util.concurrent.ExecutionException; public class KSqlDBTest { public static String KSQLDB_SERVER_HOST = "localhost"; public static int KSQLDB_SERVER_HOST_PORT = 8088; @Test public void kSqlDbTest() { ClientOptions options = ClientOptions.create() .setHost(KSQLDB_SERVER_HOST) .setPort(KSQLDB_SERVER_HOST_PORT) .setUseTls(false) .setUseAlpn(false) .setExecuteQueryMaxResultRows(1000000000); Properties properties = new Properties(); properties.put("auto.offset.reset", "earliest"); Client client = Client.create(options); try { String pullQuery = "SELECT * FROM QUERYABLE_CLIENTS_VEHICLES;"; BatchedQueryResult batchedQueryResult = client.executeQuery(pullQuery, (Map) properties); List<Row> resultRows = batchedQueryResult.get(); // 代码在此处挂起 System.out.println("Received results. Num rows: " + resultRows.size()); for (Row row : resultRows) { System.out.println("Row: " + row.values()); } } catch (InterruptedException e) { e.printStackTrace(); } catch (ExecutionException e) { e.printStackTrace(); } catch (Exception e) { e.printStackTrace(); } } }
可能的原因与解决方案
1. Quarkus事件循环线程被阻塞(最可能的根源)
Quarkus基于Vert.x构建,默认情况下测试方法和请求处理都会在事件循环线程上执行。而batchedQueryResult.get()是一个强制阻塞的调用,会占用事件循环线程直到Future完成。但kSqlDB Java客户端本身也是基于Vert.x实现的,它需要事件循环线程来处理服务器返回的响应数据流——当你阻塞了事件循环线程,客户端无法处理服务器的响应,导致Future永远无法完成,代码自然一直挂起。
解决办法:
- 改用异步处理方式:避免直接调用
get(),改用CompletableFuture的异步回调逻辑:
batchedQueryResult.thenAccept(resultRows -> { System.out.println("Received results. Num rows: " + resultRows.size()); for (Row row : resultRows) { System.out.println("Row: " + row.values()); } }).exceptionally(e -> { e.printStackTrace(); return null; }); // 测试场景下需要添加等待逻辑,比如用CountDownLatch确保异步任务完成
- 标记方法为阻塞执行:在Quarkus测试或业务方法上添加
@Blocking注解,让Quarkus把这个方法放到Worker线程池中执行,避免阻塞事件循环:
@Test @Blocking // 添加该注解 public void kSqlDbTest() { // 原代码保持不变 }
2. ALPN配置与HTTP/2协议冲突
从服务器日志可以看到请求用的是HTTP/2.0协议,但你在客户端配置里设置了setUseAlpn(false)。ALPN是HTTP/2协议协商的关键机制,虽然服务器返回了200响应,但后续的数据流传输可能因为ALPN禁用出现异常,导致客户端无法接收完整响应,进而Future无法完成。
解决办法:
尝试开启ALPN支持:
ClientOptions options = ClientOptions.create() // 保留其他配置 .setUseAlpn(true); // 改为true
注意:如果你的JDK版本较低(比如OpenJDK 8),可能需要额外添加ALPN依赖包,例如org.mortbay.jetty.alpn:alpn-boot。
3. 超大结果集导致的等待阻塞
你设置了setExecuteQueryMaxResultRows(1000000000),如果目标表的数据量极大,客户端会一直等待全量数据接收完成。虽然控制台能正常查询,但控制台是流式输出,而批量查询的get()方法会等待所有数据传输完毕才返回。
解决办法:
先测试限制结果行数,比如把参数改为10,验证代码是否能正常返回:
.setExecuteQueryMaxResultRows(10)
如果能正常返回,说明是数据量过大的问题,可以改用流式查询(client.streamQuery())来处理大量数据,而非批量查询。
4. 配置参数传递异常
你把Properties强转为Map传递给executeQuery,虽然语法合法,但可能导致配置参数未被正确解析。另外auto.offset.reset对于pull查询是否需要传递存疑,或许应该通过ClientOptions来统一配置。
解决办法:
尝试移除Properties参数,直接调用无参的executeQuery:
BatchedQueryResult batchedQueryResult = client.executeQuery(pullQuery);
建议你优先尝试第一种方案(异步处理或@Blocking注解),这是Quarkus环境下使用异步客户端最常见的坑,大概率能解决你的问题。
备注:内容来源于stack exchange,提问作者Drakem

