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

ServiceBusSessionReceiverAsyncClient关闭时抛出IllegalStateException问题

问题分析与解决方案

这个异常的核心原因是接收链路关闭的时序冲突:当take(1)触发流结束后,同步执行receiver.close()会立即关闭AMQP链路,但此时Receiver内部的背压逻辑可能还在尝试向链路添加额度(credits),导致操作失败抛出异常。

具体修复方案

1. 使用异步关闭方法替代同步关闭

将Mono.fromRunnable(receiver::close)替换为receiver.closeAsync(),因为closeAsync()是Azure SDK提供的异步关闭API,能保证关闭操作与Reactor流的生命周期正确对齐,避免同步关闭导致的时序竞争。

2. 简化流操作逻辑

take(1).next()可以直接简化为next(),因为next()本身就会订阅流并获取第一个元素,之后自动取消订阅,效果与take(1)一致但更简洁。

3. 复用ServiceBus客户端实例(可选但推荐)

每次调用方法都创建新的ServiceBusSessionReceiverAsyncClient是低效的,ServiceBus客户端是线程安全的,应该全局复用一个实例,避免频繁创建/销毁连接带来的资源开销和潜在问题。

修改后的代码

private <T extends IWocTransaction> Mono<Optional<T>> responseAsync(String transactionId, Class<T> clazz) {
    // 注意:这里的asyncClient应该是全局复用的实例,而非每次创建
    var msgStream = Flux.usingWhen(asyncClient.acceptSession(transactionId),
            receiver -> receiver.receiveMessages(),
            receiver -> receiver.closeAsync() // 改用异步关闭API
    );

    return msgStream
            .timeout(timeout)
            .next() // 直接用next()获取第一个元素
            .map(message -> {
                var json = message.getBody().toString();
                try {
                    var val = objectMapper.readValue(json, clazz);
                    return Optional.ofNullable(val);
                } catch (Exception e) {
                    log.error("Error deserializing response from string {}", json, e);
                    return Optional.empty();
                }
            })
            .doOnError(t -> {
                if (t instanceof TimeoutException) {
                    log.error("Timeout error waiting on API callback {}", kv("ApiTimeout", timeout.toString()), t);
                } else {
                    log.error("Error waiting for async callback", t);
                }
            })
            .onErrorReturn(Optional.empty());
}

额外说明

如果必须每次创建客户端,确保asyncClient也被正确关闭,可以在usingWhen外层再包裹一层来管理asyncClient的生命周期:

return Flux.usingWhen(
        Mono.fromCallable(() -> sbClientBuilder.connectionString(sbConnectionString)
                .sessionReceiver()
                .queueName("my-callback-queue")
                .receiveMode(ServiceBusReceiveMode.RECEIVE_AND_DELETE)
                .buildAsyncClient()),
        client -> Flux.usingWhen(client.acceptSession(transactionId),
                receiver -> receiver.receiveMessages(),
                receiver -> receiver.closeAsync()),
        client -> client.closeAsync()
)
// 后续的map、doOnError等逻辑...

内容的提问来源于stack exchange,提问作者LPal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 11:35:14