{"record":{"id":"902667e9a36ecbb8","repo":"apache/seatunnel","slug":"upsert-data-failed-retry-in-smaller-chunks","errorCode":null,"errorMessage":"upsert data failed, retry in smaller chunks: {} ","messagePattern":"upsert data failed, retry in smaller chunks: (.+?) ","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-milvus/src/main/java/org/apache/seatunnel/connectors/seatunnel/milvus/sink/MilvusBufferBatchWriter.java","lineNumber":289,"sourceCode":"            throw new MilvusConnectorException(MilvusConnectionErrorCode.WRITE_DATA_FAIL);\n        }\n        writeCount.addAndGet(this.writeCache.get());\n    }\n\n    private void upsertWrite(String partitionName, List<JsonObject> data)\n            throws InterruptedException {\n        UpsertReq upsertReq =\n                UpsertReq.builder().collectionName(this.collectionName).data(data).build();\n        if (StringUtils.isNotEmpty(partitionName)) {\n            upsertReq.setPartitionName(partitionName);\n        }\n        try {\n            milvusClient.upsert(upsertReq);\n        } catch (Exception e) {\n            if (e.getMessage().contains(\"rate limit exceeded\")\n                    || e.getMessage().contains(\"received message larger than max\")) {\n                if (data.size() > 10) {\n                    log.warn(\"upsert data failed, retry in smaller chunks: {} \", data.size() / 2);\n                    this.batchSize = this.batchSize / 2;\n                    log.info(\"sleep 1 minute to avoid rate limit\");\n                    // sleep 1 minute to avoid rate limit\n                    Thread.sleep(60000);\n                    log.info(\"sleep 1 minute success\");\n                    // Split the data and retry in smaller chunks\n                    List<JsonObject> firstHalf = data.subList(0, data.size() / 2);\n                    List<JsonObject> secondHalf = data.subList(data.size() / 2, data.size());\n                    upsertWrite(partitionName, firstHalf);\n                    upsertWrite(partitionName, secondHalf);\n                } else {\n                    // If the data size is 10, throw the exception to avoid infinite recursion\n                    throw new MilvusConnectorException(\n                            MilvusConnectionErrorCode.WRITE_DATA_FAIL,\n                            \"upsert data failed,\" + \" size down to 10, break\",\n                            e);\n                }\n            } else {","sourceCodeStart":271,"sourceCodeEnd":307,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-milvus/src/main/java/org/apache/seatunnel/connectors/seatunnel/milvus/sink/MilvusBufferBatchWriter.java#L271-L307","documentation":"An upsert to Milvus failed with a rate-limit or message-too-large error. The writer halves the current batch size, warns with the new chunk size, sleeps 60 seconds, and recursively retries the data split in two halves. This is a self-healing backoff, but it adds latency and mutates shared batchSize state.","triggerScenarios":"upsertWrite catches an exception whose message contains 'rate limit exceeded' or 'received message larger than max' while data.size() > 10.","commonSituations":"High-throughput streams hitting Milvus rate limits; batch sizes configured too large for the Milvus gRPC max message size.","solutions":["Lower the sink batch_size option so initial batches fit within Milvus limits.","Increase Milvus server-side rate limits (quotaAndLimits) or gRPC maxMessageSize.","Reduce sink parallelism to lower request pressure.","If latency from the 60s sleep is unacceptable, apply upstream rate limiting."],"exampleFix":"// before\nsink {\n  Milvus {\n    batch_size = 5000\n  }\n}\n// after\nsink {\n  Milvus {\n    batch_size = 1000\n  }\n}","handlingStrategy":"retry","validationCode":"// size the batch under Milvus limits\nlong estBytes = rows.stream().mapToLong(r -> estimateSize(r)).sum();\nassert estBytes < 16 * 1024 * 1024;","typeGuard":null,"tryCatchPattern":"// mirror the library's chunked retry\ntry { client.upsert(req); }\ncatch (Exception e) {\n    if (isRateLimit(e)) { Thread.sleep(60_000); upsertHalves(data); }\n    else throw e;\n}","preventionTips":["Start with conservative batch_size (e.g. 500-1000) and scale up","Raise Milvus quotaAndLimits rate limits for write-heavy jobs","Avoid deep recursion by configuring sane initial batch sizes","Load-test against Milvus with production-shaped data sizes"],"tags":["milvus","rate-limit","retry","backpressure"],"backgroundTag":"rate-limit-exceeded","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}