This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 7e24d0f42b [Improve][Connector-V2][Airtable] Add jitter to rate-limit
backoff (#11654)
7e24d0f42b is described below
commit 7e24d0f42b8ec06c81e478ef74ac9d61796d9a77
Author: Santhosh Kumar Somarapu <[email protected]>
AuthorDate: Tue Aug 11 10:07:37 2026 -0500
[Improve][Connector-V2][Airtable] Add jitter to rate-limit backoff (#11654)
Signed-off-by: Santhosh Kumar Somarapu <[email protected]>
---
.../airtable/sink/AirtableSinkWriter.java | 44 +++++-
.../airtable/source/AirtableSourceReader.java | 46 +++++-
.../airtable/sink/AirtableSinkWriterTest.java | 146 +++++++++++++++++++
.../airtable/source/AirtableSourceReaderTest.java | 161 +++++++++++++++++++++
4 files changed, 391 insertions(+), 6 deletions(-)
diff --git
a/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/airtable/sink/AirtableSinkWriter.java
b/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/airtable/sink/AirtableSinkWriter.java
index 97cd959893..a338e2f80c 100644
---
a/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/airtable/sink/AirtableSinkWriter.java
+++
b/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/airtable/sink/AirtableSinkWriter.java
@@ -21,6 +21,7 @@ import
org.apache.seatunnel.shade.com.fasterxml.jackson.databind.JsonNode;
import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.ObjectMapper;
import
org.apache.seatunnel.shade.com.fasterxml.jackson.databind.node.ArrayNode;
import
org.apache.seatunnel.shade.com.fasterxml.jackson.databind.node.ObjectNode;
+import
org.apache.seatunnel.shade.com.google.common.annotations.VisibleForTesting;
import org.apache.seatunnel.api.sink.SupportMultiTableSinkWriter;
import org.apache.seatunnel.api.table.type.SeaTunnelRow;
@@ -39,6 +40,7 @@ import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
+import java.util.concurrent.ThreadLocalRandom;
@Slf4j
public class AirtableSinkWriter extends AbstractSinkWriter<SeaTunnelRow, Void>
@@ -175,13 +177,49 @@ public class AirtableSinkWriter extends
AbstractSinkWriter<SeaTunnelRow, Void>
lastRequestTimeMillis = System.currentTimeMillis();
}
- private long calculateBackoffMillis(int retryCount) {
+ @VisibleForTesting
+ long calculateBackoffMillis(int retryCount) {
if (rateLimitBackoffMs <= 0) {
return 0L;
}
long exponential = 1L << Math.min(20, Math.max(0, retryCount - 1));
- long waitMillis = rateLimitBackoffMs * exponential;
- return Math.min(waitMillis, MAX_BACKOFF_MILLIS);
+ long waitMillis = Math.min(rateLimitBackoffMs * exponential,
MAX_BACKOFF_MILLIS);
+
+ // Spread the delay by adding a random amount on top of it. Without
this
+ // the delay is a pure function of the retry count, so every reader and
+ // writer that hits the rate limit at the same moment retries at the
same
+ // instants and the burst that caused the 429 reforms on each attempt.
+ //
+ // The jitter is added rather than centred so the wait is never shorter
+ // than rateLimitBackoffMs asked for: this fires on 429, so retrying
+ // sooner than configured would work against the setting's purpose. The
+ // result stays capped at MAX_BACKOFF_MILLIS.
+ long extra = Math.min(waitMillis, MAX_BACKOFF_MILLIS - waitMillis);
+ if (extra > 0) {
+ return waitMillis + ThreadLocalRandom.current().nextLong(extra +
1);
+ }
+
+ // Once the wait reaches MAX_BACKOFF_MILLIS there is no headroom left
to
+ // add into, so every retry past that point would come back unjittered
+ // and the callers would be back in lockstep exactly when the rate
limit
+ // is at its most persistent. Spread the wait downwards instead. The
cap
+ // is an upper bound rather than a target, so drawing below it breaks
+ // nothing.
+ //
+ // The floor is the last scheduled wait that still fitted under the
cap,
+ // or half the wait when the very first retry is already capped.
Flooring
+ // there keeps the minimum from dropping as the schedule crosses the
cap:
+ // half of MAX can be less than the previous retry's wait, which would
let
+ // a later retry sleep for less than an earlier one. It also keeps the
+ // wait at or above rateLimitBackoffMs for free, since the last
uncapped
+ // wait is never smaller than the configured backoff.
+ long floor = waitMillis / 2;
+ for (long scheduled = rateLimitBackoffMs; scheduled <
MAX_BACKOFF_MILLIS; scheduled <<= 1) {
+ if (scheduled > floor) {
+ floor = scheduled;
+ }
+ }
+ return waitMillis - ThreadLocalRandom.current().nextLong(waitMillis -
floor + 1);
}
@Override
diff --git
a/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/airtable/source/AirtableSourceReader.java
b/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/airtable/source/AirtableSourceReader.java
index 8f22b7a50b..cc2ba40162 100644
---
a/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/airtable/source/AirtableSourceReader.java
+++
b/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/main/java/org/apache/seatunnel/connectors/seatunnel/airtable/source/AirtableSourceReader.java
@@ -17,6 +17,8 @@
package org.apache.seatunnel.connectors.seatunnel.airtable.source;
+import
org.apache.seatunnel.shade.com.google.common.annotations.VisibleForTesting;
+
import org.apache.seatunnel.api.serialization.DeserializationSchema;
import org.apache.seatunnel.api.table.type.SeaTunnelRow;
import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
@@ -28,6 +30,8 @@ import
org.apache.seatunnel.connectors.seatunnel.http.source.HttpSourceReader;
import lombok.extern.slf4j.Slf4j;
+import java.util.concurrent.ThreadLocalRandom;
+
@Slf4j
public class AirtableSourceReader extends HttpSourceReader {
@@ -109,12 +113,48 @@ public class AirtableSourceReader extends
HttpSourceReader {
lastRequestTimeMillis = System.currentTimeMillis();
}
- private long calculateBackoffMillis(int retryCount) {
+ @VisibleForTesting
+ long calculateBackoffMillis(int retryCount) {
if (rateLimitBackoffMs <= 0) {
return 0L;
}
long exponential = 1L << Math.min(20, Math.max(0, retryCount - 1));
- long waitMillis = rateLimitBackoffMs * exponential;
- return Math.min(waitMillis, MAX_BACKOFF_MILLIS);
+ long waitMillis = Math.min(rateLimitBackoffMs * exponential,
MAX_BACKOFF_MILLIS);
+
+ // Spread the delay by adding a random amount on top of it. Without
this
+ // the delay is a pure function of the retry count, so every reader and
+ // writer that hits the rate limit at the same moment retries at the
same
+ // instants and the burst that caused the 429 reforms on each attempt.
+ //
+ // The jitter is added rather than centred so the wait is never shorter
+ // than rateLimitBackoffMs asked for: this fires on 429, so retrying
+ // sooner than configured would work against the setting's purpose. The
+ // result stays capped at MAX_BACKOFF_MILLIS.
+ long extra = Math.min(waitMillis, MAX_BACKOFF_MILLIS - waitMillis);
+ if (extra > 0) {
+ return waitMillis + ThreadLocalRandom.current().nextLong(extra +
1);
+ }
+
+ // Once the wait reaches MAX_BACKOFF_MILLIS there is no headroom left
to
+ // add into, so every retry past that point would come back unjittered
+ // and the callers would be back in lockstep exactly when the rate
limit
+ // is at its most persistent. Spread the wait downwards instead. The
cap
+ // is an upper bound rather than a target, so drawing below it breaks
+ // nothing.
+ //
+ // The floor is the last scheduled wait that still fitted under the
cap,
+ // or half the wait when the very first retry is already capped.
Flooring
+ // there keeps the minimum from dropping as the schedule crosses the
cap:
+ // half of MAX can be less than the previous retry's wait, which would
let
+ // a later retry sleep for less than an earlier one. It also keeps the
+ // wait at or above rateLimitBackoffMs for free, since the last
uncapped
+ // wait is never smaller than the configured backoff.
+ long floor = waitMillis / 2;
+ for (long scheduled = rateLimitBackoffMs; scheduled <
MAX_BACKOFF_MILLIS; scheduled <<= 1) {
+ if (scheduled > floor) {
+ floor = scheduled;
+ }
+ }
+ return waitMillis - ThreadLocalRandom.current().nextLong(waitMillis -
floor + 1);
}
}
diff --git
a/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/airtable/sink/AirtableSinkWriterTest.java
b/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/airtable/sink/AirtableSinkWriterTest.java
index 236c0821f8..3403fd82d8 100644
---
a/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/airtable/sink/AirtableSinkWriterTest.java
+++
b/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/airtable/sink/AirtableSinkWriterTest.java
@@ -38,7 +38,9 @@ import org.mockito.MockitoAnnotations;
import java.io.IOException;
import java.lang.reflect.Field;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.Map;
+import java.util.Set;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
@@ -114,4 +116,148 @@ public class AirtableSinkWriterTest {
// 1 initial + 3 retries = 4 calls
verify(httpClient, times(4)).doPost(anyString(), any(), anyString());
}
+
+ private AirtableSinkWriter writerWithBackoff(int rateLimitBackoffMs, int
rateLimitMaxRetries) {
+ HttpParameter param = new HttpParameter();
+ param.setUrl("https://api.airtable.com/v0/appXXX/tblYYY");
+ return new AirtableSinkWriter(
+ rowType, param, 1, false, 0, rateLimitBackoffMs,
rateLimitMaxRetries);
+ }
+
+ // calculateBackoffMillis is duplicated verbatim in AirtableSourceReader.
These
+ // mirror the reader's cases so the two copies cannot drift apart
unnoticed.
+
+ @Test
+ public void testBackoffIsZeroWhenDisabled() {
+ AirtableSinkWriter writer = writerWithBackoff(0, 3);
+
+ Assertions.assertEquals(0L, writer.calculateBackoffMillis(1));
+ Assertions.assertEquals(0L, writer.calculateBackoffMillis(5));
+ }
+
+ @Test
+ public void testBackoffNeverShorterThanConfigured() {
+ int base = 1000;
+ AirtableSinkWriter writer = writerWithBackoff(base, 5);
+
+ for (int retry = 1; retry <= 5; retry++) {
+ long scheduled = (long) base * (1L << (retry - 1));
+ long ceiling = Math.min(2 * scheduled, 300000L);
+ for (int i = 0; i < 200; i++) {
+ long actual = writer.calculateBackoffMillis(retry);
+ Assertions.assertTrue(
+ actual >= scheduled && actual <= ceiling,
+ "retry "
+ + retry
+ + " produced "
+ + actual
+ + ", expected ["
+ + scheduled
+ + ", "
+ + ceiling
+ + "]");
+ }
+ }
+ }
+
+ @Test
+ public void testBackoffIsNotDeterministic() {
+ AirtableSinkWriter writer = writerWithBackoff(1000, 3);
+
+ Set<Long> observed = new HashSet<>();
+ for (int i = 0; i < 200; i++) {
+ observed.add(writer.calculateBackoffMillis(3));
+ }
+
+ Assertions.assertTrue(
+ observed.size() > 1,
+ "backoff must vary between calls so that concurrent writers do
not "
+ + "retry in lockstep, but observed only "
+ + observed);
+ }
+
+ @Test
+ public void testBackoffRespectsMaximum() {
+ AirtableSinkWriter writer = writerWithBackoff(60000, 30);
+
+ for (int i = 0; i < 100; i++) {
+ Assertions.assertTrue(
+ writer.calculateBackoffMillis(20) <= 300000L,
+ "jittered backoff must never exceed MAX_BACKOFF_MILLIS");
+ }
+ }
+
+ @Test
+ public void testBackoffStillVariesAtTheMaximum() {
+ AirtableSinkWriter writer = writerWithBackoff(60000, 30);
+
+ Set<Long> observed = new HashSet<>();
+ for (int i = 0; i < 200; i++) {
+ observed.add(writer.calculateBackoffMillis(20));
+ }
+
+ Assertions.assertTrue(
+ observed.size() > 1,
+ "backoff must still vary once it reaches MAX_BACKOFF_MILLIS,
but observed only "
+ + observed);
+ }
+
+ @Test
+ public void testBackoffAtTheMaximumNeverDropsBelowConfiguredBase() {
+ int base = 200000;
+ AirtableSinkWriter writer = writerWithBackoff(base, 30);
+
+ for (int i = 0; i < 200; i++) {
+ long actual = writer.calculateBackoffMillis(20);
+ Assertions.assertTrue(
+ actual >= base && actual <= 300000L,
+ "backoff at the maximum produced "
+ + actual
+ + ", expected ["
+ + base
+ + ", 300000]");
+ }
+ }
+
+ @Test
+ public void testBackoffMinimumNeverDropsAtTheCapBoundary() {
+ // The scheduled wait doubles until it hits the cap. Half the cap can
be
+ // less than the wait the previous retry already guaranteed, so
without a
+ // floor the first capped retry could sleep for less than the one
before
+ // it, which is not what anyone expects from exponential backoff.
+ int base = 100000;
+ AirtableSinkWriter writer = writerWithBackoff(base, 30);
+
+ // With a 100000ms base the schedule is 100000, 200000, then capped at
+ // 300000. The last wait that fitted under the cap was 200000, so no
+ // capped retry should ever draw below that. Half the cap, 150000,
would.
+ for (int retry = 3; retry <= 8; retry++) {
+ for (int i = 0; i < 200; i++) {
+ long actual = writer.calculateBackoffMillis(retry);
+ Assertions.assertTrue(
+ actual >= 200000L && actual <= 300000L,
+ "retry " + retry + " produced " + actual + ", expected
[200000, 300000]");
+ }
+ }
+ }
+
+ @Test
+ public void testBackoffStillVariesWhenBaseExceedsTheMaximum() {
+ AirtableSinkWriter writer = writerWithBackoff(300000, 30);
+
+ Set<Long> observed = new HashSet<>();
+ for (int i = 0; i < 200; i++) {
+ long actual = writer.calculateBackoffMillis(1);
+ Assertions.assertTrue(
+ actual >= 150000L && actual <= 300000L,
+ "backoff produced " + actual + ", expected [150000,
300000]");
+ observed.add(actual);
+ }
+
+ Assertions.assertTrue(
+ observed.size() > 1,
+ "backoff must still vary when the configured base is at or
above the maximum, "
+ + "but observed only "
+ + observed);
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/airtable/source/AirtableSourceReaderTest.java
b/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/airtable/source/AirtableSourceReaderTest.java
index a26be6763f..774d3d988d 100644
---
a/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/airtable/source/AirtableSourceReaderTest.java
+++
b/seatunnel-connectors-v2/connector-http/connector-http-airtable/src/test/java/org/apache/seatunnel/connectors/seatunnel/airtable/source/AirtableSourceReaderTest.java
@@ -33,6 +33,9 @@ import org.junit.jupiter.api.Test;
import org.mockito.Mock;
import org.mockito.MockitoAnnotations;
+import java.util.HashSet;
+import java.util.Set;
+
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.ArgumentMatchers.anyString;
@@ -99,4 +102,162 @@ public class AirtableSourceReaderTest {
verify(httpClient, times(2))
.execute(anyString(), anyString(), any(), any(), any(),
anyBoolean());
}
+
+ @Test
+ public void testBackoffIsZeroWhenDisabled() {
+ AirtableSourceReader reader =
+ new AirtableSourceReader(parameter, context, schema, null,
null, null, 0, 0, 3);
+
+ Assertions.assertEquals(0L, reader.calculateBackoffMillis(1));
+ Assertions.assertEquals(0L, reader.calculateBackoffMillis(5));
+ }
+
+ @Test
+ public void testBackoffNeverShorterThanConfigured() {
+ int base = 1000;
+ AirtableSourceReader reader =
+ new AirtableSourceReader(parameter, context, schema, null,
null, null, 0, base, 5);
+
+ for (int retry = 1; retry <= 5; retry++) {
+ long base_ = (long) base * (1L << (retry - 1));
+ long ceiling = Math.min(2 * base_, 300000L);
+ for (int i = 0; i < 200; i++) {
+ long actual = reader.calculateBackoffMillis(retry);
+ Assertions.assertTrue(
+ actual >= base_ && actual <= ceiling,
+ "retry "
+ + retry
+ + " produced "
+ + actual
+ + ", expected ["
+ + base_
+ + ", "
+ + ceiling
+ + "]");
+ }
+ }
+ }
+
+ @Test
+ public void testBackoffIsNotDeterministic() {
+ AirtableSourceReader reader =
+ new AirtableSourceReader(parameter, context, schema, null,
null, null, 0, 1000, 3);
+
+ Set<Long> observed = new HashSet<>();
+ for (int i = 0; i < 200; i++) {
+ observed.add(reader.calculateBackoffMillis(3));
+ }
+
+ // Without jitter every call returns the same value. At retry 3 with a
+ // 1000ms base the range is [4000, 8000], so seeing a single value
across
+ // 200 draws is not chance.
+ Assertions.assertTrue(
+ observed.size() > 1,
+ "backoff must vary between calls so that concurrent clients do
not "
+ + "retry in lockstep, but observed only "
+ + observed);
+ }
+
+ @Test
+ public void testBackoffRespectsMaximum() {
+ int base = 60000;
+ AirtableSourceReader reader =
+ new AirtableSourceReader(parameter, context, schema, null,
null, null, 0, base, 30);
+
+ for (int i = 0; i < 100; i++) {
+ Assertions.assertTrue(
+ reader.calculateBackoffMillis(20) <= 300000L,
+ "jittered backoff must never exceed MAX_BACKOFF_MILLIS");
+ }
+ }
+
+ @Test
+ public void testBackoffStillVariesAtTheMaximum() {
+ int base = 60000;
+ AirtableSourceReader reader =
+ new AirtableSourceReader(parameter, context, schema, null,
null, null, 0, base, 30);
+
+ Set<Long> observed = new HashSet<>();
+ for (int i = 0; i < 200; i++) {
+ observed.add(reader.calculateBackoffMillis(20));
+ }
+
+ // There is no headroom to add into once the wait reaches the maximum,
so
+ // the spread has to go downwards. Leaving it flat would put every
caller
+ // back in lockstep for exactly the retries that matter most.
+ Assertions.assertTrue(
+ observed.size() > 1,
+ "backoff must still vary once it reaches MAX_BACKOFF_MILLIS,
but "
+ + "observed only "
+ + observed);
+ }
+
+ @Test
+ public void testBackoffMinimumNeverDropsAtTheCapBoundary() {
+ // The scheduled wait doubles until it hits the cap. Half the cap can
be
+ // less than the wait the previous retry already guaranteed, so
without a
+ // floor the first capped retry could sleep for less than the one
before
+ // it, which is not what anyone expects from exponential backoff.
+ int base = 100000;
+ AirtableSourceReader reader =
+ new AirtableSourceReader(parameter, context, schema, null,
null, null, 0, base, 30);
+
+ // With a 100000ms base the schedule is 100000, 200000, then capped at
+ // 300000. The last wait that fitted under the cap was 200000, so no
+ // capped retry should ever draw below that. Half the cap, 150000,
would.
+ for (int retry = 3; retry <= 8; retry++) {
+ for (int i = 0; i < 200; i++) {
+ long actual = reader.calculateBackoffMillis(retry);
+ Assertions.assertTrue(
+ actual >= 200000L && actual <= 300000L,
+ "retry " + retry + " produced " + actual + ", expected
[200000, 300000]");
+ }
+ }
+ }
+
+ @Test
+ public void testBackoffStillVariesWhenBaseExceedsTheMaximum() {
+ // A base at or above MAX_BACKOFF_MILLIS pins waitMillis to the cap
from
+ // the first retry. Flooring at the base would then leave no room to
+ // jitter, which is the lockstep this change exists to remove, and it
+ // would bite the operators running the most conservative backoff.
+ int base = 300000;
+ AirtableSourceReader reader =
+ new AirtableSourceReader(parameter, context, schema, null,
null, null, 0, base, 30);
+
+ Set<Long> observed = new HashSet<>();
+ for (int i = 0; i < 200; i++) {
+ long actual = reader.calculateBackoffMillis(1);
+ Assertions.assertTrue(
+ actual >= 150000L && actual <= 300000L,
+ "backoff produced " + actual + ", expected [150000,
300000]");
+ observed.add(actual);
+ }
+
+ Assertions.assertTrue(
+ observed.size() > 1,
+ "backoff must still vary when the configured base is at or
above the maximum, "
+ + "but observed only "
+ + observed);
+ }
+
+ @Test
+ public void testBackoffAtTheMaximumNeverDropsBelowConfiguredBase() {
+ // A base above half the maximum makes the downward spread collide with
+ // the configured backoff, which Airtable's 429 handling depends on.
+ int base = 200000;
+ AirtableSourceReader reader =
+ new AirtableSourceReader(parameter, context, schema, null,
null, null, 0, base, 30);
+
+ for (int i = 0; i < 200; i++) {
+ long actual = reader.calculateBackoffMillis(20);
+ Assertions.assertTrue(
+ actual >= base && actual <= 300000L,
+ "backoff at the maximum produced "
+ + actual
+ + ", expected ["
+ + base
+ + ", 300000]");
+ }
+ }
}