Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -2421,6 +2421,7 @@ definitions:
- RESET_PAGINATION
- RATE_LIMITED
- REFRESH_TOKEN_THEN_RETRY
- REDUCE_PAGE_SIZE
examples:
- SUCCESS
- FAIL
Expand All @@ -2429,6 +2430,7 @@ definitions:
- RESET_PAGINATION
- RATE_LIMITED
- REFRESH_TOKEN_THEN_RETRY
- REDUCE_PAGE_SIZE
failure_type:
title: Failure Type
description: Failure type of traced exception if a response matches the filter.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -557,6 +557,7 @@ class Action(Enum):
RESET_PAGINATION = "RESET_PAGINATION"
RATE_LIMITED = "RATE_LIMITED"
REFRESH_TOKEN_THEN_RETRY = "REFRESH_TOKEN_THEN_RETRY"
REDUCE_PAGE_SIZE = "REDUCE_PAGE_SIZE"


class FailureType(Enum):
Expand All @@ -578,6 +579,7 @@ class HttpResponseFilter(BaseModel):
"RESET_PAGINATION",
"RATE_LIMITED",
"REFRESH_TOKEN_THEN_RETRY",
"REDUCE_PAGE_SIZE",
],
title="Action",
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ def __post_init__(self, parameters: Mapping[str, Any]) -> None:
if not isinstance(page_size, int):
raise Exception(f"{page_size} is of type {type(page_size)}. Expected {int}")
self._page_size = page_size
self._default_page_size = self._page_size

