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

使用Java SDK向Cosmos DB插入15亿条数据的提速方案咨询

Cosmos DB Java SDK 批量插入性能优化方案

问题背景

需向Cosmos DB的DOCUMENT和DOCUMENT_ATT两个容器插入约15亿条记录,当前采用循环生成数据后批量插入的方式,插入速度仅约每秒300条,Cosmos DB已配置自动缩放吞吐量(最大4000RU/s),需优化插入流程提升速度。

核心优化方案

1. 优化Cosmos客户端初始化配置

Java SDK通过CosmosClientBuilder替代.NET的CosmosClientOptions完成客户端配置,关键调整包括:

  • 使用DIRECT连接模式:跳过网关代理,降低网络开销,更适合大规模批量操作
  • 增大连接池上限:提升并发请求处理能力
  • 延长请求超时:避免大批次操作因超时中断

2. 优化批量执行参数

调整CosmosBulkExecutionOptions最大化批量效率:

  • 设置匹配吞吐量的最大并发数:让SDK充分利用Cosmos DB的RU配额
  • 启用批量分区键分组:SDK自动按分区键合并操作,减少请求次数

3. 增大单批次数据量

当前每批仅生成100条Document和50条DocumentAttribute,远低于目标的7000条/批。建议将单批次数据量调整至7000左右(需根据单条数据大小和RU消耗微调,避免单批RU超过吞吐量上限),大幅减少请求次数。

4. 异步并行执行批量操作

同步调用executeBulkOperations会阻塞线程,改用异步API(返回Mono<Void>),同时并行处理两个容器的插入任务,充分利用CPU和网络资源。

5. 数据生成逻辑优化

  • 预先初始化集合容量:避免ArrayList动态扩容带来的性能损耗
  • 复用对象或高效生成工具:比如用ThreadLocal缓存UUID生成器,减少重复对象创建开销

修改后的示例代码

import com.azure.cosmos.CosmosClient;
import com.azure.cosmos.CosmosClientBuilder;
import com.azure.cosmos.CosmosContainer;
import com.azure.cosmos.CosmosDatabase;
import com.azure.cosmos.models.CosmosBulkExecutionOptions;
import com.azure.cosmos.models.CosmosItemOperation;
import com.azure.cosmos.models.PartitionKey;
import reactor.core.publisher.Mono;

import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;

public class CosmosBulkInsertOptimized {

    public void bulkInsert(CosmosDatabase db, Date startDate, int maxFileSize, int totalNumber, int totalIntNumber, String idPrefix, List<String> list1, List<String> list2) {
        // 优化客户端配置:DIRECT模式+连接池调整
        CosmosClient cosmosClient = new CosmosClientBuilder()
                .endpoint("<你的Cosmos端点>")
                .key("<你的Cosmos密钥>")
                .directMode()
                .connectionPolicy(policy -> policy
                        .setMaxConnectionsPerEndpoint(1000)
                        .setRequestTimeoutInMillis(60000))
                .buildClient();

        CosmosContainer documentContainer = db.getContainer("DOCUMENT");
        CosmosContainer attributeContainer = db.getContainer("DOCUMENT_ATT");

        // 优化批量执行配置
        CosmosBulkExecutionOptions bulkOptions = new CosmosBulkExecutionOptions()
                .setMaxConcurrency(100)
                .setEnableBulkExecution(true);

        AtomicInteger number = new AtomicInteger(1);
        AtomicInteger extUserNumber = new AtomicInteger(1);
        AtomicInteger intNumber = new AtomicInteger(1);
        Date currentDate = startDate;

        // 按目标批次拆分总数据量
        long totalDocs = 1500000000L;
        int batchSize = 7000;
        long totalBatches = (totalDocs + batchSize - 1) / batchSize;

        for (long batch = 0; batch < totalBatches; batch++) {
            List<Document> docInsert = new ArrayList<>(batchSize);
            List<DocumentAttribute> docAttr = new ArrayList<>(batchSize / 2);

            // 生成单批次数据
            while (docInsert.size() < batchSize) {
                String docId = UUID.randomUUID().toString();
                Date expiryTime = DateUtils.addYears(currentDate, 10);

                // 添加2条Document
                docInsert.add(new Document(docId, docId, Math.floor(Math.random() * maxFileSize) + 1, number.get()));
                docInsert.add(new Document(docId, docId, Math.floor(Math.random() * maxFileSize) + 1, number.incrementAndGet()));
                number.set(number.get() + 1 > totalNumber ? 1 : number.get() + 1);

                List<User> users = new ArrayList<>(7);
                users.add(new User(idPrefix + extUserNumber.get(), "EXT", list1));
                users.add(new User(idPrefix + extUserNumber.incrementAndGet(), "EXT", list1));
                extUserNumber.set(extUserNumber.get() + 1 > totalNumber ? 1 : extUserNumber.get() + 1);

                for (int u = 0; u < 5; u++) {
                    users.add(new User(idPrefix + (intNumber.get() + u), "INT", list2));
                }
                intNumber.set(intNumber.get() + 5 > totalIntNumber ? 1 : intNumber.get() + 5);
                docAttr.add(new DocumentAttribute(docId, users, currentDate));

                // 按需更新时间戳
                if (docInsert.size() % 100 == 0) {
                    currentDate = DateUtils.addSeconds(currentDate, 1);
                }
            }

            // 异步生成两个容器的批量操作任务
            Mono<Void> docBulkTask = documentContainer.executeBulkOperations(
                    docInsert.stream()
                            .map(doc -> CosmosBulkOperations.getCreateItemOperation(doc, new PartitionKey(doc.getId())))
                            .collect(Collectors.toList()),
                    bulkOptions
            );

            Mono<Void> attrBulkTask = attributeContainer.executeBulkOperations(
                    docAttr.stream()
                            .map(attr -> CosmosBulkOperations.getCreateItemOperation(attr, new PartitionKey(attr.getDocId())))
                            .collect(Collectors.toList()),
                    bulkOptions
            );

            // 并行执行两个任务,等待批次完成
            Mono.when(docBulkTask, attrBulkTask).block();
        }

        cosmosClient.close();
    }

    // 实体类需根据实际业务定义调整
    private static class Document {
        private String id;
        private String docId;
        private double fileSize;
        private int number;

        public Document(String id, String docId, double fileSize, int number) {
            this.id = id;
            this.docId = docId;
            this.fileSize = fileSize;
            this.number = number;
        }

        public String getId() { return id; }
    }

    private static class DocumentAttribute {
        private String docId;
        private List<User> users;
        private Date date;

        public DocumentAttribute(String docId, List<User> users, Date date) {
            this.docId = docId;
            this.users = users;
            this.date = date;
        }

        public String getDocId() { return docId; }
    }

    private static class User {
        private String id;
        private String type;
        private List<String> list;

        public User(String id, String type, List<String> list) {
            this.id = id;
            this.type = type;
            this.list = list;
        }
    }
}

额外建议

  • 临时提升自动缩放最大RU:4000RU/s对于15亿条数据插入可能不足,可根据单条数据RU消耗计算所需吞吐量,临时调高上限(如10000RU/s),插入完成后再调回
  • 监控Cosmos DB指标:查看吞吐量利用率、分区分布、请求延迟等,排查是否存在热点分区或吞吐量瓶颈
  • 考虑专用批量导入工具:若Java SDK优化后仍不满足需求,可使用Azure Data Factory或Cosmos DB Bulk Executor等专用工具,针对大规模导入做深度优化

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 00:57:48