apache/seatunnel · error · ElasticsearchConnectorException

BULK_RESPONSE_ERROR

BULK_RESPONSE_ERROR

Error message

bulk es error: 

What it means

After sending a bulk request, ElasticsearchSinkWriter checks BulkResponse.isErrors(). If Elasticsearch reports per-item failures, it throws ElasticsearchConnectorException with BULK_RESPONSE_ERROR and the full response body. The batch is retried per the retry policy before failing.

Solutions

  1. Inspect the embedded response JSON for per-item error reasons (type/reason per item)
  2. Fix mapping conflicts: delete/reindex the conflicting index or align the schema
  3. Reduce batch size if the cluster rejected the bulk (429/oversized)
  4. Ensure the target index exists and is open before writing
  5. Check ES cluster health and disk watermarks (read-only indices block writes)
Defensive patterns

Strategy: retry

Try / catch

try {
    writer.prepareCommit();
} catch (ElasticsearchConnectorException e) {
    if (e.getErrorCode() == ElasticsearchConnectorErrorCode.BULK_RESPONSE_ERROR) {
        String body = e.getMessage(); // parse per-item errors and fix mappings/batch size
        log.error("Bulk item failures: " + body);
    }
    throw e;
}

Prevention

When it happens

Trigger: esRestClient.bulk(requestBody) returns a response with errors:true — e.g. mapping conflicts, index_not_found, version conflicts, rejected by cluster (429/oversized bulk), or invalid documents.

Common situations: Writing documents whose fields conflict with existing index mappings; bulk payloads exceeding http.max_content_length; index missing or closed; ES cluster under load rejecting bulk shards.

Understand the failure class

Background: "API error: {status}" and "HTTP 401/403/404/429/5xx" errors: non-2xx HTTP responses explained — this error's family across 27 libraries.

Related errors


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

Appendix: source

Thrown at seatunnel-connectors-v2/connector-elasticsearch/src/main/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchSinkWriter.java:250

     * <p>The action is registered before multi-table resource injection but invoked after startup.
     */
    private void timerFlush() {
        bulkEsWithRetry(this.esRestClient, this.requestEsList);
    }

    @Override
    public void abortPrepare() {}

    public synchronized void bulkEsWithRetry(
            EsRestClient esRestClient, List<String> requestEsList) {
        try {
            RetryUtils.retryWithException(
                    () -> {
                        if (!requestEsList.isEmpty()) {
                            String requestBody = String.join("\n", requestEsList) + "\n";
                            BulkResponse bulkResponse = esRestClient.bulk(requestBody);
                            if (bulkResponse.isErrors()) {
                                throw new ElasticsearchConnectorException(
                                        ElasticsearchConnectorErrorCode.BULK_RESPONSE_ERROR,
                                        "bulk es error: " + bulkResponse.getResponse());
                            }
                            return bulkResponse;
                        }
                        return null;
                    },
                    retryMaterial);
            requestEsList.clear();
        } catch (Exception e) {
            throw new ElasticsearchConnectorException(
                    CommonErrorCodeDeprecated.SQL_OPERATION_FAILED,
                    "ElasticSearch execute batch statement error",
                    e);
        }
    }

    @Override

View on GitHub (pinned to cf67b549a7)