@property
def initial_token(self) -> Optional[Any]:
Expand Down Expand Up @@ -104,3 +105,10 @@ def next_page_token(

def get_page_size(self) -> Optional[int]:
return self._page_size

def reduce_page_size(self) -> None:
if self._page_size is not None and self._page_size > 1:
self._page_size = max(1, self._page_size // 2)

def reset_page_size(self) -> None:
self._page_size = self._default_page_size
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,8 @@ def __post_init__(self, parameters: Mapping[str, Any]) -> None:
)
else:
self._page_size = None
self._default_page_size = self._page_size
self._effective_page_size: Optional[int] = None

@property
def initial_token(self) -> Optional[Any]:
Expand Down Expand Up @@ -103,10 +105,28 @@ def next_page_token(
return last_page_token_value + last_page_size

def get_page_size(self) -> Optional[int]:
if self._effective_page_size is not None:
return self._effective_page_size
if self._page_size:
page_size = self._page_size.eval(self.config)
if not isinstance(page_size, int):
raise Exception(f"{page_size} is of type {type(page_size)}. Expected {int}")
return page_size
else:
return None

def _get_default_page_size(self) -> Optional[int]:
if self._default_page_size:
page_size = self._default_page_size.eval(self.config)
if not isinstance(page_size, int):
raise Exception(f"{page_size} is of type {type(page_size)}. Expected {int}")
return page_size
return None

def reduce_page_size(self) -> None:
current = self.get_page_size()
if current is not None and current > 1:
self._effective_page_size = max(1, current // 2)

def reset_page_size(self) -> None:
self._effective_page_size = None
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ def __post_init__(self, parameters: Mapping[str, Any]) -> None:
if not isinstance(page_size, int):
raise Exception(f"{page_size} is of type {type(page_size)}. Expected {int}")
self._page_size = page_size
self._default_page_size = self._page_size

@property
def initial_token(self) -> Optional[Any]:
Expand Down Expand Up @@ -69,3 +70,10 @@ def next_page_token(

def get_page_size(self) -> Optional[int]:
return self._page_size

def reduce_page_size(self) -> None:
if self._page_size is not None and self._page_size > 1:
self._page_size = max(1, self._page_size // 2)

def reset_page_size(self) -> None:
self._page_size = self._default_page_size
Original file line number Diff line number Diff line change
Expand Up @@ -46,3 +46,17 @@ def get_page_size(self) -> Optional[int]:
"""
:return: page size: The number of records to fetch in a page. Returns None if unspecified
"""

def reduce_page_size(self) -> None:
"""Halve the current effective page size (floored at 1).

Called by `SimpleRetriever` when a `REDUCE_PAGE_SIZE` response action is received.
Subclasses that support dynamic page-size reduction should override this method.
"""

def reset_page_size(self) -> None:
"""Restore the page size to the originally configured default.

Called by `SimpleRetriever` after a successful page fetch following a reduction.
Subclasses that support dynamic page-size reduction should override this method.
"""
29 changes: 29 additions & 0 deletions airbyte_cdk/sources/declarative/retrievers/simple_retriever.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,9 @@
from airbyte_cdk.sources.declarative.partition_routers.single_partition_router import (
SinglePartitionRouter,
)
from airbyte_cdk.sources.declarative.requesters.paginators.default_paginator import (
DefaultPaginator,
)
from airbyte_cdk.sources.declarative.requesters.paginators.no_pagination import NoPagination
from airbyte_cdk.sources.declarative.requesters.paginators.paginator import Paginator
from airbyte_cdk.sources.declarative.requesters.query_properties import QueryProperties
Expand All @@ -41,6 +44,9 @@
from airbyte_cdk.sources.declarative.stream_slicers.stream_slicer import StreamSlicer
from airbyte_cdk.sources.source import ExperimentalClassWarning
from airbyte_cdk.sources.streams.core import StreamData
from airbyte_cdk.sources.streams.http.page_size_reduction_exception import (
PageSizeReductionRequiredException,
)
from airbyte_cdk.sources.streams.http.pagination_reset_exception import (
PaginationResetRequiredException,
)
Expand Down Expand Up @@ -403,7 +409,11 @@ def _read_pages(
yield current_record
except PaginationResetRequiredException:
reset_pagination = True
except PageSizeReductionRequiredException:
self._reduce_paginator_page_size()
continue
else:
self._reset_paginator_page_size()
if not response:
break

Expand Down Expand Up @@ -433,6 +443,25 @@ def _read_pages(
# Always return an empty generator just in case no records were ever yielded
yield from []

def _reduce_paginator_page_size(self) -> None:
"""Delegate page-size reduction to the paginator's `PaginationStrategy`, if available."""
if isinstance(self._paginator, DefaultPaginator):
strategy = self._paginator.pagination_strategy
previous = strategy.get_page_size()
strategy.reduce_page_size()
current = strategy.get_page_size()
LOGGER.info(
"Reducing page size for stream '%s' from %s to %s due to server error.",
self.name,
previous,
current,
)

def _reset_paginator_page_size(self) -> None:
"""Restore the paginator's page size to the configured default after a successful fetch."""
if isinstance(self._paginator, DefaultPaginator):
self._paginator.pagination_strategy.reset_page_size()

def _get_initial_next_page_token(self) -> Optional[Mapping[str, Any]]:
initial_token = self._paginator.get_initial_token()
next_page_token = {"next_page_token": initial_token} if initial_token is not None else None
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ class ResponseAction(Enum):
RESET_PAGINATION = "RESET_PAGINATION"
RATE_LIMITED = "RATE_LIMITED"
REFRESH_TOKEN_THEN_RETRY = "REFRESH_TOKEN_THEN_RETRY"
REDUCE_PAGE_SIZE = "REDUCE_PAGE_SIZE"


@dataclass
Expand Down
6 changes: 6 additions & 0 deletions airbyte_cdk/sources/streams/http/http_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,9 @@
RequestBodyException,
UserDefinedBackoffException,
)
from airbyte_cdk.sources.streams.http.page_size_reduction_exception import (
PageSizeReductionRequiredException,
)
from airbyte_cdk.sources.streams.http.pagination_reset_exception import (
PaginationResetRequiredException,
)
Expand Down Expand Up @@ -441,6 +444,9 @@ def _handle_error_resolution(
if error_resolution.response_action == ResponseAction.RESET_PAGINATION:
raise PaginationResetRequiredException()

if error_resolution.response_action == ResponseAction.REDUCE_PAGE_SIZE:
raise PageSizeReductionRequiredException()

# Emit stream status RUNNING with the reason RATE_LIMITED to log that the rate limit has been reached
if error_resolution.response_action == ResponseAction.RATE_LIMITED:
# TODO: Update to handle with message repository when concurrent message repository is ready
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
class PageSizeReductionRequiredException(Exception):
pass
Original file line number Diff line number Diff line change
Expand Up @@ -152,3 +152,51 @@ def test_interpolated_page_size_raises_on_non_integer():
config={"page_size": "invalid"},
parameters={},
)


@pytest.mark.parametrize(
"initial_page_size,expected_after_reduce",
[
pytest.param(100, 50, id="halve_100"),
pytest.param(10, 5, id="halve_10"),
pytest.param(3, 1, id="halve_3_floors_to_1"),
pytest.param(2, 1, id="halve_2_to_1"),
pytest.param(1, 1, id="already_at_minimum"),
],
)
def test_reduce_page_size(initial_page_size, expected_after_reduce):
strategy = CursorPaginationStrategy(
page_size=initial_page_size, cursor_value="token", config={}, parameters={}
)
strategy.reduce_page_size()
assert strategy.get_page_size() == expected_after_reduce


def test_reduce_page_size_multiple_times():
strategy = CursorPaginationStrategy(
page_size=100, cursor_value="token", config={}, parameters={}
)
strategy.reduce_page_size()
assert strategy.get_page_size() == 50
strategy.reduce_page_size()
assert strategy.get_page_size() == 25
strategy.reduce_page_size()
assert strategy.get_page_size() == 12


def test_reset_page_size_restores_default():
strategy = CursorPaginationStrategy(
page_size=100, cursor_value="token", config={}, parameters={}
)
strategy.reduce_page_size()
assert strategy.get_page_size() == 50
strategy.reset_page_size()
assert strategy.get_page_size() == 100


def test_reduce_page_size_noop_when_none():
strategy = CursorPaginationStrategy(
page_size=None, cursor_value="token", config={}, parameters={}
)
strategy.reduce_page_size()
assert strategy.get_page_size() is None
Original file line number Diff line number Diff line change
Expand Up @@ -146,3 +146,42 @@ def test_offset_increment_paginator_strategy_initial_token(
)

assert paginator_strategy.initial_token == expected_initial_token


@pytest.mark.parametrize(
"initial_page_size,expected_after_reduce",
[
pytest.param(100, 50, id="halve_100"),
pytest.param(10, 5, id="halve_10"),
pytest.param(3, 1, id="halve_3_floors_to_1"),
pytest.param(1, 1, id="already_at_minimum"),
],
)
def test_reduce_page_size(initial_page_size, expected_after_reduce):
strategy = OffsetIncrement(
page_size=initial_page_size, parameters={}, config={}, extractor=None
)
strategy.reduce_page_size()
assert strategy.get_page_size() == expected_after_reduce


def test_reduce_page_size_multiple_times():
strategy = OffsetIncrement(page_size=100, parameters={}, config={}, extractor=None)
strategy.reduce_page_size()
assert strategy.get_page_size() == 50
strategy.reduce_page_size()
assert strategy.get_page_size() == 25


def test_reset_page_size_restores_default():
strategy = OffsetIncrement(page_size=100, parameters={}, config={}, extractor=None)
strategy.reduce_page_size()
assert strategy.get_page_size() == 50
strategy.reset_page_size()
assert strategy.get_page_size() == 100


def test_reduce_page_size_noop_when_none():
strategy = OffsetIncrement(page_size=None, parameters={}, config={}, extractor=None)
strategy.reduce_page_size()
assert strategy.get_page_size() is None
Original file line number Diff line number Diff line change
Expand Up @@ -105,3 +105,40 @@ def test_page_increment_paginator_strategy_initial_token(
)

assert paginator_strategy.initial_token == expected_initial_token


@pytest.mark.parametrize(
"initial_page_size,expected_after_reduce",
[
pytest.param(100, 50, id="halve_100"),
pytest.param(10, 5, id="halve_10"),
pytest.param(3, 1, id="halve_3_floors_to_1"),
pytest.param(1, 1, id="already_at_minimum"),
],
)
def test_reduce_page_size(initial_page_size, expected_after_reduce):
strategy = PageIncrement(page_size=initial_page_size, parameters={}, config={})
strategy.reduce_page_size()
assert strategy.get_page_size() == expected_after_reduce


def test_reduce_page_size_multiple_times():
strategy = PageIncrement(page_size=100, parameters={}, config={})
strategy.reduce_page_size()
assert strategy.get_page_size() == 50
strategy.reduce_page_size()
assert strategy.get_page_size() == 25


def test_reset_page_size_restores_default():
strategy = PageIncrement(page_size=100, parameters={}, config={})
strategy.reduce_page_size()
assert strategy.get_page_size() == 50
strategy.reset_page_size()
assert strategy.get_page_size() == 100


def test_reduce_page_size_noop_when_none():
strategy = PageIncrement(page_size=None, parameters={}, config={})
strategy.reduce_page_size()
assert strategy.get_page_size() is None
Loading
Loading