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

  1. Inspect the wrapped cause for the concrete gRPC status (NOT_FOUND, UNAVAILABLE, DEADLINE_EXCEEDED, etc.)
  2. Verify the collection name exists in Qdrant (create it or enable auto-creation in the connector config)
  3. Check network connectivity and Qdrant URL/credentials from all SeaTunnel worker nodes
  4. Reduce batch_size to stay under gRPC max message size
  5. 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

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


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)