apache/seatunnel · error · RuntimeException
Upsert failed
Error message
Upsert failed
What it means
QdrantBatchWriter.upsert submits a batched Points.UpsertPoints gRPC call asynchronously and blocks on the future's get(). If the future is interrupted or completes exceptionally, the writer wraps the cause in a RuntimeException('Upsert failed'). The actual reason (timeout, collection missing, connection failure, unauthenticated) is always in the cause chain.
Solutions
- Inspect the wrapped cause for the concrete gRPC status (NOT_FOUND, UNAVAILABLE, DEADLINE_EXCEEDED, etc.)
- Verify the collection name exists in Qdrant (create it or enable auto-creation in the connector config)
- Check network connectivity and Qdrant URL/credentials from all SeaTunnel worker nodes
- Reduce batch_size to stay under gRPC max message size
- Retry the job if the cause was a transient interruption or brief unavailability
Example fix
// before
} catch (InterruptedException | ExecutionException e) {
throw new RuntimeException("Upsert failed", e);
}
// after
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("Upsert failed (interrupted)", e);
} catch (ExecutionException e) {
LOG.error("Qdrant upsert failed; check collection existence, connectivity and batch size", e);
throw new RuntimeException("Upsert failed", e);
} Defensive patterns
Strategy: retry
Validate before calling
// pre-flight: collection must exist and be reachable
boolean ok = client.collectionExistsAsync(collectionName).get(5, TimeUnit.SECONDS);
if (!ok) throw new IllegalStateException("Qdrant collection missing: " + collectionName); Try / catch
try {
writer.flush(); // triggers upsert
} catch (RuntimeException e) {
if ("Upsert failed".equals(e.getMessage())) {
Throwable cause = e.getCause();
if (cause instanceof InterruptedException) {
Thread.currentThread().interrupt();
}
// inspect gRPC status in cause; retry on transient UNAVAILABLE/DEADLINE_EXCEEDED
} else throw e;
} Prevention
- Create the Qdrant collection before running the job or enable auto-creation
- Validate host/port/API key connectivity from all worker nodes pre-job
- Keep batch sizes under the gRPC max message size (default 4MB)
- Use modest write concurrency to avoid broker-side overload and timeouts
When it happens
Trigger: flush() calls upsert(); the gRPC upsert future throws ExecutionException (server error, connection loss, deadline exceeded) or the waiting thread is interrupted.
Common situations: Qdrant collection does not exist or was deleted mid-job; Qdrant host/port wrong or unreachable from worker nodes; batch too large exceeding gRPC message limits; auth/API key mismatch; network timeouts under load.
Related errors
- COMMON-07
- COMMON-07
- Failed to open catalog
- Failed to write vertices to NebulaGraph tag ' '. The writer…
- FLUSH_DATA_FAILED
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/8b8c9c58d8a121a5.
Report an issue: GitHub.
Appendix: source
Thrown at seatunnel-connectors-v2/connector-qdrant/src/main/java/org/apache/seatunnel/connectors/seatunnel/qdrant/sink/QdrantBatchWriter.java:138
if (!point.hasId()) {
point.setId(id(UUID.randomUUID()));
}
point.setVectors(Points.Vectors.newBuilder().setVectors(namedVectors).build());
return point.build();
}
private void upsert() {
try {
qdrantClient
.upsertAsync(
Points.UpsertPoints.newBuilder()
.setCollectionName(collectionName)
.addAllPoints(qdrantDataCache)
.build())
.get();
} catch (InterruptedException | ExecutionException e) {
throw new RuntimeException("Upsert failed", e);
}
}
public static Points.PointId pointId(SeaTunnelDataType<?> fieldType, Object value) {
SqlType sqlType = fieldType.getSqlType();
switch (sqlType) {
case INT:
return id(Integer.parseInt(value.toString()));
case STRING:
return id(UUID.fromString(value.toString()));
default:
throw new QdrantConnectorException(
CommonErrorCode.UNSUPPORTED_DATA_TYPE,
"Unexpected value type for point ID: " + sqlType.name());
}
}
public static JsonWithInt.Value buildPayload(SeaTunnelDataType<?> fieldType, Object value) {View on GitHub (pinned to cf67b549a7)