airbyte_cdk.sources.declarative.retrievers.page_size_reducer

  1# Copyright (c) 2026 Airbyte, Inc., all rights reserved.
  2
  3import logging
  4import time
  5from dataclasses import dataclass
  6from enum import Enum
  7from typing import Callable, Optional
  8
  9from airbyte_cdk.models import FailureType
 10from airbyte_cdk.utils.traced_exception import AirbyteTracedException
 11
 12LOGGER = logging.getLogger("airbyte")
 13
 14
 15class PageSizeResetPolicy(Enum):
 16    NEVER = "NEVER"
 17    AFTER_SUCCESSFUL_PAGE = "AFTER_SUCCESSFUL_PAGE"
 18
 19
 20@dataclass(frozen=True)
 21class PageSizeReduction:
 22    """
 23    How much to shrink the page size when an error handler resolves to `ResponseAction.REDUCE_PAGE_SIZE`.
 24
 25    This is immutable configuration: it is created once per stream and shared, while the page size in effect
 26    lives in a `PageSizeReducer` created per partition read.
 27    """
 28
 29    reduction_factor: float = 2.0
 30    minimum_page_size: int = 1
 31    max_attempts: int = 5
 32    reset_policy: PageSizeResetPolicy = PageSizeResetPolicy.NEVER
 33    # Base wait before a page is re-issued, multiplied by the number of attempts made in a row. A
 34    # `REDUCE_PAGE_SIZE` response never reaches the HTTP retry budget - the exception raised for it is
 35    # deliberately not a backoff exception - so this is the only thing spacing these requests out, and an
 36    # API whose 502 means "we are briefly unwell" rather than "your page is too big" needs it to be more
 37    # than a token pause.
 38    backoff_seconds: float = 0.5
 39    # How many times the same page is re-issued unchanged once the page size cannot be shrunk any further,
 40    # before the read gives up. Zero keeps the strict behaviour: the first response that cannot be answered
 41    # with a smaller page fails the stream. It is what gives an API whose error is transient a budget at the
 42    # floor, where reducing is no longer an option but waiting still is.
 43    retries_at_minimum_page_size: int = 0
 44    # Appended to the two messages raised once the page size cannot be reduced any further. Those are
 45    # `transient_error`s the CDK has no remediation for - it only knows that the API rejected every page size
 46    # asked for - while the connector knows what narrows a query down on this particular API.
 47    failure_message: Optional[str] = None
 48
 49    def __post_init__(self) -> None:
 50        if self.reduction_factor <= 1:
 51            raise ValueError(
 52                f"The page size reduction factor needs to be greater than 1. Got {self.reduction_factor}"
 53            )
 54        if self.minimum_page_size < 1:
 55            raise ValueError(
 56                f"The minimum page size needs to be strictly positive. Got {self.minimum_page_size}"
 57            )
 58        if self.max_attempts < 1:
 59            raise ValueError(
 60                f"The maximum number of page size reductions needs to be strictly positive. Got {self.max_attempts}"
 61            )
 62        if self.backoff_seconds < 0:
 63            raise ValueError(
 64                f"The wait between page size reductions cannot be negative. Got {self.backoff_seconds}"
 65            )
 66        if self.retries_at_minimum_page_size < 0:
 67            raise ValueError(
 68                f"The number of retries at the minimum page size cannot be negative. Got {self.retries_at_minimum_page_size}"
 69            )
 70
 71
 72class PageSizeReducer:
 73    """
 74    Tracks the page size to use while reading one partition when the API asks for smaller pages.
 75
 76    One instance is created per `SimpleRetriever._read_pages` call. This is deliberate: a retriever and its
 77    paginator are shared by every partition of a stream and partitions are read concurrently, so the reduced
 78    page size must not be stored on the paginator or on the retriever.
 79    """
 80
 81    def __init__(
 82        self,
 83        config: PageSizeReduction,
 84        configured_page_size: Optional[int],
 85        stream_name: str = "",
 86        sleep: Optional[Callable[[float], None]] = None,
 87    ) -> None:
 88        self._config = config
 89        self._configured_page_size = configured_page_size
 90        self._stream_name = stream_name
 91        # Resolved on each call rather than bound here: a default of `time.sleep` would capture
 92        # the function object, and a test that patches `time.sleep` to keep a run of reductions
 93        # from taking its two minutes for real would have no effect on an already-bound default.
 94        self._sleep_override = sleep
 95        self._current_page_size: Optional[int] = None
 96        self._attempts = 0
 97        self._total_reductions = 0
 98        self._retries_at_minimum_page_size = 0
 99
