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
16class PageSizeResetPolicy(Enum): 17 NEVER = "NEVER" 18 AFTER_SUCCESSFUL_PAGE = "AFTER_SUCCESSFUL_PAGE"
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.
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.
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
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
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.
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.