apache/seatunnel · critical · ElasticsearchConnectorException

COMMON_SQL_OPERATION_FAILED

COMMON_SQL_OPERATION_FAILED

Error message

ElasticSearch execute batch statement error

What it means

bulkEsWithRetry wraps the entire retry-and-bulk execution in a try/catch; any Exception (I/O failure, exhausted retries, client errors) is rethrown as ElasticsearchConnectorException with COMMON_SQL_OPERATION_FAILED and message 'ElasticSearch execute batch statement error', preserving the original cause.

Solutions

  1. Check the 'Caused by' chain: if it ends in BULK_RESPONSE_ERROR, fix the per-item bulk issues; if network, fix connectivity
  2. Verify ES host/port/credentials in sink options and test with curl
  3. Increase retry policy (attempts/backoff) for transient cluster issues
  4. Check ES cluster health, disk watermarks, and load
Defensive patterns

Strategy: try-catch

Validate before calling

// preflight: verify ES reachability before the job
// curl -u user:pass http://host:9200/_cluster/health

Try / catch

try {
    writer.prepareCommit();
} catch (ElasticsearchConnectorException e) {
    Throwable cause = e.getCause();
    log.error("ES batch write failed; root cause: " + (cause == null ? "n/a" : cause.toString()), e);
    throw e;
}

Prevention

When it happens

Trigger: esRestClient.bulk throws or retries are exhausted during write, prepareCommit, timerFlush, or close — e.g. connection failures, timeouts, or the BULK_RESPONSE_ERROR from [997] after all retries.

Common situations: ES node unreachable/wrong host:port, authentication failures, network partitions, cluster overload, or persistent per-item bulk errors that survive retries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/2fbc7f2c5c90989a. 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:261

        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
    public void close() {
        if (esRestClient == null) {
            return;
        }
        try {
            bulkEsWithRetry(this.esRestClient, this.requestEsList);
        } finally {
            releaseSharedClientResource();
            closeOwnedClient();
        }
    }

View on GitHub (pinned to cf67b549a7)