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]");
+        }
+    }
 }

Reply via email to