apache/seatunnel · warning

upsert data failed, retry in smaller chunks

Error message

upsert data failed, retry in smaller chunks: {} 

What it means

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.

Solutions

  1. Lower the sink batch_size option so initial batches fit within Milvus limits.
  2. Increase Milvus server-side rate limits (quotaAndLimits) or gRPC maxMessageSize.
  3. Reduce sink parallelism to lower request pressure.
  4. If latency from the 60s sleep is unacceptable, apply upstream rate limiting.

Example fix

// before
sink {
  Milvus {
    batch_size = 5000
  }
}
// after
sink {
  Milvus {
    batch_size = 1000
  }
}
Defensive patterns

Strategy: retry

Validate before calling

// size the batch under Milvus limits
long estBytes = rows.stream().mapToLong(r -> estimateSize(r)).sum();
assert estBytes < 16 * 1024 * 1024;

Try / catch

// mirror the library's chunked retry
try { client.upsert(req); }
catch (Exception e) {
    if (isRateLimit(e)) { Thread.sleep(60_000); upsertHalves(data); }
    else throw e;
}

Prevention

When it happens

Trigger: upsertWrite catches an exception whose message contains 'rate limit exceeded' or 'received message larger than max' while data.size() > 10.

Common situations: High-throughput streams hitting Milvus rate limits; batch sizes configured too large for the Milvus gRPC max message size.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/902667e9a36ecbb8. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-milvus/src/main/java/org/apache/seatunnel/connectors/seatunnel/milvus/sink/MilvusBufferBatchWriter.java:289

            throw new MilvusConnectorException(MilvusConnectionErrorCode.WRITE_DATA_FAIL);
        }
        writeCount.addAndGet(this.writeCache.get());
    }

    private void upsertWrite(String partitionName, List<JsonObject> data)
            throws InterruptedException {
        UpsertReq upsertReq =
                UpsertReq.builder().collectionName(this.collectionName).data(data).build();
        if (StringUtils.isNotEmpty(partitionName)) {
            upsertReq.setPartitionName(partitionName);
        }
        try {
            milvusClient.upsert(upsertReq);
        } catch (Exception e) {
            if (e.getMessage().contains("rate limit exceeded")
                    || e.getMessage().contains("received message larger than max")) {
                if (data.size() > 10) {
                    log.warn("upsert data failed, retry in smaller chunks: {} ", data.size() / 2);
                    this.batchSize = this.batchSize / 2;
                    log.info("sleep 1 minute to avoid rate limit");
                    // sleep 1 minute to avoid rate limit
                    Thread.sleep(60000);
                    log.info("sleep 1 minute success");
                    // Split the data and retry in smaller chunks
                    List<JsonObject> firstHalf = data.subList(0, data.size() / 2);
                    List<JsonObject> secondHalf = data.subList(data.size() / 2, data.size());
                    upsertWrite(partitionName, firstHalf);
                    upsertWrite(partitionName, secondHalf);
                } else {
                    // If the data size is 10, throw the exception to avoid infinite recursion
                    throw new MilvusConnectorException(
                            MilvusConnectionErrorCode.WRITE_DATA_FAIL,
                            "upsert data failed," + " size down to 10, break",
                            e);
                }
            } else {

View on GitHub (pinned to cf67b549a7)