Skip to content

Commit 4c357d9

Browse files
committed
fix: remove BulkIngester backoffPolicy to prevent compound retries
BulkIngester.backoffPolicy and DispatchingTransport both used config.maxRetries(), causing up to (maxRetries+1)^2 HTTP requests per document: each BulkIngester per-item 429 retry issued a new bulk request that re-entered DispatchingTransport with a fresh retry budget. Removing backoffPolicy gives DispatchingTransport sole ownership of the retry budget. HTTP-level 429s (whole-request) continue to be retried as exceptions by DispatchingTransport; per-item 429s in a 200 response go to the listener's error path (DLQ + task failure), matching pre-migration HLRC behavior.
1 parent 80672d1 commit 4c357d9

1 file changed

Lines changed: 0 additions & 3 deletions

File tree

src/main/java/io/confluent/connect/elasticsearch/ElasticsearchClient.java

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,6 @@
2525
import co.elastic.clients.elasticsearch.indices.get_mapping.IndexMappingRecord;
2626
import co.elastic.clients.json.JsonpMapper;
2727
import co.elastic.clients.json.jackson.JacksonJsonpMapper;
28-
import co.elastic.clients.transport.BackoffPolicy;
2928
import co.elastic.clients.transport.ElasticsearchTransport;
3029
import co.elastic.clients.transport.Endpoint;
3130
import co.elastic.clients.transport.TransportOptions;
@@ -195,8 +194,6 @@ public ElasticsearchClient(
195194
.maxConcurrentRequests(config.maxInFlightRequests())
196195
.flushInterval(config.lingerMs(), TimeUnit.MILLISECONDS)
197196
.scheduler(this.bulkScheduler)
198-
.backoffPolicy(BackoffPolicy.exponentialBackoff(
199-
config.retryBackoffMs(), config.maxRetries()))
200197
.listener(buildListener(afterBulkCallback));
201198
if (config.bulkSize() > 0) {
202199
b.maxSize(config.bulkSize());

0 commit comments

Comments
 (0)