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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 01:24:50