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

