You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Quarkus应用中Java版kSqlDB客户端问题:BatchedQueryResult的completableFuture get()方法挂起

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-stream REST接口也能正常拿到结果

你的代码参考

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.21 12:58:09