Azure Cosmos DB Java SDK批量写入请求超限及错误处理优化咨询
通过Spring Scheduler以300条为批次,异步向Azure Cosmos DB(NoSQL API)批量Upsert 10万条数据。单条数据平均消耗60-70 Request Units(RU),容器配置了10000 RU的自动缩放功能。
当前代码
AtomicInteger retry = new AtomicInteger(retryCount); Flux<Data> gtinContainerFlux = Flux.fromIterable(dataset); Flux<CosmosItemOperation> cosmosItemOperationFlux = gtinContainerFlux.map(data -> CosmosBulkOperations.getUpsertItemOperation(data, new PartitionKey(data.getKey()))); cosmosAsyncContainer.executeBulkOperations(cosmosItemOperationFlux, new CosmosBulkExecutionOptions()) .subscribe(cosmosBulkOperationResponse -> { CosmosBulkItemResponse cosmosBulkItemResponse = cosmosBulkOperationResponse.getResponse(); CosmosItemOperation cosmosItemOperation = cosmosBulkOperationResponse.getOperation(); if (cosmosBulkOperationResponse.getException() != null) { log.warn("ERROR : {}", cosmosItemOperation.<Data>getItem().getId()); log.error("Bulk operation failed", cosmosBulkOperationResponse.getException()); } else if (cosmosBulkItemResponse == null || !cosmosBulkOperationResponse.getResponse().isSuccessStatusCode()) { log.error("The operation for Item ID: [{}] Item PartitionKey Value: [{}] did not complete successfully with a {} response code.", cosmosItemOperation.<Data>getItem().getId(), cosmosItemOperation.<Data>getItem().getGtinKey(), cosmosBulkItemResponse != null ? cosmosBulkItemResponse.getStatusCode() : "n/a"); errorEvents.add(cosmosItemOperation.<Data>getItem()); } else { log.debug("Data posted successfully with RUs used : " + cosmosBulkItemResponse.getRequestCharge() + " completed in " + cosmosBulkItemResponse.getCosmosDiagnostics().getDuration().getNano()); } });
执行错误
Request rate is large. More Request Units may be needed, so no changes were made. Please retry this request later
堆栈信息
at com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdRequestManager.messageReceived(RntbdRequestManager.java:1121) 2024-01-15T11:09:54.963941386Z at com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdRequestManager.channelRead(RntbdRequestManager.java:214) 2024-01-15T11:09:54.963946986Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:444) 2024-01-15T11:09:54.963972387Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) 2024-01-15T11:09:54.963978787Z at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) 2024-01-15T11:09:54.963992287Z at io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:346) 2024-01-15T11:09:54.963997587Z at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:318) 2024-01-15T11:09:54.964002887Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:444) 2024-01-15T11:09:54.964007787Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) 2024-01-15T11:09:54.964012587Z at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) 2024-01-15T11:09:54.964018187Z at io.netty.channel.CombinedChannelDuplexHandler$DelegatingChannelHandlerContext.fireChannelRead(CombinedChannelDuplexHandler.java:436) 2024-01-15T11:09:54.964022987Z at io.netty.channel.CombinedChannelDuplexHandler.channelRead(CombinedChannelDuplexHandler.java:253) 2024-01-15T11:09:54.964028088Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:442) 2024-01-15T11:09:54.964033288Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) 2024-01-15T11:09:54.964038288Z at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) 2024-01-15T11:09:54.964043688Z at io.netty.handler.timeout.IdleStateHandler.channelRead(IdleStateHandler.java:286) 2024-01-15T11:09:54.964048788Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:442) 2024-01-15T11:09:54.964053688Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) 2024-01-15T11:09:54.964058488Z at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) 2024-01-15T11:09:54.964063888Z at io.netty.handler.ssl.SslHandler.unwrap(SslHandler.java:1475) 2024-01-15T11:09:54.964081088Z at io.netty.handler.ssl.SslHandler.decodeJdkCompatible(SslHandler.java:1338) 2024-01-15T11:09:54.964086788Z at io.netty.handler.ssl.SslHandler.decode(SslHandler.java:1387) 2024-01-15T11:09:54.964091689Z at io.netty.handler.codec.ByteToMessageDecoder.decodeRemovalReentryProtection(ByteToMessageDecoder.java:529) 2024-01-15T11:09:54.964096589Z at io.netty.handler.codec.ByteToMessageDecoder.callDecode(ByteToMessageDecoder.java:468) 2024-01-15T11:09:54.964101589Z at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:290) 2024-01-15T11:09:54.964106389Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:444) 2024-01-15T11:09:54.964111589Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) 2024-01-15T11:09:54.964116489Z at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) 2024-01-15T11:09:54.964121489Z at io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410) 2024-01-15T11:09:54.964126289Z at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:440) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) 2024-01-15T11:09:54.964135789Z at io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919) 2024-01-15T11:09:54.964140789Z at io.netty.channel.epoll.AbstractEpollStreamChannel$EpollStreamUnsafe.epollInReady(AbstractEpollStreamChannel.java:800) 2024-01-15T11:09:54.964145589Z at io.netty.channel.epoll.EpollEventLoop.processReady(EpollEventLoop.java:509) 2024-01-15T11:09:54.964155690Z at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:407) 2024-01-15T11:09:54.964160990Z at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:997) 2024-01-15T11:09:54.964166190Z at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74) 2024-01-15T11:09:54.964171090Z at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30) 2024-01-15T11:09:54.964176090Z at java.base/java.lang.Thread.run(Thread.java:833)
核心问题
当前代码中cosmosBulkOperationResponse.getException() != null分支的日志无法打印,无法正确处理该错误,同时需要解决请求速率过大导致的限流问题。
一、错误处理优化
1. 完善Reactive订阅的错误处理
仅使用subscribe(Consumer)会忽略Flux的全局错误,需添加onError分支捕获批量操作的全局异常,同时调整单条操作的错误判断逻辑:
cosmosAsyncContainer.executeBulkOperations(cosmosItemOperationFlux, new CosmosBulkExecutionOptions()) .doOnNext(cosmosBulkOperationResponse -> { CosmosBulkItemResponse cosmosBulkItemResponse = cosmosBulkOperationResponse.getResponse(); CosmosItemOperation cosmosItemOperation = cosmosBulkOperationResponse.getOperation(); // 优先判断限流错误 if (cosmosBulkItemResponse != null && cosmosBulkItemResponse.getStatusCode() == 429) { log.warn("Item [{}] hit rate limit", cosmosItemOperation.<Data>getItem().getId()); errorEvents.add(cosmosItemOperation.<Data>getItem()); return; } if (cosmosBulkOperationResponse.getException() != null) { log.warn("ERROR : {}", cosmosItemOperation.<Data>getItem().getId()); log.error("Bulk operation failed", cosmosBulkOperationResponse.getException()); errorEvents.add(cosmosItemOperation.<Data>getItem()); } else if (cosmosBulkItemResponse == null || !cosmosBulkItemResponse.isSuccessStatusCode()) { log.error("Item ID: [{}], PartitionKey: [{}] failed with code: {}", cosmosItemOperation.<Data>getItem().getId(), cosmosItemOperation.<Data>getItem().getGtinKey(), cosmosBulkItemResponse != null ? cosmosBulkItemResponse.getStatusCode() : "n/a"); errorEvents.add(cosmosItemOperation.<Data>getItem()); } else { log.debug("Item [{}] upserted successfully, RU used: {}, duration: {}ns", cosmosItemOperation.<Data>getItem().getId(), cosmosBulkItemResponse.getRequestCharge(), cosmosBulkItemResponse.getCosmosDiagnostics().getDuration().getNano()); } }) .doOnError(globalException -> { log.error("Global bulk operation failed", globalException); // 处理全局错误,比如标记整个批次失败 }) .subscribe();
2. 区分限流异常和其他异常
对于429限流错误,SDK会自动重试,单独标记日志可方便后续排查;同时确保errorEvents正确收集所有失败条目,用于后续重试。
二、代码调整:控制批次速率
1. 调整批次大小
当前300条/批次需18000-21000 RU,远超容器10000 RU上限,建议将批次大小调整为140-160条(140*70=9800 RU,接近上限):
// 将数据集按150条拆分 List<List<Data>> batches = Lists.partition(dataset, 150); for (List<Data> batch : batches) { Flux<Data> batchFlux = Flux.fromIterable(batch); // 后续批量操作逻辑不变 }
2. 添加批次间隔
在批次之间添加短暂延迟,给自动缩放足够时间调整RU,避免持续压满:
Flux.fromIterable(batches) .delayElements(Duration.ofMillis(500)) // 每个批次间隔500ms .flatMap(batch -> { Flux<Data> batchFlux = Flux.fromIterable(batch); Flux<CosmosItemOperation> operationFlux = batchFlux.map(data -> CosmosBulkOperations.getUpsertItemOperation(data, new PartitionKey(data.getKey()))); return cosmosAsyncContainer.executeBulkOperations(operationFlux, new CosmosBulkExecutionOptions()) .doOnNext(...) // 错误处理逻辑 .doOnError(...); }) .subscribe();
3. 配置批量执行选项
启用SDK自动重试,并设置合理的重试策略:
CosmosBulkExecutionOptions bulkOptions = new CosmosBulkExecutionOptions() .setMaxRetryAttemptsOnThrottledRequests(5) // 限流时最多重试5次 .setMaxRetryWaitTime(Duration.ofSeconds(30)); // 重试总等待时间30秒 cosmosAsyncContainer.executeBulkOperations(cosmosItemOperationFlux, bulkOptions) // 后续订阅逻辑
三、配置设置建议
1. 调整容器自动缩放参数
将自动缩放的最小RU设置为5000 RU以上,避免自动缩放反应不及时;若持续限流,可考虑将最大RU提升至20000 RU。
2. 检查分区键分布
确保数据分区键分布均匀,避免单一分区被持续压满导致限流;若分区键倾斜严重,需重新设计分区策略。
3. 启用SDK诊断日志
在application.yml中添加Cosmos SDK日志配置,方便排查限流细节:
logging: level: com.azure.cosmos: DEBUG
内容的提问来源于stack exchange,提问作者Prashant Aghara