100    def _sleep(self, seconds: float) -> None:
101        if self._sleep_override is not None:
102            self._sleep_override(seconds)
103            return
104        time.sleep(seconds)
105
106    @property
107    def page_size_override(self) -> Optional[int]:
108        """
109        :return: the reduced page size to request, or None while the configured page size is in effect
110        """
111        return self._current_page_size
112
113    def reduce(self) -> None:
114        """
115        Shrink the page size used for the next request. Raises once the page size cannot be shrunk any further
116        so that an endpoint that keeps failing does not loop forever.
117        """
118        if self._configured_page_size is None:
119            raise AirbyteTracedException(
120                internal_message=f"Stream {self._stream_name} received a REDUCE_PAGE_SIZE response action but its paginator does not inject a page size",
121                message=f"Stream {self._stream_name} is set up to reduce its page size on error but does not define "
122                f"one. Set `page_size` on the pagination strategy and `page_size_option` on the paginator.",
123                failure_type=FailureType.config_error,
124            )
125
126        current_page_size = (
127            self._current_page_size
128            if self._current_page_size is not None
129            else self._configured_page_size
130        )
131        if not isinstance(current_page_size, int) or isinstance(current_page_size, bool):
132            # A custom pagination strategy can return anything from `get_page_size`. Reducing
133            # is arithmetic, so a non-integer would otherwise fail with a bare TypeError in
134            # the middle of a sync.
135            raise AirbyteTracedException(
136                internal_message=f"Stream {self._stream_name} has a page size of type {type(current_page_size).__name__}: {current_page_size!r}",
137                message=f"The page size of stream {self._stream_name} is not a whole number, so the connector cannot "
138                f"reduce it. Make sure the pagination strategy's `page_size` is a number.",
139                failure_type=FailureType.config_error,
140            )
141
142        reduced_page_size = max(
143            self._config.minimum_page_size,
144            int(current_page_size // self._config.reduction_factor),
145        )
146        if reduced_page_size >= current_page_size:
147            # Nothing left to give up on the page size. Whether that is the end of the read is
148            # `retries_at_minimum_page_size`'s call, not this branch's: the response may still be transient.
149            self._retry_at_minimum_page_size(current_page_size)
150            return
151
152        self._attempts += 1
153        self._total_reductions += 1
154        if self._attempts > self._config.max_attempts:
155            # The budget counts the reductions that did *not* get a page through, which is what separates a
156            # partition that is stuck from one that is merely expensive. A partition where pages keep
157            # succeeding restarts this counter on each of them, under either reset policy, and reads to the
158            # end however many pages it has; a partition where nothing gets through burns the budget here.
159            raise AirbyteTracedException(
160                internal_message=f"Stream {self._stream_name} reduced its page size {self._attempts - 1} times in a row without a single page succeeding, which is the configured maximum of {self._config.max_attempts} ({self._total_reductions - 1} reductions so far while reading this partition)",
161                # `transient_error`, so the only remediation is the connector's own, if it defined one.
162                message=self._with_failure_message(
163                    f"The source keeps rejecting pages of stream {self._stream_name} at every page size the "
164                    f"connector requested, down to {current_page_size} records per page."
165                ),
166                failure_type=FailureType.transient_error,
167            )
168
169        backoff = self._config.backoff_seconds * self._attempts
170        LOGGER.info(
171            f"Reducing the page size of stream {self._stream_name} from {current_page_size} to {reduced_page_size} "
172            f"and retrying the same page in {backoff}s."
173        )
174        self._current_page_size = reduced_page_size
175        self._sleep(backoff)
176
177    def _retry_at_minimum_page_size(self, current_page_size: int) -> None:
178        """
179        Handle a `REDUCE_PAGE_SIZE` response that arrives when the page size is already as small as the
180        connector is allowed to request.
181
182        Reducing is out of options here, but re-issuing the page is not: an API that answers 502 to a page it
183        considers too heavy answers the same 502 when it is briefly unwell, and the error handler cannot tell
184        the two apart. `REDUCE_PAGE_SIZE` bypasses the HTTP retry budget, so without this budget the second
185        kind of 502 ends the stream on the first response once the floor is reached - fewer attempts than the
186        same connector got before it adopted the reduction.
187
188        The budget applies however the page size arrived at the floor, whether by reduction or because the
189        configured page size was already there. Only once it is spent does the failure depend on that: a page
190        size the connector's own floor blocks from ever being reduced is a configuration error, while a floor
191        the reduction walked down to means the API kept rejecting every size, which is transient.
192        """
193        # The retries come first, including on a stream whose page size started at the floor. How the page size
194        # arrived there says nothing about the response, so letting the misconfiguration branch below decide it
195        # would give the same manifest a retry budget or none depending on the user's `page_size`.
196        if self._retries_at_minimum_page_size < self._config.retries_at_minimum_page_size:
197            self._retries_at_minimum_page_size += 1
198            backoff = self._config.backoff_seconds * self._retries_at_minimum_page_size
199            LOGGER.info(
200                f"Stream {self._stream_name} cannot request a page smaller than {current_page_size} records, "
201                f"so the same page is retried unchanged in {backoff}s "
202                f"({self._retries_at_minimum_page_size} of {self._config.retries_at_minimum_page_size})."
203            )
204            self._sleep(backoff)
205            return
206
207        if self._current_page_size is None and self._config.minimum_page_size > 1:
208            # No reduction was ever applied and the connector's own floor is what blocks it, so the page size
209            # can never be reduced on this stream however the API behaves. That is a configuration error, and
210            # it is actionable: both numbers in the message are the connector's to change. It is reported once
211            # the retries above are spent, so a manifest that asked for them still gets them.
212            raise AirbyteTracedException(
213                internal_message=f"Stream {self._stream_name} has a configured page size of {current_page_size} which is not greater than the configured minimum page size of {self._config.minimum_page_size}, so it can never be reduced"
214                + (
215                    f" ({self._retries_at_minimum_page_size} retries at that size were spent first)"
216                    if self._retries_at_minimum_page_size
217                    else ""
218                ),
219                message=f"The page size of stream {self._stream_name} ({current_page_size}) is already at or below "
220                f"the configured minimum of {self._config.minimum_page_size}, so the connector cannot reduce it. "
221                f"Raise the page size of the stream, or lower `minimum_page_size`.",
222                failure_type=FailureType.config_error,
223            )
224
225        raise AirbyteTracedException(
226            internal_message=f"Stream {self._stream_name} still fails with a page size of {current_page_size}, which is the smallest page size allowed by the configured minimum of {self._config.minimum_page_size}"
227            + (
228                f", after {self._retries_at_minimum_page_size} retries at that size"
229                if self._retries_at_minimum_page_size
230                else ""
231            ),
232            # `transient_error`, so the only remediation is the connector's own, if it defined one.
233            message=self._with_failure_message(
234                f"The source keeps rejecting pages of stream {self._stream_name} at the smallest page size "
235                f"the connector is allowed to request ({current_page_size} records per page)."
236            ),
237            failure_type=FailureType.transient_error,
238        )
239
240    def _with_failure_message(self, message: str) -> str:
241        """
242        :return: the message followed by the connector's `failure_message`, when it defined one
243        """
244        failure_message = (self._config.failure_message or "").strip()
245        if not failure_message:
246            return message
247        return f"{message} {failure_message}"
248
249    def on_successful_page(self) -> None:
250        """
251        Called after each page that did not require a reduction.
252
253        The `max_attempts` budget restarts under both policies. It counts the reductions made *in a row*
254        without a single page succeeding, which is what separates a partition that is stuck from one that is
255        merely expensive: on a stream whose per-page cost varies - the GraphQL case this feature exists for -
256        a handful of heavy pages spread over a long partition is a healthy read, and a budget spanning the
257        whole partition would fail it at the `max_attempts + 1`-th heavy page while every reduction so far had
258        been followed by a successful page. The terminal message says the source rejected every page size the
259        connector asked for, so the budget has to mean exactly that.
260
261        The budget still terminates the read, because only a page that succeeded restarts it and only
262        `_read_pages` calls this, once per page it consumed. So between any two restarts the partition made one
263        page of progress, and the reductions that make no progress are bounded by `max_attempts`. Under `NEVER`
264        the reduced page size is never restored either, so it strictly decreases and `minimum_page_size` bounds
265        the reductions of the whole partition on its own.
266
267        Only `AFTER_SUCCESSFUL_PAGE` restores the page size. That policy is for an API that rejects the
268        configured page size on every page, so every page legitimately costs one reduction; `NEVER` keeps the
269        reduced size for the rest of the partition, which is the right behaviour when reductions are one-off.
270        """
271        self._attempts = 0
272        self._retries_at_minimum_page_size = 0
273
274        if self._config.reset_policy != PageSizeResetPolicy.AFTER_SUCCESSFUL_PAGE:
275            return
276
277        if self._current_page_size is not None:
278            LOGGER.info(
279                f"Restoring the page size of stream {self._stream_name} from {self._current_page_size} to {self._configured_page_size}."
280            )
281            self._current_page_size = None
LOGGER = <Logger airbyte (INFO)>
class PageSizeResetPolicy(enum.Enum):
16class PageSizeResetPolicy(Enum):
17    NEVER = "NEVER"
18    AFTER_SUCCESSFUL_PAGE = "AFTER_SUCCESSFUL_PAGE"
NEVER = <PageSizeResetPolicy.NEVER: 'NEVER'>
AFTER_SUCCESSFUL_PAGE = <PageSizeResetPolicy.AFTER_SUCCESSFUL_PAGE: 'AFTER_SUCCESSFUL_PAGE'>
@dataclass(frozen=True)
class PageSizeReduction:
21@dataclass(frozen=True)
22class PageSizeReduction:
23    """
24    How much to shrink the page size when an error handler resolves to `ResponseAction.REDUCE_PAGE_SIZE`.
25
26    This is immutable configuration: it is created once per stream and shared, while the page size in effect
27    lives in a `PageSizeReducer` created per partition read.
28    """
29
30    reduction_factor: float = 2.0
31    minimum_page_size: int = 1
32    max_attempts: int = 5
33    reset_policy: PageSizeResetPolicy = PageSizeResetPolicy.NEVER
34    # Base wait before a page is re-issued, multiplied by the number of attempts made in a row. A
35    # `REDUCE_PAGE_SIZE` response never reaches the HTTP retry budget - the exception raised for it is
36    # deliberately not a backoff exception - so this is the only thing spacing these requests out, and an
37    # API whose 502 means "we are briefly unwell" rather than "your page is too big" needs it to be more
38    # than a token pause.
39    backoff_seconds: float = 0.5
40    # How many times the same page is re-issued unchanged once the page size cannot be shrunk any further,
41    # before the read gives up. Zero keeps the strict behaviour: the first response that cannot be answered
42    # with a smaller page fails the stream. It is what gives an API whose error is transient a budget at the
43    # floor, where reducing is no longer an option but waiting still is.
44    retries_at_minimum_page_size: int = 0
45    # Appended to the two messages raised once the page size cannot be reduced any further. Those are
46    # `transient_error`s the CDK has no remediation for - it only knows that the API rejected every page size
47    # asked for - while the connector knows what narrows a query down on this particular API.
48    failure_message: Optional[str] = None
49
50    def __post_init__(self) -> None:
51        if self.reduction_factor <= 1:
52            raise ValueError(
53                f"The page size reduction factor needs to be greater than 1. Got {self.reduction_factor}"
54            )
55        if self.minimum_page_size < 1:
56            raise ValueError(
57                f"The minimum page size needs to be strictly positive. Got {self.minimum_page_size}"
58            )
59        if self.max_attempts < 1:
60            raise ValueError(
61                f"The maximum number of page size reductions needs to be strictly positive. Got {self.max_attempts}"
62            )
63        if self.backoff_seconds < 0:
64            raise ValueError(
65                f"The wait between page size reductions cannot be negative. Got {self.backoff_seconds}"
66            )
67        if self.retries_at_minimum_page_size < 0:
68            raise ValueError(
69                f"The number of retries at the minimum page size cannot be negative. Got {self.retries_at_minimum_page_size}"
70            )

How much to shrink the page size when an error handler resolves to ResponseAction.REDUCE_PAGE_SIZE.

This is immutable configuration: it is created once per stream and shared, while the page size in effect lives in a PageSizeReducer created per partition read.

PageSizeReduction( reduction_factor: float = 2.0, minimum_page_size: int = 1, max_attempts: int = 5, reset_policy: PageSizeResetPolicy = <PageSizeResetPolicy.NEVER: 'NEVER'>, backoff_seconds: float = 0.5, retries_at_minimum_page_size: int = 0, failure_message: Optional[str] = None)
reduction_factor: float = 2.0
minimum_page_size: int = 1
max_attempts: int = 5
reset_policy: PageSizeResetPolicy = <PageSizeResetPolicy.NEVER: 'NEVER'>
backoff_seconds: float = 0.5
retries_at_minimum_page_size: int = 0
failure_message: Optional[str] = None
class PageSizeReducer:
 73class PageSizeReducer:
 74    """
 75    Tracks the page size to use while reading one partition when the API asks for smaller pages.
 76
 77    One instance is created per `SimpleRetriever._read_pages` call. This is deliberate: a retriever and its
 78    paginator are shared by every partition of a stream and partitions are read concurrently, so the reduced
 79    page size must not be stored on the paginator or on the retriever.
 80    """
 81
 82    def __init__(
 83        self,
 84        config: PageSizeReduction,
 85        configured_page_size: Optional[int],
 86        stream_name: str = "",
 87        sleep: Optional[Callable[[float], None]] = None,
 88    ) -> None:
 89        self._config = config
 90        self._configured_page_size = configured_page_size
 91        self._stream_name = stream_name
 92        # Resolved on each call rather than bound here: a default of `time.sleep` would capture
 93        # the function object, and a test that patches `time.sleep` to keep a run of reductions
 94        # from taking its two minutes for real would have no effect on an already-bound default.
 95        self._sleep_override = sleep
 96        self._current_page_size: Optional[int] = None
 97        self._attempts = 0
 98        self._total_reductions = 0
 99        self._retries_at_minimum_page_size = 0
100
101    def _sleep(self, seconds: float) -> None:
102        if self._sleep_override is not None:
103            self._sleep_override(seconds)
104            return
105        time.sleep(seconds)
106
107    @property
108    def page_size_override(self) -> Optional[int]:
109        """
110        :return: the reduced page size to request, or None while the configured page size is in effect
111        """
112        return self._current_page_size
113
114    def reduce(self) -> None:
115        """
116        Shrink the page size used for the next request. Raises once the page size cannot be shrunk any further
117        so that an endpoint that keeps failing does not loop forever.
118        """
119        if self._configured_page_size is None:
120            raise AirbyteTracedException(
121                internal_message=f"Stream {self._stream_name} received a REDUCE_PAGE_SIZE response action but its paginator does not inject a page size",
122                message=f"Stream {self._stream_name} is set up to reduce its page size on error but does not define "
123                f"one. Set `page_size` on the pagination strategy and `page_size_option` on the paginator.",
124                failure_type=FailureType.config_error,
125            )
126
127        current_page_size = (
128            self._current_page_size
129            if self._current_page_size is not None
130            else self._configured_page_size
131        )
132        if not isinstance(current_page_size, int) or isinstance(current_page_size, bool):
133            # A custom pagination strategy can return anything from `get_page_size`. Reducing
134            # is arithmetic, so a non-integer would otherwise fail with a bare TypeError in
135            # the middle of a sync.
136            raise AirbyteTracedException(
137                internal_message=f"Stream {self._stream_name} has a page size of type {type(current_page_size).__name__}: {current_page_size!r}",
138                message=f"The page size of stream {self._stream_name} is not a whole number, so the connector cannot "
139                f"reduce it. Make sure the pagination strategy's `page_size` is a number.",
140                failure_type=FailureType.config_error,
141            )
142
143        reduced_page_size = max(
144            self._config.minimum_page_size,
145            int(current_page_size // self._config.reduction_factor),
146        )
147        if reduced_page_size >= current_page_size:
148            # Nothing left to give up on the page size. Whether that is the end of the read is
149            # `retries_at_minimum_page_size`'s call, not this branch's: the response may still be transient.
150            self._retry_at_minimum_page_size(current_page_size)
151            return
152
153        self._attempts += 1
154        self._total_reductions += 1
155        if self._attempts > self._config.max_attempts:
156            # The budget counts the reductions that did *not* get a page through, which is what separates a
157            # partition that is stuck from one that is merely expensive. A partition where pages keep
158            # succeeding restarts this counter on each of them, under either reset policy, and reads to the
159            # end however many pages it has; a partition where nothing gets through burns the budget here.
160            raise AirbyteTracedException(
161                internal_message=f"Stream {self._stream_name} reduced its page size {self._attempts - 1} times in a row without a single page succeeding, which is the configured maximum of {self._config.max_attempts} ({self._total_reductions - 1} reductions so far while reading this partition)",
162                # `transient_error`, so the only remediation is the connector's own, if it defined one.
163                message=self._with_failure_message(
164                    f"The source keeps rejecting pages of stream {self._stream_name} at every page size the "
165                    f"connector requested, down to {current_page_size} records per page."
166                ),
167                failure_type=FailureType.transient_error,
168            )
169
170        backoff = self._config.backoff_seconds * self._attempts
171        LOGGER.info(
172            f"Reducing the page size of stream {self._stream_name} from {current_page_size} to {reduced_page_size} "
173            f"and retrying the same page in {backoff}s."
174        )
175        self._current_page_size = reduced_page_size
176        self._sleep(backoff)
177
178    def _retry_at_minimum_page_size(self, current_page_size: int) -> None:
179        """
180        Handle a `REDUCE_PAGE_SIZE` response that arrives when the page size is already as small as the
181        connector is allowed to request.
182
183        Reducing is out of options here, but re-issuing the page is not: an API that answers 502 to a page it
184        considers too heavy answers the same 502 when it is briefly unwell, and the error handler cannot tell
185        the two apart. `REDUCE_PAGE_SIZE` bypasses the HTTP retry budget, so without this budget the second
186        kind of 502 ends the stream on the first response once the floor is reached - fewer attempts than the
187        same connector got before it adopted the reduction.
188
189        The budget applies however the page size arrived at the floor, whether by reduction or because the
190        configured page size was already there. Only once it is spent does the failure depend on that: a page
191        size the connector's own floor blocks from ever being reduced is a configuration error, while a floor
192        the reduction walked down to means the API kept rejecting every size, which is transient.
193        """
194        # The retries come first, including on a stream whose page size started at the floor. How the page size
195        # arrived there says nothing about the response, so letting the misconfiguration branch below decide it
196        # would give the same manifest a retry budget or none depending on the user's `page_size`.
197        if self._retries_at_minimum_page_size < self._config.retries_at_minimum_page_size:
198            self._retries_at_minimum_page_size += 1
199            backoff = self._config.backoff_seconds * self._retries_at_minimum_page_size
200            LOGGER.info(
201                f"Stream {self._stream_name} cannot request a page smaller than {current_page_size} records, "
202                f"so the same page is retried unchanged in {backoff}s "
203                f"({self._retries_at_minimum_page_size} of {self._config.retries_at_minimum_page_size})."
204            )
205            self._sleep(backoff)
206            return
207
208        if self._current_page_size is None and self._config.minimum_page_size > 1:
209            # No reduction was ever applied and the connector's own floor is what blocks it, so the page size
210            # can never be reduced on this stream however the API behaves. That is a configuration error, and
211            # it is actionable: both numbers in the message are the connector's to change. It is reported once
212            # the retries above are spent, so a manifest that asked for them still gets them.
213            raise AirbyteTracedException(
214                internal_message=f"Stream {self._stream_name} has a configured page size of {current_page_size} which is not greater than the configured minimum page size of {self._config.minimum_page_size}, so it can never be reduced"
215                + (
216                    f" ({self._retries_at_minimum_page_size} retries at that size were spent first)"
217                    if self._retries_at_minimum_page_size
218                    else ""
219                ),
220                message=f"The page size of stream {self._stream_name} ({current_page_size}) is already at or below "
221                f"the configured minimum of {self._config.minimum_page_size}, so the connector cannot reduce it. "
222                f"Raise the page size of the stream, or lower `minimum_page_size`.",
223                failure_type=FailureType.config_error,
224            )
225
226        raise AirbyteTracedException(
227            internal_message=f"Stream {self._stream_name} still fails with a page size of {current_page_size}, which is the smallest page size allowed by the configured minimum of {self._config.minimum_page_size}"
228            + (
229                f", after {self._retries_at_minimum_page_size} retries at that size"
230                if self._retries_at_minimum_page_size
231                else ""
232            ),
233            # `transient_error`, so the only remediation is the connector's own, if it defined one.
234            message=self._with_failure_message(
235                f"The source keeps rejecting pages of stream {self._stream_name} at the smallest page size "
236                f"the connector is allowed to request ({current_page_size} records per page)."
237            ),
238            failure_type=FailureType.transient_error,
239        )
240
241    def _with_failure_message(self, message: str) -> str:
242        """
243        :return: the message followed by the connector's `failure_message`, when it defined one
244        """
245        failure_message = (self._config.failure_message or "").strip()
246        if not failure_message:
247            return message
248        return f"{message} {failure_message}"
249
250    def on_successful_page(self) -> None:
251        """
252        Called after each page that did not require a reduction.
253
254        The `max_attempts` budget restarts under both policies. It counts the reductions made *in a row*
255        without a single page succeeding, which is what separates a partition that is stuck from one that is
256        merely expensive: on a stream whose per-page cost varies - the GraphQL case this feature exists for -
257        a handful of heavy pages spread over a long partition is a healthy read, and a budget spanning the
258        whole partition would fail it at the `max_attempts + 1`-th heavy page while every reduction so far had
259        been followed by a successful page. The terminal message says the source rejected every page size the
260        connector asked for, so the budget has to mean exactly that.
261
262        The budget still terminates the read, because only a page that succeeded restarts it and only
263        `_read_pages` calls this, once per page it consumed. So between any two restarts the partition made one
264        page of progress, and the reductions that make no progress are bounded by `max_attempts`. Under `NEVER`
265        the reduced page size is never restored either, so it strictly decreases and `minimum_page_size` bounds
266        the reductions of the whole partition on its own.
267
268        Only `AFTER_SUCCESSFUL_PAGE` restores the page size. That policy is for an API that rejects the
269        configured page size on every page, so every page legitimately costs one reduction; `NEVER` keeps the
270        reduced size for the rest of the partition, which is the right behaviour when reductions are one-off.
271        """
272        self._attempts = 0
273        self._retries_at_minimum_page_size = 0
274
275        if self._config.reset_policy != PageSizeResetPolicy.AFTER_SUCCESSFUL_PAGE:
276            return
277
278        if self._current_page_size is not None:
279            LOGGER.info(
280                f"Restoring the page size of stream {self._stream_name} from {self._current_page_size} to {self._configured_page_size}."
281            )
282            self._current_page_size = None

Tracks the page size to use while reading one partition when the API asks for smaller pages.

One instance is created per SimpleRetriever._read_pages call. This is deliberate: a retriever and its paginator are shared by every partition of a stream and partitions are read concurrently, so the reduced page size must not be stored on the paginator or on the retriever.

PageSizeReducer( config: PageSizeReduction, configured_page_size: Optional[int], stream_name: str = '', sleep: Optional[Callable[[float], NoneType]] = None)
82    def __init__(
83        self,
84        config: PageSizeReduction,
85        configured_page_size: Optional[int],
86        stream_name: str = "",
87        sleep: Optional[Callable[[float], None]] = None,
88    ) -> None:
89        self._config = config
90        self._configured_page_size = configured_page_size
91        self._stream_name = stream_name
92        # Resolved on each call rather than bound here: a default of `time.sleep` would capture
93        # the function object, and a test that patches `time.sleep` to keep a run of reductions
94        # from taking its two minutes for real would have no effect on an already-bound default.
95        self._sleep_override = sleep
96        self._current_page_size: Optional[int] = None
97        self._attempts = 0
98        self._total_reductions = 0
99        self._retries_at_minimum_page_size = 0
page_size_override: Optional[int]
107    @property
108    def page_size_override(self) -> Optional[int]:
109        """
110        :return: the reduced page size to request, or None while the configured page size is in effect
111        """
112        return self._current_page_size
Returns

the reduced page size to request, or None while the configured page size is in effect

def reduce(self) -> None:
114    def reduce(self) -> None:
115        """
116        Shrink the page size used for the next request. Raises once the page size cannot be shrunk any further
117        so that an endpoint that keeps failing does not loop forever.
118        """
119        if self._configured_page_size is None:
120            raise AirbyteTracedException(
121                internal_message=f"Stream {self._stream_name} received a REDUCE_PAGE_SIZE response action but its paginator does not inject a page size",
122                message=f"Stream {self._stream_name} is set up to reduce its page size on error but does not define "
123                f"one. Set `page_size` on the pagination strategy and `page_size_option` on the paginator.",
124                failure_type=FailureType.config_error,
125            )
126
127        current_page_size = (
128            self._current_page_size
129            if self._current_page_size is not None
130            else self._configured_page_size
131        )
132        if not isinstance(current_page_size, int) or isinstance(current_page_size, bool):
133            # A custom pagination strategy can return anything from `get_page_size`. Reducing
134            # is arithmetic, so a non-integer would otherwise fail with a bare TypeError in
135            # the middle of a sync.
136            raise AirbyteTracedException(
137                internal_message=f"Stream {self._stream_name} has a page size of type {type(current_page_size).__name__}: {current_page_size!r}",
138                message=f"The page size of stream {self._stream_name} is not a whole number, so the connector cannot "
139                f"reduce it. Make sure the pagination strategy's `page_size` is a number.",
140                failure_type=FailureType.config_error,
141            )
142
143        reduced_page_size = max(
144            self._config.minimum_page_size,
145            int(current_page_size // self._config.reduction_factor),
146        )
147        if reduced_page_size >= current_page_size:
148            # Nothing left to give up on the page size. Whether that is the end of the read is
149            # `retries_at_minimum_page_size`'s call, not this branch's: the response may still be transient.
150            self._retry_at_minimum_page_size(current_page_size)
151            return
152
153        self._attempts += 1
154        self._total_reductions += 1
155        if self._attempts > self._config.max_attempts:
156            # The budget counts the reductions that did *not* get a page through, which is what separates a
157            # partition that is stuck from one that is merely expensive. A partition where pages keep
158            # succeeding restarts this counter on each of them, under either reset policy, and reads to the
159            # end however many pages it has; a partition where nothing gets through burns the budget here.
160            raise AirbyteTracedException(
161                internal_message=f"Stream {self._stream_name} reduced its page size {self._attempts - 1} times in a row without a single page succeeding, which is the configured maximum of {self._config.max_attempts} ({self._total_reductions - 1} reductions so far while reading this partition)",
162                # `transient_error`, so the only remediation is the connector's own, if it defined one.
163                message=self._with_failure_message(
164                    f"The source keeps rejecting pages of stream {self._stream_name} at every page size the "
165                    f"connector requested, down to {current_page_size} records per page."
166                ),
167                failure_type=FailureType.transient_error,
168            )
169
170        backoff = self._config.backoff_seconds * self._attempts
171        LOGGER.info(
172            f"Reducing the page size of stream {self._stream_name} from {current_page_size} to {reduced_page_size} "
173            f"and retrying the same page in {backoff}s."
174        )
175        self._current_page_size = reduced_page_size
176        self._sleep(backoff)

Shrink the page size used for the next request. Raises once the page size cannot be shrunk any further so that an endpoint that keeps failing does not loop forever.

def on_successful_page(self) -> None:
250    def on_successful_page(self) -> None:
251        """
252        Called after each page that did not require a reduction.
253
254        The `max_attempts` budget restarts under both policies. It counts the reductions made *in a row*
255        without a single page succeeding, which is what separates a partition that is stuck from one that is
256        merely expensive: on a stream whose per-page cost varies - the GraphQL case this feature exists for -
257        a handful of heavy pages spread over a long partition is a healthy read, and a budget spanning the
258        whole partition would fail it at the `max_attempts + 1`-th heavy page while every reduction so far had
259        been followed by a successful page. The terminal message says the source rejected every page size the
260        connector asked for, so the budget has to mean exactly that.
261
262        The budget still terminates the read, because only a page that succeeded restarts it and only
263        `_read_pages` calls this, once per page it consumed. So between any two restarts the partition made one
264        page of progress, and the reductions that make no progress are bounded by `max_attempts`. Under `NEVER`
265        the reduced page size is never restored either, so it strictly decreases and `minimum_page_size` bounds
266        the reductions of the whole partition on its own.
267
268        Only `AFTER_SUCCESSFUL_PAGE` restores the page size. That policy is for an API that rejects the
269        configured page size on every page, so every page legitimately costs one reduction; `NEVER` keeps the
270        reduced size for the rest of the partition, which is the right behaviour when reductions are one-off.
271        """
272        self._attempts = 0
273        self._retries_at_minimum_page_size = 0
274
275        if self._config.reset_policy != PageSizeResetPolicy.AFTER_SUCCESSFUL_PAGE:
276            return
277
278        if self._current_page_size is not None:
279            LOGGER.info(
280                f"Restoring the page size of stream {self._stream_name} from {self._current_page_size} to {self._configured_page_size}."
281            )
282            self._current_page_size = None

Called after each page that did not require a reduction.

The max_attempts budget restarts under both policies. It counts the reductions made in a row without a single page succeeding, which is what separates a partition that is stuck from one that is merely expensive: on a stream whose per-page cost varies - the GraphQL case this feature exists for - a handful of heavy pages spread over a long partition is a healthy read, and a budget spanning the whole partition would fail it at the max_attempts + 1-th heavy page while every reduction so far had been followed by a successful page. The terminal message says the source rejected every page size the connector asked for, so the budget has to mean exactly that.

The budget still terminates the read, because only a page that succeeded restarts it and only _read_pages calls this, once per page it consumed. So between any two restarts the partition made one page of progress, and the reductions that make no progress are bounded by max_attempts. Under NEVER the reduced page size is never restored either, so it strictly decreases and minimum_page_size bounds the reductions of the whole partition on its own.

Only AFTER_SUCCESSFUL_PAGE restores the page size. That policy is for an API that rejects the configured page size on every page, so every page legitimately costs one reduction; NEVER keeps the reduced size for the rest of the partition, which is the right behaviour when reductions are one-off.