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
相关产品推荐
相关产品推荐

