airbyte_cdk.sources.declarative.retrievers.simple_retriever
1# 2# Copyright (c) 2025 Airbyte, Inc., all rights reserved. 3# 4 5import datetime 6import json 7import logging 8from collections import defaultdict 9from dataclasses import InitVar, dataclass, field 10from functools import partial 11from typing import ( 12 Any, 13 Callable, 14 Iterable, 15 List, 16 Mapping, 17 MutableMapping, 18 Optional, 19 Set, 20 Tuple, 21 Union, 22) 23 24import requests 25from typing_extensions import deprecated 26 27from airbyte_cdk.models import FailureType 28from airbyte_cdk.sources.declarative.extractors.http_selector import HttpSelector 29from airbyte_cdk.sources.declarative.extractors.record_filter import ( 30 ClientSideIncrementalRecordFilterDecorator, 31) 32from airbyte_cdk.sources.declarative.interpolation import InterpolatedString 33from airbyte_cdk.sources.declarative.partition_routers.single_partition_router import ( 34 SinglePartitionRouter, 35) 36from airbyte_cdk.sources.declarative.requesters.paginators.no_pagination import NoPagination 37from airbyte_cdk.sources.declarative.requesters.paginators.paginator import ( 38 Paginator, 39 page_size_override_kwargs, 40 stream_slice_kwargs, 41) 42from airbyte_cdk.sources.declarative.requesters.query_properties import QueryProperties 43from airbyte_cdk.sources.declarative.requesters.request_options import ( 44 DefaultRequestOptionsProvider, 45 RequestOptionsProvider, 46) 47from airbyte_cdk.sources.declarative.requesters.requester import Requester 48from airbyte_cdk.sources.declarative.retrievers.page_size_reducer import ( 49 PageSizeReducer, 50 PageSizeReduction, 51) 52from airbyte_cdk.sources.declarative.retrievers.pagination_tracker import PaginationTracker 53from airbyte_cdk.sources.declarative.retrievers.request_window_splitting import ( 54 RequestWindowSplitting, 55) 56from airbyte_cdk.sources.declarative.retrievers.retriever import Retriever 57from airbyte_cdk.sources.declarative.stream_slicers.stream_slicer import StreamSlicer 58from airbyte_cdk.sources.source import ExperimentalClassWarning 59from airbyte_cdk.sources.streams.core import StreamData 60from airbyte_cdk.sources.streams.http.page_size_reduction_exception import ( 61 PageSizeReductionNotSupportedException, 62 PageSizeReductionRequiredException, 63) 64from airbyte_cdk.sources.streams.http.pagination_reset_exception import ( 65 PaginationResetRequiredException, 66) 67from airbyte_cdk.sources.streams.http.request_window_split_exception import ( 68 RequestWindowSplitNotSupportedException, 69 RequestWindowSplitRequiredException, 70) 71from airbyte_cdk.sources.types import Config, Record, StreamSlice 72from airbyte_cdk.utils.mapping_helpers import combine_mappings 73from airbyte_cdk.utils.traced_exception import AirbyteTracedException 74 75FULL_REFRESH_SYNC_COMPLETE_KEY = "__ab_full_refresh_sync_complete" 76LOGGER = logging.getLogger("airbyte") 77 78# A defense-in-depth bound independent of any `request_window_splitter`'s own no-progress guard: protects 79# against a misbehaving custom cursor whose children don't actually shrink the window. Not user-configurable. 80_MAX_REQUEST_WINDOW_SPLIT_DEPTH = 10 81 82 83@dataclass 84class SimpleRetriever(Retriever): 85 """ 86 Retrieves records by synchronously sending requests to fetch records. 87 88 The retriever acts as an orchestrator between the requester, the record selector, the paginator, and the stream slicer. 89 90 For each stream slice, submit requests until there are no more pages of records to fetch. 91 92 This retriever currently inherits from HttpStream to reuse the request submission and pagination machinery. 93 As a result, some of the parameters passed to some methods are unused. 94 The two will be decoupled in a future release. 95 96 Attributes: 97 stream_name (str): The stream's name 98 stream_primary_key (Optional[Union[str, List[str], List[List[str]]]]): The stream's primary key 99 requester (Requester): The HTTP requester 100 record_selector (HttpSelector): The record selector 101 paginator (Optional[Paginator]): The paginator 102 stream_slicer (Optional[StreamSlicer]): The stream slicer 103 parameters (Mapping[str, Any]): Additional runtime parameters to be used for string interpolation 104 post_pagination_filter (Optional[ClientSideIncrementalRecordFilterDecorator]): Set for data feed streams only. 105 Records the cursor considers already synced are dropped once pagination has observed them 106 page_size_reduction (Optional[PageSizeReduction]): How much to shrink the page size when an error handler 107 resolves to `ResponseAction.REDUCE_PAGE_SIZE`. `None` disables page size reduction entirely. 108 It is immutable configuration; the page size in effect lives in a `PageSizeReducer` that 109 `_read_pages` creates per call, so the retriever and its paginator - both shared by every 110 partition of the stream, read concurrently - stay stateless. When `page_size_reduction` is 111 `None` no reducer is created and `_read_pages` keeps its previous behaviour 112 request_window_splitting (Optional[RequestWindowSplitting]): Policy applied when an error handler 113 resolves to `ResponseAction.SPLIT_REQUEST_WINDOW`, or when `RequestWindowSplitRequiredException` 114 is raised directly by custom code. `None` disables request-window splitting entirely: `read_records` 115 converts the exception into `RequestWindowSplitNotSupportedException` instead of attempting a 116 split it cannot perform. 117 request_window_splitter: Bound method - normally the stream's own cursor's `split_request_window` - 118 asked to replace a failing `StreamSlice` with smaller children. `None` has the same effect as 119 `request_window_splitting` being `None`. 120 """ 121 122 requester: Requester 123 record_selector: HttpSelector 124 config: Config 125 parameters: InitVar[Mapping[str, Any]] 126 name: str 127 _name: Union[InterpolatedString, str] = field(init=False, repr=False, default="") 128 primary_key: Optional[Union[str, List[str], List[List[str]]]] 129 _primary_key: str = field(init=False, repr=False, default="") 130 paginator: Optional[Paginator] = None 131 stream_slicer: StreamSlicer = field( 132 default_factory=lambda: SinglePartitionRouter(parameters={}) 133 ) 134 request_option_provider: RequestOptionsProvider = field( 135 default_factory=lambda: DefaultRequestOptionsProvider(parameters={}) 136 ) 137 ignore_stream_slicer_parameters_on_paginated_requests: bool = False 138 additional_query_properties: Optional[QueryProperties] = None 139 log_formatter: Optional[Callable[[requests.Response], Any]] = None 140 pagination_tracker_factory: Callable[[], PaginationTracker] = field( 141 default_factory=lambda: lambda: PaginationTracker() 142 ) 143 post_pagination_filter: Optional[ClientSideIncrementalRecordFilterDecorator] = None 144 page_size_reduction: Optional[PageSizeReduction] = None 145 request_window_splitting: Optional[RequestWindowSplitting] = None 146 request_window_splitter: Optional[ 147 Callable[[StreamSlice, Optional[datetime.timedelta]], Optional[List[StreamSlice]]] 148 ] = None 149 150 def __post_init__(self, parameters: Mapping[str, Any]) -> None: 151 self._paginator = self.paginator or NoPagination(parameters=parameters) 152 self._parameters = parameters 153 self._name = ( 154 InterpolatedString(self._name, parameters=parameters) 155 if isinstance(self._name, str) 156 else self._name 157 ) 158 159 @property # type: ignore 160 def name(self) -> str: 161 """ 162 :return: Stream name 163 """ 164 return ( 165 str(self._name.eval(self.config)) 166 if isinstance(self._name, InterpolatedString) 167 else self._name 168 ) 169 170 @name.setter 171 def name(self, value: str) -> None: 172 if not isinstance(value, property): 173 self._name = value 174 175 def _get_mapping( 176 self, method: Callable[..., Optional[Union[Mapping[str, Any], str]]], **kwargs: Any 177 ) -> Tuple[Union[Mapping[str, Any], str], Set[str]]: 178 """ 179 Get mapping from the provided method, and get the keys of the mapping. 180 If the method returns a string, it will return the string and an empty set. 181 If the method returns a dict, it will return the dict and its keys. 182 """ 183 mapping = method(**kwargs) or {} 184 keys = set(mapping.keys()) if not isinstance(mapping, str) else set() 185 return mapping, keys 186 187 def _get_request_options( 188 self, 189 stream_slice: Optional[StreamSlice], 190 next_page_token: Optional[Mapping[str, Any]], 191 paginator_method: Callable[..., Optional[Union[Mapping[str, Any], str]]], 192 stream_slicer_method: Callable[..., Optional[Union[Mapping[str, Any], str]]], 193 page_size_override: Optional[int] = None, 194 ) -> Union[Mapping[str, Any], str]: 195 """ 196 Get the request_option from the paginator and the stream slicer. 197 Raise a ValueError if there's a key collision 198 Returned merged mapping otherwise 199 """ 200 is_body_json = paginator_method.__name__ == "get_request_body_json" 201 202 mappings = [ 203 paginator_method( 204 stream_slice=stream_slice, 205 next_page_token=next_page_token, 206 **page_size_override_kwargs(page_size_override), 207 ), 208 ] 209 if not next_page_token or not self.ignore_stream_slicer_parameters_on_paginated_requests: 210 mappings.append( 211 stream_slicer_method( 212 stream_slice=stream_slice, 213 next_page_token=next_page_token, 214 ) 215 ) 216 return combine_mappings(mappings, allow_same_value_merge=is_body_json) 217 218 def _request_headers( 219 self, 220 stream_slice: Optional[StreamSlice] = None, 221 next_page_token: Optional[Mapping[str, Any]] = None, 222 page_size_override: Optional[int] = None, 223 ) -> Mapping[str, Any]: 224 """ 225 Specifies request headers. 226 Authentication headers will overwrite any overlapping headers returned from this method. 227 """ 228 headers = self._get_request_options( 229 stream_slice, 230 next_page_token, 231 self._paginator.get_request_headers, 232 self.request_option_provider.get_request_headers, 233 **page_size_override_kwargs(page_size_override), 234 ) 235 if isinstance(headers, str): 236 raise ValueError("Request headers cannot be a string") 237 return {str(k): str(v) for k, v in headers.items()} 238 239 def _request_params( 240 self, 241 stream_slice: Optional[StreamSlice] = None, 242 next_page_token: Optional[Mapping[str, Any]] = None, 243 page_size_override: Optional[int] = None, 244 ) -> Mapping[str, Any]: 245 """ 246 Specifies the query parameters that should be set on an outgoing HTTP request given the inputs. 247 248 E.g: you might want to define query parameters for paging if next_page_token is not None. 249 """ 250 params = self._get_request_options( 251 stream_slice, 252 next_page_token, 253 self._paginator.get_request_params, 254 self.request_option_provider.get_request_params, 255 **page_size_override_kwargs(page_size_override), 256 ) 257 if isinstance(params, str): 258 raise ValueError("Request params cannot be a string") 259 return params 260 261 def _request_body_data( 262 self, 263 stream_slice: Optional[StreamSlice] = None, 264 next_page_token: Optional[Mapping[str, Any]] = None, 265 page_size_override: Optional[int] = None, 266 ) -> Union[Mapping[str, Any], str]: 267 """ 268 Specifies how to populate the body of the request with a non-JSON payload. 269 270 If returns a ready text that it will be sent as is. 271 If returns a dict that it will be converted to a urlencoded form. 272 E.g. {"key1": "value1", "key2": "value2"} => "key1=value1&key2=value2" 273 274 At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden. 275 """ 276 return self._get_request_options( 277 stream_slice, 278 next_page_token, 279 self._paginator.get_request_body_data, 280 self.request_option_provider.get_request_body_data, 281 **page_size_override_kwargs(page_size_override), 282 ) 283 284 def _request_body_json( 285 self, 286 stream_slice: Optional[StreamSlice] = None, 287 next_page_token: Optional[Mapping[str, Any]] = None, 288 page_size_override: Optional[int] = None, 289 ) -> Optional[Mapping[str, Any]]: 290 """ 291 Specifies how to populate the body of the request with a JSON payload. 292 293 At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden. 294 """ 295 body_json = self._get_request_options( 296 stream_slice, 297 next_page_token, 298 self._paginator.get_request_body_json, 299 self.request_option_provider.get_request_body_json, 300 **page_size_override_kwargs(page_size_override), 301 ) 302 if isinstance(body_json, str): 303 raise ValueError("Request body json cannot be a string") 304 return body_json 305 306 def _paginator_path( 307 self, 308 next_page_token: Optional[Mapping[str, Any]] = None, 309 stream_slice: Optional[StreamSlice] = None, 310 ) -> Optional[str]: 311 """ 312 If the paginator points to a path, follow it, else return nothing so the requester is used. 313 :param next_page_token: 314 :return: 315 """ 316 return self._paginator.path( 317 next_page_token=next_page_token, 318 stream_state={}, # stream_state as an interpolation context is deprecated 319 stream_slice=stream_slice, 320 ) 321 322 def _parse_response( 323 self, 324 response: Optional[requests.Response], 325 records_schema: Mapping[str, Any], 326 stream_slice: Optional[StreamSlice] = None, 327 next_page_token: Optional[Mapping[str, Any]] = None, 328 ) -> Iterable[Record]: 329 if not response: 330 yield from [] 331 else: 332 yield from self.record_selector.select_records( 333 response=response, 334 stream_state={}, # stream_state as an interpolation context is deprecated 335 records_schema=records_schema, 336 stream_slice=stream_slice, 337 next_page_token=next_page_token, 338 ) 339 340 @property # type: ignore 341 def primary_key(self) -> Optional[Union[str, List[str], List[List[str]]]]: 342 """The stream's primary key""" 343 return self._primary_key 344 345 @primary_key.setter 346 def primary_key(self, value: str) -> None: 347 if not isinstance(value, property): 348 self._primary_key = value 349 350 def _next_page_token( 351 self, 352 response: requests.Response, 353 last_page_size: int, 354 last_record: Optional[Record], 355 last_page_token_value: Optional[Any], 356 page_size_override: Optional[int] = None, 357 stream_slice: Optional[StreamSlice] = None, 358 ) -> Optional[Mapping[str, Any]]: 359 """ 360 Specifies a pagination strategy. 361 362 The value returned from this method is passed to most other methods in this class. Use it to form a request e.g: set headers or query params. 363 364 :return: The token for the next page from the input response object. Returning None means there are no more pages to read in this response. 365 """ 366 return self._paginator.next_page_token( 367 response=response, 368 last_page_size=last_page_size, 369 last_record=last_record, 370 last_page_token_value=last_page_token_value, 371 **page_size_override_kwargs(page_size_override), 372 **stream_slice_kwargs(self._paginator.next_page_token, stream_slice), 373 ) 374 375 def _fetch_next_page( 376 self, 377 stream_slice: StreamSlice, 378 next_page_token: Optional[Mapping[str, Any]] = None, 379 page_size_override: Optional[int] = None, 380 ) -> Optional[requests.Response]: 381 return self.requester.send_request( 382 path=self._paginator_path( 383 next_page_token=next_page_token, 384 stream_slice=stream_slice, 385 ), 386 stream_state={}, # stream_state as an interpolation context is deprecated 387 stream_slice=stream_slice, 388 next_page_token=next_page_token, 389 request_headers=self._request_headers( 390 stream_slice=stream_slice, 391 next_page_token=next_page_token, 392 **page_size_override_kwargs(page_size_override), 393 ), 394 request_params=self._request_params( 395 stream_slice=stream_slice, 396 next_page_token=next_page_token, 397 **page_size_override_kwargs(page_size_override), 398 ), 399 request_body_data=self._request_body_data( 400 stream_slice=stream_slice, 401 next_page_token=next_page_token, 402 **page_size_override_kwargs(page_size_override), 403 ), 404 request_body_json=self._request_body_json( 405 stream_slice=stream_slice, 406 next_page_token=next_page_token, 407 **page_size_override_kwargs(page_size_override), 408 ), 409 log_formatter=self.log_formatter, 410 ) 411 412 # This logic is similar to _read_pages in the HttpStream class. When making changes here, consider making changes there as well. 413 def _read_pages( 414 self, 415 records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]], 416 stream_slice: StreamSlice, 417 ) -> Iterable[Record]: 418 original_stream_slice = stream_slice 419 pagination_tracker = self.pagination_tracker_factory() 420 page_size_reducer = ( 421 PageSizeReducer( 422 self.page_size_reduction, 423 self._paginator.get_page_size(), 424 stream_name=self.name, 425 ) 426 if self.page_size_reduction 427 else None 428 ) 429 reset_pagination = False 430 reduce_page_size = False 431 next_page_token = self._get_initial_next_page_token() 432 while True: 433 page_size_override = page_size_reducer.page_size_override if page_size_reducer else None 434 merged_records: MutableMapping[str, Any] = defaultdict(dict) 435 last_page_size = 0 436 last_record: Optional[Record] = None 437 438 response = None 439 try: 440 if self.additional_query_properties: 441 for ( 442 properties 443 ) in self.additional_query_properties.get_request_property_chunks(): 444 stream_slice = StreamSlice( 445 partition=stream_slice.partition or {}, 446 cursor_slice=stream_slice.cursor_slice or {}, 447 extra_fields={"query_properties": properties}, 448 ) 449 response = self._fetch_next_page( 450 stream_slice, 451 next_page_token, 452 **page_size_override_kwargs(page_size_override), 453 ) 454 455 for current_record in records_generator_fn(response): 456 if self.additional_query_properties.property_chunking: 457 merge_key = self.additional_query_properties.property_chunking.get_merge_key( 458 current_record 459 ) 460 if merge_key: 461 _deep_merge(merged_records[merge_key], current_record) 462 else: 463 # We should still emit records even if the record did not have a merge key 464 pagination_tracker.observe(current_record) 465 last_page_size += 1 466 last_record = current_record 467 yield current_record 468 else: 469 pagination_tracker.observe(current_record) 470 last_page_size += 1 471 last_record = current_record 472 yield current_record 473 474 for merged_record in merged_records.values(): 475 record = Record( 476 data=merged_record, stream_name=self.name, associated_slice=stream_slice 477 ) 478 pagination_tracker.observe(record) 479 last_page_size += 1 480 last_record = record 481 yield record 482 else: 483 response = self._fetch_next_page( 484 stream_slice, 485 next_page_token, 486 **page_size_override_kwargs(page_size_override), 487 ) 488 for current_record in records_generator_fn(response): 489 pagination_tracker.observe(current_record) 490 last_page_size += 1 491 last_record = current_record 492 yield current_record 493 except PaginationResetRequiredException: 494 reset_pagination = True 495 except PageSizeReductionRequiredException: 496 if page_size_reducer is None: 497 # The action can be attached to a requester we cannot validate at config time, such as one 498 # built by a custom error handler or a CustomRequester. Re-raise as the misconfiguration it 499 # is: the exception being handled is the neutral "the API asked for a smaller page" signal. 500 raise PageSizeReductionNotSupportedException(stream_name=self.name) 501 if last_page_size: 502 # The reduction is safe only because it re-issues a page whose records were not emitted. 503 # Every in-CDK way of reaching this raises from `_fetch_next_page`, before the record loop, 504 # and the factory rejects the manifest constructs that would not, but a custom extractor, 505 # filter or transformation can issue its own request from inside the record generator. 506 # Re-issuing the page then duplicates the records already yielded, so fail instead. 507 raise AirbyteTracedException( 508 internal_message=f"Stream {self.name} requested a page size reduction after {last_page_size} records of the page had already been emitted", 509 message=f"Stream {self.name} asked for a smaller page size in the middle of a page. The page cannot be requested again without duplicating the records already read from it. Move the REDUCE_PAGE_SIZE action to the error handler of the stream's main requester.", 510 failure_type=FailureType.config_error, 511 ) 512 # Raises once the page size cannot be reduced any further and the retries allowed at that 513 # floor are spent, which is what stops the loop when the API keeps failing. 514 page_size_reducer.reduce() 515 reduce_page_size = True 516 else: 517 if page_size_reducer: 518 page_size_reducer.on_successful_page() 519 if not response: 520 break 521 522 if reduce_page_size: 523 # Retry the very same page: neither the token nor the slice change, only the page size does - 524 # and not even that once the reducer is at its floor and only waiting is left. 525 reduce_page_size = False 526 continue 527 528 if reset_pagination or pagination_tracker.has_reached_limit(): 529 next_page_token = self._get_initial_next_page_token() 530 previous_slice = stream_slice 531 stream_slice = pagination_tracker.reduce_slice_range_if_possible( 532 stream_slice, original_stream_slice 533 ) 534 LOGGER.info( 535 f"Hitting PaginationReset event. StreamSlice used will go from {previous_slice} to {stream_slice}" 536 ) 537 reset_pagination = False 538 else: 539 last_page_token_value = ( 540 next_page_token.get("next_page_token") if next_page_token else None 541 ) 542 next_page_token = self._next_page_token( 543 response=response, # type:ignore # we are breaking from the loop on the try/else if there are no response so this should be fine 544 last_page_size=last_page_size, 545 last_record=last_record, 546 last_page_token_value=last_page_token_value, 547 **page_size_override_kwargs(page_size_override), 548 **stream_slice_kwargs(self._next_page_token, stream_slice), 549 ) 550 if not next_page_token: 551 break 552 553 # Always return an empty generator just in case no records were ever yielded 554 yield from [] 555 556 def _get_initial_next_page_token(self) -> Optional[Mapping[str, Any]]: 557 initial_token = self._paginator.get_initial_token() 558 next_page_token = {"next_page_token": initial_token} if initial_token is not None else None 559 return next_page_token 560 561 def read_records( 562 self, 563 records_schema: Mapping[str, Any], 564 stream_slice: Optional[StreamSlice] = None, 565 ) -> Iterable[StreamData]: 566 """ 567 Fetch a stream's records from an HTTP API source 568 569 :param records_schema: json schema to describe record 570 :param stream_slice: The stream slice to read data for 571 :return: The records read from the API source 572 """ 573 _slice = stream_slice or StreamSlice(partition={}, cursor_slice={}) # None-check 574 yield from self._read_records_or_split_request_window( 575 records_schema, _slice, _slice, depth=0 576 ) 577 578 def _read_records_or_split_request_window( 579 self, 580 records_schema: Mapping[str, Any], 581 stream_slice: StreamSlice, 582 original_slice: StreamSlice, 583 depth: int, 584 ) -> Iterable[StreamData]: 585 """ 586 Read `stream_slice` to completion, replacing it with smaller children and recursing into each of them in 587 turn when a `RequestWindowSplitRequiredException` is raised while reading it. 588 589 Each recursive call re-enters this method - and, through it, `_read_pages` - from scratch for the child 590 slice it is given, so every child gets its own paginator token, `PaginationTracker`, and 591 `PageSizeReducer` for free: nothing here needs to reset that state explicitly. Because the whole 592 recursion lives inside the single generator `DeclarativePartition.read()` consumes, 593 `PartitionReader.process_partition()` only calls `cursor.close_partition()` once, after this generator is 594 fully exhausted - so a failure anywhere in the recursion (a child that cannot be read, or a window that 595 cannot be split any further) propagates out without ever checkpointing the original partition. 596 597 `original_slice` is threaded through the recursion unchanged so every record yielded, however deep the 598 recursion went to produce it, is re-stamped with the partition's own slice rather than the child slice it 599 was actually read against - see `_reassociate_with_original_slice`. `depth` bounds the recursion 600 independently of `request_window_splitter`'s own no-progress guard: it is enforced here so a 601 misbehaving custom cursor cannot recurse indefinitely regardless of what that guard does or does not 602 catch. 603 """ 604 record_generator = partial( 605 self._parse_records, 606 stream_slice=stream_slice, 607 records_schema=records_schema, 608 ) 609 610 emitted_count = 0 611 try: 612 records: Iterable[Mapping[str, Any]] = self._read_pages(record_generator, stream_slice) 613 if self.post_pagination_filter: 614 # A data feed paginates until it reaches a record older than the cursor, so the page that triggers the stop 615 # condition still holds already-synced records. Those are filtered here rather than in the record selector 616 # so that the paginator keeps seeing the whole page: the stop condition is evaluated on the last record of 617 # the page, which is precisely one of the records being dropped. Two consequences of filtering this late: 618 # the pagination tracker observes the dropped records, and a `file_uploader` on the record selector has 619 # already uploaded their files by the time they are dropped. 620 records = self.post_pagination_filter.filter_records( 621 records, 622 # the filter is only used for its cursor comparison, which does not read the stream state 623 stream_state={}, 624 stream_slice=stream_slice, 625 ) 626 for record in records: 627 emitted_count += 1 628 yield self._reassociate_with_original_slice(record, original_slice) 629 except RequestWindowSplitRequiredException as exception: 630 if self.request_window_splitting is None or self.request_window_splitter is None: 631 raise RequestWindowSplitNotSupportedException(stream_name=self.name) from exception 632 633 if emitted_count: 634 # Splitting and re-reading the window re-emits these records: a failed partition is never 635 # checkpointed, so the next attempt would re-emit them anyway. Always replaying makes forward 636 # progress instead of retrying the same oversized window forever; primary-key dedup downstream 637 # handles the duplicates either way. 638 LOGGER.warning( 639 f"Stream {self.name} already emitted {emitted_count} record(s) from {stream_slice} before " 640 f"it was rejected; splitting and re-reading the window may re-emit them." 641 ) 642 643 if depth >= _MAX_REQUEST_WINDOW_SPLIT_DEPTH: 644 # Logged separately from the exception below so it stays visible even if the trace message's 645 # internal_message isn't surfaced by whatever catches it. 646 LOGGER.warning( 647 f"Stream {self.name} hit the maximum request window split depth " 648 f"({_MAX_REQUEST_WINDOW_SPLIT_DEPTH}) while splitting {original_slice}." 649 ) 650 raise AirbyteTracedException( 651 internal_message=f"Stream {self.name} exceeded the maximum request window split depth of {_MAX_REQUEST_WINDOW_SPLIT_DEPTH} while splitting {original_slice}", 652 # `transient_error`, so the only remediation is the connector's own, if it defined one. 653 message=self._with_request_window_failure_message( 654 f"Stream {self.name} could not split its request window to a size the API accepts " 655 f"within {_MAX_REQUEST_WINDOW_SPLIT_DEPTH} splits. This usually means the stream's " 656 f"cursor is not actually shrinking the window on each split; if it uses a custom " 657 f"cursor, check its `split_request_window` implementation." 658 ), 659 failure_type=FailureType.transient_error, 660 ) from exception 661 662 children = self.request_window_splitter( 663 stream_slice, self.request_window_splitting.min_split_window 664 ) 665 if children is None: 666 min_split_window_note = ( 667 f", its configured `min_split_window` ({self.request_window_splitting.min_split_window})" 668 if self.request_window_splitting.min_split_window 669 else "" 670 ) 671 raise AirbyteTracedException( 672 internal_message=f"Stream {self.name} could not split its request window {stream_slice} any further", 673 # `transient_error` regardless of how the triggering response was classified: exhaustion is 674 # never the user's fault, matching `PageSizeReducer`'s equivalent exhaustion branch. 675 message=self._with_request_window_failure_message( 676 f"The API kept rejecting stream {self.name}'s request window even at the smallest window " 677 f"its cursor allows{min_split_window_note}, or splitting the window further would not " 678 f"make progress." 679 ), 680 failure_type=FailureType.transient_error, 681 ) from exception 682 683 LOGGER.info( 684 f"Stream {self.name}: the API rejected request window {stream_slice} (split depth {depth}); " 685 f"reducing it to {children} and reading each in turn." 686 ) 687 for child in children: 688 yield from self._read_records_or_split_request_window( 689 records_schema, child, original_slice, depth + 1 690 ) 691 692 def _with_request_window_failure_message(self, message: str) -> str: 693 """ 694 :return: the message followed by `request_window_splitting`'s configured `failure_message`, when it 695 defined one - mirroring `PageSizeReducer._with_failure_message` 696 """ 697 assert self.request_window_splitting is not None 698 failure_message = (self.request_window_splitting.failure_message or "").strip() 699 if not failure_message: 700 return message 701 return f"{message} {failure_message}" 702 703 @staticmethod 704 def _reassociate_with_original_slice( 705 record: StreamData, original_slice: StreamSlice 706 ) -> StreamData: 707 """ 708 Records read while recursing into a split child window are parsed against that child's `StreamSlice`, 709 so `RecordSelector` stamps them with the child as `associated_slice`. `DeclarativePartition.read()` 710 passes an already-built `Record` through unchanged rather than re-wrapping it, so left uncorrected the 711 child slice - not the partition's own slice - is what `ConcurrentCursor.observe()` would key its 712 per-partition bookkeeping by. `close_partition()` always looks that bookkeeping up by the partition's own 713 slice, so it would never find it, silently losing the precise most-recently-observed cursor value for a 714 split partition (recoverable in practice today only because callers of that lookup fall back to the 715 slice's end boundary when it is missing). Re-stamping every record with `original_slice` here restores 716 the same tracking a non-split read already gets, rather than relying on that fallback. 717 """ 718 if isinstance(record, Record) and record.associated_slice is not original_slice: 719 return Record( 720 data=record.data, 721 stream_name=record.stream_name, 722 associated_slice=original_slice, 723 file_reference=record.file_reference, 724 ) 725 return record 726 727 def _parse_records( 728 self, 729 response: Optional[requests.Response], 730 records_schema: Mapping[str, Any], 731 stream_slice: Optional[StreamSlice], 732 ) -> Iterable[Record]: 733 yield from self._parse_response( 734 response, 735 stream_slice=stream_slice, 736 records_schema=records_schema, 737 ) 738 739 def must_deduplicate_query_params(self) -> bool: 740 return True 741 742 @staticmethod 743 def _to_partition_key(to_serialize: Any) -> str: 744 # separators have changed in Python 3.4. To avoid being impacted by further change, we explicitly specify our own value 745 return json.dumps(to_serialize, indent=None, separators=(",", ":"), sort_keys=True) 746 747 748def _deep_merge( 749 target: MutableMapping[str, Any], source: Union[Record, MutableMapping[str, Any]] 750) -> None: 751 """ 752 Recursively merge two dictionaries, combining nested dictionaries instead of overwriting them. 753 754 :param target: The dictionary to merge into (modified in place) 755 :param source: The dictionary to merge from 756 """ 757 for key, value in source.items(): 758 if ( 759 key in target 760 and isinstance(target[key], MutableMapping) 761 and isinstance(value, MutableMapping) 762 ): 763 _deep_merge(target[key], value) 764 else: 765 target[key] = value 766 767 768@deprecated( 769 "This class is experimental. Use at your own risk.", 770 category=ExperimentalClassWarning, 771) 772@dataclass 773class LazySimpleRetriever(SimpleRetriever): 774 """ 775 A retriever that supports lazy loading from parent streams. 776 """ 777 778 def _read_pages( 779 self, 780 records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]], 781 stream_slice: StreamSlice, 782 ) -> Iterable[Record]: 783 response = stream_slice.extra_fields["child_response"] 784 if response: 785 last_page_size, last_record = 0, None 786 for record in records_generator_fn(response): # type: ignore[call-arg] # only _parse_records expected as a func 787 last_page_size += 1 788 last_record = record 789 yield record 790 791 next_page_token = self._next_page_token( 792 response, 793 last_page_size, 794 last_record, 795 None, 796 **stream_slice_kwargs(self._next_page_token, stream_slice), 797 ) 798 if next_page_token: 799 yield from self._paginate( 800 next_page_token, 801 records_generator_fn, 802 stream_slice, 803 ) 804 805 yield from [] 806 else: 807 # coderabbit detected an interesting bug/gap where if we were to not get a child_response, we 808 # might recurse forever. This might not be the case, but it is worth noting that this code path 809 # isn't comprehensively tested. 810 yield from self._read_pages(records_generator_fn, stream_slice) 811 812 def _paginate( 813 self, 814 next_page_token: Any, 815 records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]], 816 stream_slice: StreamSlice, 817 ) -> Iterable[Record]: 818 """Handle pagination by fetching subsequent pages.""" 819 pagination_complete = False 820 821 while not pagination_complete: 822 response = self._fetch_next_page(stream_slice, next_page_token) 823 last_page_size, last_record = 0, None 824 825 for record in records_generator_fn(response): # type: ignore[call-arg] # only _parse_records expected as a func 826 last_page_size += 1 827 last_record = record 828 yield record 829 830 if not response: 831 pagination_complete = True 832 else: 833 last_page_token_value = ( 834 next_page_token.get("next_page_token") if next_page_token else None 835 ) 836 next_page_token = self._next_page_token( 837 response, 838 last_page_size, 839 last_record, 840 last_page_token_value, 841 **stream_slice_kwargs(self._next_page_token, stream_slice), 842 ) 843 844 if not next_page_token: 845 pagination_complete = True
FULL_REFRESH_SYNC_COMPLETE_KEY =
'__ab_full_refresh_sync_complete'
LOGGER =
<Logger airbyte (INFO)>
84@dataclass 85class SimpleRetriever(Retriever): 86 """ 87 Retrieves records by synchronously sending requests to fetch records. 88 89 The retriever acts as an orchestrator between the requester, the record selector, the paginator, and the stream slicer. 90 91 For each stream slice, submit requests until there are no more pages of records to fetch. 92 93 This retriever currently inherits from HttpStream to reuse the request submission and pagination machinery. 94 As a result, some of the parameters passed to some methods are unused. 95 The two will be decoupled in a future release. 96 97 Attributes: 98 stream_name (str): The stream's name 99 stream_primary_key (Optional[Union[str, List[str], List[List[str]]]]): The stream's primary key 100 requester (Requester): The HTTP requester 101 record_selector (HttpSelector): The record selector 102 paginator (Optional[Paginator]): The paginator 103 stream_slicer (Optional[StreamSlicer]): The stream slicer 104 parameters (Mapping[str, Any]): Additional runtime parameters to be used for string interpolation 105 post_pagination_filter (Optional[ClientSideIncrementalRecordFilterDecorator]): Set for data feed streams only. 106 Records the cursor considers already synced are dropped once pagination has observed them 107 page_size_reduction (Optional[PageSizeReduction]): How much to shrink the page size when an error handler 108 resolves to `ResponseAction.REDUCE_PAGE_SIZE`. `None` disables page size reduction entirely. 109 It is immutable configuration; the page size in effect lives in a `PageSizeReducer` that 110 `_read_pages` creates per call, so the retriever and its paginator - both shared by every 111 partition of the stream, read concurrently - stay stateless. When `page_size_reduction` is 112 `None` no reducer is created and `_read_pages` keeps its previous behaviour 113 request_window_splitting (Optional[RequestWindowSplitting]): Policy applied when an error handler 114 resolves to `ResponseAction.SPLIT_REQUEST_WINDOW`, or when `RequestWindowSplitRequiredException` 115 is raised directly by custom code. `None` disables request-window splitting entirely: `read_records` 116 converts the exception into `RequestWindowSplitNotSupportedException` instead of attempting a 117 split it cannot perform. 118 request_window_splitter: Bound method - normally the stream's own cursor's `split_request_window` - 119 asked to replace a failing `StreamSlice` with smaller children. `None` has the same effect as 120 `request_window_splitting` being `None`. 121 """ 122 123 requester: Requester 124 record_selector: HttpSelector 125 config: Config 126 parameters: InitVar[Mapping[str, Any]] 127 name: str 128 _name: Union[InterpolatedString, str] = field(init=False, repr=False, default="") 129 primary_key: Optional[Union[str, List[str], List[List[str]]]] 130 _primary_key: str = field(init=False, repr=False, default="") 131 paginator: Optional[Paginator] = None 132 stream_slicer: StreamSlicer = field( 133 default_factory=lambda: SinglePartitionRouter(parameters={}) 134 ) 135 request_option_provider: RequestOptionsProvider = field( 136 default_factory=lambda: DefaultRequestOptionsProvider(parameters={}) 137 ) 138 ignore_stream_slicer_parameters_on_paginated_requests: bool = False 139 additional_query_properties: Optional[QueryProperties] = None 140 log_formatter: Optional[Callable[[requests.Response], Any]] = None 141 pagination_tracker_factory: Callable[[], PaginationTracker] = field( 142 default_factory=lambda: lambda: PaginationTracker() 143 ) 144 post_pagination_filter: Optional[ClientSideIncrementalRecordFilterDecorator] = None 145 page_size_reduction: Optional[PageSizeReduction] = None 146 request_window_splitting: Optional[RequestWindowSplitting] = None 147 request_window_splitter: Optional[ 148 Callable[[StreamSlice, Optional[datetime.timedelta]], Optional[List[StreamSlice]]] 149 ] = None 150 151 def __post_init__(self, parameters: Mapping[str, Any]) -> None: 152 self._paginator = self.paginator or NoPagination(parameters=parameters) 153 self._parameters = parameters 154 self._name = ( 155 InterpolatedString(self._name, parameters=parameters) 156 if isinstance(self._name, str) 157 else self._name 158 ) 159 160 @property # type: ignore 161 def name(self) -> str: 162 """ 163 :return: Stream name 164 """ 165 return ( 166 str(self._name.eval(self.config)) 167 if isinstance(self._name, InterpolatedString) 168 else self._name 169 ) 170 171 @name.setter 172 def name(self, value: str) -> None: 173 if not isinstance(value, property): 174 self._name = value 175 176 def _get_mapping( 177 self, method: Callable[..., Optional[Union[Mapping[str, Any], str]]], **kwargs: Any 178 ) -> Tuple[Union[Mapping[str, Any], str], Set[str]]: 179 """ 180 Get mapping from the provided method, and get the keys of the mapping. 181 If the method returns a string, it will return the string and an empty set. 182 If the method returns a dict, it will return the dict and its keys. 183 """ 184 mapping = method(**kwargs) or {} 185 keys = set(mapping.keys()) if not isinstance(mapping, str) else set() 186 return mapping, keys 187 188 def _get_request_options( 189 self, 190 stream_slice: Optional[StreamSlice], 191 next_page_token: Optional[Mapping[str, Any]], 192 paginator_method: Callable[..., Optional[Union[Mapping[str, Any], str]]], 193 stream_slicer_method: Callable[..., Optional[Union[Mapping[str, Any], str]]], 194 page_size_override: Optional[int] = None, 195 ) -> Union[Mapping[str, Any], str]: 196 """ 197 Get the request_option from the paginator and the stream slicer. 198 Raise a ValueError if there's a key collision 199 Returned merged mapping otherwise 200 """ 201 is_body_json = paginator_method.__name__ == "get_request_body_json" 202 203 mappings = [ 204 paginator_method( 205 stream_slice=stream_slice, 206 next_page_token=next_page_token, 207 **page_size_override_kwargs(page_size_override), 208 ), 209 ] 210 if not next_page_token or not self.ignore_stream_slicer_parameters_on_paginated_requests: 211 mappings.append( 212 stream_slicer_method( 213 stream_slice=stream_slice, 214 next_page_token=next_page_token, 215 ) 216 ) 217 return combine_mappings(mappings, allow_same_value_merge=is_body_json) 218 219 def _request_headers( 220 self, 221 stream_slice: Optional[StreamSlice] = None, 222 next_page_token: Optional[Mapping[str, Any]] = None, 223 page_size_override: Optional[int] = None, 224 ) -> Mapping[str, Any]: 225 """ 226 Specifies request headers. 227 Authentication headers will overwrite any overlapping headers returned from this method. 228 """ 229 headers = self._get_request_options( 230 stream_slice, 231 next_page_token, 232 self._paginator.get_request_headers, 233 self.request_option_provider.get_request_headers, 234 **page_size_override_kwargs(page_size_override), 235 ) 236 if isinstance(headers, str): 237 raise ValueError("Request headers cannot be a string") 238 return {str(k): str(v) for k, v in headers.items()} 239 240 def _request_params( 241 self, 242 stream_slice: Optional[StreamSlice] = None, 243 next_page_token: Optional[Mapping[str, Any]] = None, 244 page_size_override: Optional[int] = None, 245 ) -> Mapping[str, Any]: 246 """ 247 Specifies the query parameters that should be set on an outgoing HTTP request given the inputs. 248 249 E.g: you might want to define query parameters for paging if next_page_token is not None. 250 """ 251 params = self._get_request_options( 252 stream_slice, 253 next_page_token, 254 self._paginator.get_request_params, 255 self.request_option_provider.get_request_params, 256 **page_size_override_kwargs(page_size_override), 257 ) 258 if isinstance(params, str): 259 raise ValueError("Request params cannot be a string") 260 return params 261 262 def _request_body_data( 263 self, 264 stream_slice: Optional[StreamSlice] = None, 265 next_page_token: Optional[Mapping[str, Any]] = None, 266 page_size_override: Optional[int] = None, 267 ) -> Union[Mapping[str, Any], str]: 268 """ 269 Specifies how to populate the body of the request with a non-JSON payload. 270 271 If returns a ready text that it will be sent as is. 272 If returns a dict that it will be converted to a urlencoded form. 273 E.g. {"key1": "value1", "key2": "value2"} => "key1=value1&key2=value2" 274 275 At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden. 276 """ 277 return self._get_request_options( 278 stream_slice, 279 next_page_token, 280 self._paginator.get_request_body_data, 281 self.request_option_provider.get_request_body_data, 282 **page_size_override_kwargs(page_size_override), 283 ) 284 285 def _request_body_json( 286 self, 287 stream_slice: Optional[StreamSlice] = None, 288 next_page_token: Optional[Mapping[str, Any]] = None, 289 page_size_override: Optional[int] = None, 290 ) -> Optional[Mapping[str, Any]]: 291 """ 292 Specifies how to populate the body of the request with a JSON payload. 293 294 At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden. 295 """ 296 body_json = self._get_request_options( 297 stream_slice, 298 next_page_token, 299 self._paginator.get_request_body_json, 300 self.request_option_provider.get_request_body_json, 301 **page_size_override_kwargs(page_size_override), 302 ) 303 if isinstance(body_json, str): 304 raise ValueError("Request body json cannot be a string") 305 return body_json 306 307 def _paginator_path( 308 self, 309 next_page_token: Optional[Mapping[str, Any]] = None, 310 stream_slice: Optional[StreamSlice] = None, 311 ) -> Optional[str]: 312 """ 313 If the paginator points to a path, follow it, else return nothing so the requester is used. 314 :param next_page_token: 315 :return: 316 """ 317 return self._paginator.path( 318 next_page_token=next_page_token, 319 stream_state={}, # stream_state as an interpolation context is deprecated 320 stream_slice=stream_slice, 321 ) 322 323 def _parse_response( 324 self, 325 response: Optional[requests.Response], 326 records_schema: Mapping[str, Any], 327 stream_slice: Optional[StreamSlice] = None, 328 next_page_token: Optional[Mapping[str, Any]] = None, 329 ) -> Iterable[Record]: 330 if not response: 331 yield from [] 332 else: 333 yield from self.record_selector.select_records( 334 response=response, 335 stream_state={}, # stream_state as an interpolation context is deprecated 336 records_schema=records_schema, 337 stream_slice=stream_slice, 338 next_page_token=next_page_token, 339 ) 340 341 @property # type: ignore 342 def primary_key(self) -> Optional[Union[str, List[str], List[List[str]]]]: 343 """The stream's primary key""" 344 return self._primary_key 345 346 @primary_key.setter 347 def primary_key(self, value: str) -> None: 348 if not isinstance(value, property): 349 self._primary_key = value 350 351 def _next_page_token( 352 self, 353 response: requests.Response, 354 last_page_size: int, 355 last_record: Optional[Record], 356 last_page_token_value: Optional[Any], 357 page_size_override: Optional[int] = None, 358 stream_slice: Optional[StreamSlice] = None, 359 ) -> Optional[Mapping[str, Any]]: 360 """ 361 Specifies a pagination strategy. 362 363 The value returned from this method is passed to most other methods in this class. Use it to form a request e.g: set headers or query params. 364 365 :return: The token for the next page from the input response object. Returning None means there are no more pages to read in this response. 366 """ 367 return self._paginator.next_page_token( 368 response=response, 369 last_page_size=last_page_size, 370 last_record=last_record, 371 last_page_token_value=last_page_token_value, 372 **page_size_override_kwargs(page_size_override), 373 **stream_slice_kwargs(self._paginator.next_page_token, stream_slice), 374 ) 375 376 def _fetch_next_page( 377 self, 378 stream_slice: StreamSlice, 379 next_page_token: Optional[Mapping[str, Any]] = None, 380 page_size_override: Optional[int] = None, 381 ) -> Optional[requests.Response]: 382 return self.requester.send_request( 383 path=self._paginator_path( 384 next_page_token=next_page_token, 385 stream_slice=stream_slice, 386 ), 387 stream_state={}, # stream_state as an interpolation context is deprecated 388 stream_slice=stream_slice, 389 next_page_token=next_page_token, 390 request_headers=self._request_headers( 391 stream_slice=stream_slice, 392 next_page_token=next_page_token, 393 **page_size_override_kwargs(page_size_override), 394 ), 395 request_params=self._request_params( 396 stream_slice=stream_slice, 397 next_page_token=next_page_token, 398 **page_size_override_kwargs(page_size_override), 399 ), 400 request_body_data=self._request_body_data( 401 stream_slice=stream_slice, 402 next_page_token=next_page_token, 403 **page_size_override_kwargs(page_size_override), 404 ), 405 request_body_json=self._request_body_json( 406 stream_slice=stream_slice, 407 next_page_token=next_page_token, 408 **page_size_override_kwargs(page_size_override), 409 ), 410 log_formatter=self.log_formatter, 411 ) 412 413 # This logic is similar to _read_pages in the HttpStream class. When making changes here, consider making changes there as well. 414 def _read_pages( 415 self, 416 records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]], 417 stream_slice: StreamSlice, 418 ) -> Iterable[Record]: 419 original_stream_slice = stream_slice 420 pagination_tracker = self.pagination_tracker_factory() 421 page_size_reducer = ( 422 PageSizeReducer( 423 self.page_size_reduction, 424 self._paginator.get_page_size(), 425 stream_name=self.name, 426 ) 427 if self.page_size_reduction 428 else None 429 ) 430 reset_pagination = False 431 reduce_page_size = False 432 next_page_token = self._get_initial_next_page_token() 433 while True: 434 page_size_override = page_size_reducer.page_size_override if page_size_reducer else None 435 merged_records: MutableMapping[str, Any] = defaultdict(dict) 436 last_page_size = 0 437 last_record: Optional[Record] = None 438 439 response = None 440 try: 441 if self.additional_query_properties: 442 for ( 443 properties 444 ) in self.additional_query_properties.get_request_property_chunks(): 445 stream_slice = StreamSlice( 446 partition=stream_slice.partition or {}, 447 cursor_slice=stream_slice.cursor_slice or {}, 448 extra_fields={"query_properties": properties}, 449 ) 450 response = self._fetch_next_page( 451 stream_slice, 452 next_page_token, 453 **page_size_override_kwargs(page_size_override), 454 ) 455 456 for current_record in records_generator_fn(response): 457 if self.additional_query_properties.property_chunking: 458 merge_key = self.additional_query_properties.property_chunking.get_merge_key( 459 current_record 460 ) 461 if merge_key: 462 _deep_merge(merged_records[merge_key], current_record) 463 else: 464 # We should still emit records even if the record did not have a merge key 465 pagination_tracker.observe(current_record) 466 last_page_size += 1 467 last_record = current_record 468 yield current_record 469 else: 470 pagination_tracker.observe(current_record) 471 last_page_size += 1 472 last_record = current_record 473 yield current_record 474 475 for merged_record in merged_records.values(): 476 record = Record( 477 data=merged_record, stream_name=self.name, associated_slice=stream_slice 478 ) 479 pagination_tracker.observe(record) 480 last_page_size += 1 481 last_record = record 482 yield record 483 else: 484 response = self._fetch_next_page( 485 stream_slice, 486 next_page_token, 487 **page_size_override_kwargs(page_size_override), 488 ) 489 for current_record in records_generator_fn(response): 490 pagination_tracker.observe(current_record) 491 last_page_size += 1 492 last_record = current_record 493 yield current_record 494 except PaginationResetRequiredException: 495 reset_pagination = True 496 except PageSizeReductionRequiredException: 497 if page_size_reducer is None: 498 # The action can be attached to a requester we cannot validate at config time, such as one 499 # built by a custom error handler or a CustomRequester. Re-raise as the misconfiguration it 500 # is: the exception being handled is the neutral "the API asked for a smaller page" signal. 501 raise PageSizeReductionNotSupportedException(stream_name=self.name) 502 if last_page_size: 503 # The reduction is safe only because it re-issues a page whose records were not emitted. 504 # Every in-CDK way of reaching this raises from `_fetch_next_page`, before the record loop, 505 # and the factory rejects the manifest constructs that would not, but a custom extractor, 506 # filter or transformation can issue its own request from inside the record generator. 507 # Re-issuing the page then duplicates the records already yielded, so fail instead. 508 raise AirbyteTracedException( 509 internal_message=f"Stream {self.name} requested a page size reduction after {last_page_size} records of the page had already been emitted", 510 message=f"Stream {self.name} asked for a smaller page size in the middle of a page. The page cannot be requested again without duplicating the records already read from it. Move the REDUCE_PAGE_SIZE action to the error handler of the stream's main requester.", 511 failure_type=FailureType.config_error, 512 ) 513 # Raises once the page size cannot be reduced any further and the retries allowed at that 514 # floor are spent, which is what stops the loop when the API keeps failing. 515 page_size_reducer.reduce() 516 reduce_page_size = True 517 else: 518 if page_size_reducer: 519 page_size_reducer.on_successful_page() 520 if not response: 521 break 522 523 if reduce_page_size: 524 # Retry the very same page: neither the token nor the slice change, only the page size does - 525 # and not even that once the reducer is at its floor and only waiting is left. 526 reduce_page_size = False 527 continue 528 529 if reset_pagination or pagination_tracker.has_reached_limit(): 530 next_page_token = self._get_initial_next_page_token() 531 previous_slice = stream_slice 532 stream_slice = pagination_tracker.reduce_slice_range_if_possible( 533 stream_slice, original_stream_slice 534 ) 535 LOGGER.info( 536 f"Hitting PaginationReset event. StreamSlice used will go from {previous_slice} to {stream_slice}" 537 ) 538 reset_pagination = False 539 else: 540 last_page_token_value = ( 541 next_page_token.get("next_page_token") if next_page_token else None 542 ) 543 next_page_token = self._next_page_token( 544 response=response, # type:ignore # we are breaking from the loop on the try/else if there are no response so this should be fine 545 last_page_size=last_page_size, 546 last_record=last_record, 547 last_page_token_value=last_page_token_value, 548 **page_size_override_kwargs(page_size_override), 549 **stream_slice_kwargs(self._next_page_token, stream_slice), 550 ) 551 if not next_page_token: 552 break 553 554 # Always return an empty generator just in case no records were ever yielded 555 yield from [] 556 557 def _get_initial_next_page_token(self) -> Optional[Mapping[str, Any]]: 558 initial_token = self._paginator.get_initial_token() 559 next_page_token = {"next_page_token": initial_token} if initial_token is not None else None 560 return next_page_token 561 562 def read_records( 563 self, 564 records_schema: Mapping[str, Any], 565 stream_slice: Optional[StreamSlice] = None, 566 ) -> Iterable[StreamData]: 567 """ 568 Fetch a stream's records from an HTTP API source 569 570 :param records_schema: json schema to describe record 571 :param stream_slice: The stream slice to read data for 572 :return: The records read from the API source 573 """ 574 _slice = stream_slice or StreamSlice(partition={}, cursor_slice={}) # None-check 575 yield from self._read_records_or_split_request_window( 576 records_schema, _slice, _slice, depth=0 577 ) 578 579 def _read_records_or_split_request_window( 580 self, 581 records_schema: Mapping[str, Any], 582 stream_slice: StreamSlice, 583 original_slice: StreamSlice, 584 depth: int, 585 ) -> Iterable[StreamData]: 586 """ 587 Read `stream_slice` to completion, replacing it with smaller children and recursing into each of them in 588 turn when a `RequestWindowSplitRequiredException` is raised while reading it. 589 590 Each recursive call re-enters this method - and, through it, `_read_pages` - from scratch for the child 591 slice it is given, so every child gets its own paginator token, `PaginationTracker`, and 592 `PageSizeReducer` for free: nothing here needs to reset that state explicitly. Because the whole 593 recursion lives inside the single generator `DeclarativePartition.read()` consumes, 594 `PartitionReader.process_partition()` only calls `cursor.close_partition()` once, after this generator is 595 fully exhausted - so a failure anywhere in the recursion (a child that cannot be read, or a window that 596 cannot be split any further) propagates out without ever checkpointing the original partition. 597 598 `original_slice` is threaded through the recursion unchanged so every record yielded, however deep the 599 recursion went to produce it, is re-stamped with the partition's own slice rather than the child slice it 600 was actually read against - see `_reassociate_with_original_slice`. `depth` bounds the recursion 601 independently of `request_window_splitter`'s own no-progress guard: it is enforced here so a 602 misbehaving custom cursor cannot recurse indefinitely regardless of what that guard does or does not 603 catch. 604 """ 605 record_generator = partial( 606 self._parse_records, 607 stream_slice=stream_slice, 608 records_schema=records_schema, 609 ) 610 611 emitted_count = 0 612 try: 613 records: Iterable[Mapping[str, Any]] = self._read_pages(record_generator, stream_slice) 614 if self.post_pagination_filter: 615 # A data feed paginates until it reaches a record older than the cursor, so the page that triggers the stop 616 # condition still holds already-synced records. Those are filtered here rather than in the record selector 617 # so that the paginator keeps seeing the whole page: the stop condition is evaluated on the last record of 618 # the page, which is precisely one of the records being dropped. Two consequences of filtering this late: 619 # the pagination tracker observes the dropped records, and a `file_uploader` on the record selector has 620 # already uploaded their files by the time they are dropped. 621 records = self.post_pagination_filter.filter_records( 622 records, 623 # the filter is only used for its cursor comparison, which does not read the stream state 624 stream_state={}, 625 stream_slice=stream_slice, 626 ) 627 for record in records: 628 emitted_count += 1 629 yield self._reassociate_with_original_slice(record, original_slice) 630 except RequestWindowSplitRequiredException as exception: 631 if self.request_window_splitting is None or self.request_window_splitter is None: 632 raise RequestWindowSplitNotSupportedException(stream_name=self.name) from exception 633 634 if emitted_count: 635 # Splitting and re-reading the window re-emits these records: a failed partition is never 636 # checkpointed, so the next attempt would re-emit them anyway. Always replaying makes forward 637 # progress instead of retrying the same oversized window forever; primary-key dedup downstream 638 # handles the duplicates either way. 639 LOGGER.warning( 640 f"Stream {self.name} already emitted {emitted_count} record(s) from {stream_slice} before " 641 f"it was rejected; splitting and re-reading the window may re-emit them." 642 ) 643 644 if depth >= _MAX_REQUEST_WINDOW_SPLIT_DEPTH: 645 # Logged separately from the exception below so it stays visible even if the trace message's 646 # internal_message isn't surfaced by whatever catches it. 647 LOGGER.warning( 648 f"Stream {self.name} hit the maximum request window split depth " 649 f"({_MAX_REQUEST_WINDOW_SPLIT_DEPTH}) while splitting {original_slice}." 650 ) 651 raise AirbyteTracedException( 652 internal_message=f"Stream {self.name} exceeded the maximum request window split depth of {_MAX_REQUEST_WINDOW_SPLIT_DEPTH} while splitting {original_slice}", 653 # `transient_error`, so the only remediation is the connector's own, if it defined one. 654 message=self._with_request_window_failure_message( 655 f"Stream {self.name} could not split its request window to a size the API accepts " 656 f"within {_MAX_REQUEST_WINDOW_SPLIT_DEPTH} splits. This usually means the stream's " 657 f"cursor is not actually shrinking the window on each split; if it uses a custom " 658 f"cursor, check its `split_request_window` implementation." 659 ), 660 failure_type=FailureType.transient_error, 661 ) from exception 662 663 children = self.request_window_splitter( 664 stream_slice, self.request_window_splitting.min_split_window 665 ) 666 if children is None: 667 min_split_window_note = ( 668 f", its configured `min_split_window` ({self.request_window_splitting.min_split_window})" 669 if self.request_window_splitting.min_split_window 670 else "" 671 ) 672 raise AirbyteTracedException( 673 internal_message=f"Stream {self.name} could not split its request window {stream_slice} any further", 674 # `transient_error` regardless of how the triggering response was classified: exhaustion is 675 # never the user's fault, matching `PageSizeReducer`'s equivalent exhaustion branch. 676 message=self._with_request_window_failure_message( 677 f"The API kept rejecting stream {self.name}'s request window even at the smallest window " 678 f"its cursor allows{min_split_window_note}, or splitting the window further would not " 679 f"make progress." 680 ), 681 failure_type=FailureType.transient_error, 682 ) from exception 683 684 LOGGER.info( 685 f"Stream {self.name}: the API rejected request window {stream_slice} (split depth {depth}); " 686 f"reducing it to {children} and reading each in turn." 687 ) 688 for child in children: 689 yield from self._read_records_or_split_request_window( 690 records_schema, child, original_slice, depth + 1 691 ) 692 693 def _with_request_window_failure_message(self, message: str) -> str: 694 """ 695 :return: the message followed by `request_window_splitting`'s configured `failure_message`, when it 696 defined one - mirroring `PageSizeReducer._with_failure_message` 697 """ 698 assert self.request_window_splitting is not None 699 failure_message = (self.request_window_splitting.failure_message or "").strip() 700 if not failure_message: 701 return message 702 return f"{message} {failure_message}" 703 704 @staticmethod 705 def _reassociate_with_original_slice( 706 record: StreamData, original_slice: StreamSlice 707 ) -> StreamData: 708 """ 709 Records read while recursing into a split child window are parsed against that child's `StreamSlice`, 710 so `RecordSelector` stamps them with the child as `associated_slice`. `DeclarativePartition.read()` 711 passes an already-built `Record` through unchanged rather than re-wrapping it, so left uncorrected the 712 child slice - not the partition's own slice - is what `ConcurrentCursor.observe()` would key its 713 per-partition bookkeeping by. `close_partition()` always looks that bookkeeping up by the partition's own 714 slice, so it would never find it, silently losing the precise most-recently-observed cursor value for a 715 split partition (recoverable in practice today only because callers of that lookup fall back to the 716 slice's end boundary when it is missing). Re-stamping every record with `original_slice` here restores 717 the same tracking a non-split read already gets, rather than relying on that fallback. 718 """ 719 if isinstance(record, Record) and record.associated_slice is not original_slice: 720 return Record( 721 data=record.data, 722 stream_name=record.stream_name, 723 associated_slice=original_slice, 724 file_reference=record.file_reference, 725 ) 726 return record 727 728 def _parse_records( 729 self, 730 response: Optional[requests.Response], 731 records_schema: Mapping[str, Any], 732 stream_slice: Optional[StreamSlice], 733 ) -> Iterable[Record]: 734 yield from self._parse_response( 735 response, 736 stream_slice=stream_slice, 737 records_schema=records_schema, 738 ) 739 740 def must_deduplicate_query_params(self) -> bool: 741 return True 742 743 @staticmethod 744 def _to_partition_key(to_serialize: Any) -> str: 745 # separators have changed in Python 3.4. To avoid being impacted by further change, we explicitly specify our own value 746 return json.dumps(to_serialize, indent=None, separators=(",", ":"), sort_keys=True)
Retrieves records by synchronously sending requests to fetch records.
The retriever acts as an orchestrator between the requester, the record selector, the paginator, and the stream slicer.
For each stream slice, submit requests until there are no more pages of records to fetch.
This retriever currently inherits from HttpStream to reuse the request submission and pagination machinery. As a result, some of the parameters passed to some methods are unused. The two will be decoupled in a future release.
Attributes:
- stream_name (str): The stream's name
- stream_primary_key (Optional[Union[str, List[str], List[List[str]]]]): The stream's primary key
- requester (Requester): The HTTP requester
- record_selector (HttpSelector): The record selector
- paginator (Optional[Paginator]): The paginator
- stream_slicer (Optional[StreamSlicer]): The stream slicer
- parameters (Mapping[str, Any]): Additional runtime parameters to be used for string interpolation
- post_pagination_filter (Optional[ClientSideIncrementalRecordFilterDecorator]): Set for data feed streams only. Records the cursor considers already synced are dropped once pagination has observed them
- page_size_reduction (Optional[PageSizeReduction]): How much to shrink the page size when an error handler
resolves to
ResponseAction.REDUCE_PAGE_SIZE.Nonedisables page size reduction entirely. It is immutable configuration; the page size in effect lives in aPageSizeReducerthat_read_pagescreates per call, so the retriever and its paginator - both shared by every partition of the stream, read concurrently - stay stateless. Whenpage_size_reductionisNoneno reducer is created and_read_pageskeeps its previous behaviour - request_window_splitting (Optional[RequestWindowSplitting]): Policy applied when an error handler
resolves to
ResponseAction.SPLIT_REQUEST_WINDOW, or whenRequestWindowSplitRequiredExceptionis raised directly by custom code.Nonedisables request-window splitting entirely:read_recordsconverts the exception intoRequestWindowSplitNotSupportedExceptioninstead of attempting a split it cannot perform. - request_window_splitter: Bound method - normally the stream's own cursor's
split_request_window- asked to replace a failingStreamSlicewith smaller children.Nonehas the same effect asrequest_window_splittingbeingNone.
SimpleRetriever( requester: airbyte_cdk.Requester, record_selector: airbyte_cdk.sources.declarative.extractors.HttpSelector, config: Mapping[str, Any], parameters: dataclasses.InitVar[typing.Mapping[str, typing.Any]], name: str = <property object>, primary_key: Union[str, List[str], List[List[str]], NoneType] = <property object>, paginator: Optional[airbyte_cdk.sources.declarative.requesters.paginators.Paginator] = None, stream_slicer: airbyte_cdk.sources.declarative.stream_slicers.StreamSlicer = <factory>, request_option_provider: airbyte_cdk.sources.declarative.requesters.request_options.RequestOptionsProvider = <factory>, ignore_stream_slicer_parameters_on_paginated_requests: bool = False, additional_query_properties: Optional[airbyte_cdk.sources.declarative.requesters.query_properties.QueryProperties] = None, log_formatter: Optional[Callable[[requests.models.Response], Any]] = None, pagination_tracker_factory: Callable[[], airbyte_cdk.sources.declarative.retrievers.pagination_tracker.PaginationTracker] = <factory>, post_pagination_filter: Optional[airbyte_cdk.sources.declarative.extractors.record_filter.ClientSideIncrementalRecordFilterDecorator] = None, page_size_reduction: Optional[airbyte_cdk.sources.declarative.retrievers.page_size_reducer.PageSizeReduction] = None, request_window_splitting: Optional[airbyte_cdk.sources.declarative.retrievers.request_window_splitting.RequestWindowSplitting] = None, request_window_splitter: Optional[Callable[[airbyte_cdk.StreamSlice, Optional[datetime.timedelta]], Optional[List[airbyte_cdk.StreamSlice]]]] = None)
requester: airbyte_cdk.Requester
record_selector: airbyte_cdk.sources.declarative.extractors.HttpSelector
name: str
160 @property # type: ignore 161 def name(self) -> str: 162 """ 163 :return: Stream name 164 """ 165 return ( 166 str(self._name.eval(self.config)) 167 if isinstance(self._name, InterpolatedString) 168 else self._name 169 )
Returns
Stream name
primary_key: Union[str, List[str], List[List[str]], NoneType]
341 @property # type: ignore 342 def primary_key(self) -> Optional[Union[str, List[str], List[List[str]]]]: 343 """The stream's primary key""" 344 return self._primary_key
The stream's primary key
stream_slicer: airbyte_cdk.sources.declarative.stream_slicers.StreamSlicer
request_option_provider: airbyte_cdk.sources.declarative.requesters.request_options.RequestOptionsProvider
additional_query_properties: Optional[airbyte_cdk.sources.declarative.requesters.query_properties.QueryProperties] =
None
pagination_tracker_factory: Callable[[], airbyte_cdk.sources.declarative.retrievers.pagination_tracker.PaginationTracker]
post_pagination_filter: Optional[airbyte_cdk.sources.declarative.extractors.record_filter.ClientSideIncrementalRecordFilterDecorator] =
None
page_size_reduction: Optional[airbyte_cdk.sources.declarative.retrievers.page_size_reducer.PageSizeReduction] =
None
request_window_splitting: Optional[airbyte_cdk.sources.declarative.retrievers.request_window_splitting.RequestWindowSplitting] =
None
request_window_splitter: Optional[Callable[[airbyte_cdk.StreamSlice, Optional[datetime.timedelta]], Optional[List[airbyte_cdk.StreamSlice]]]] =
None
def
read_records( self, records_schema: Mapping[str, Any], stream_slice: Optional[airbyte_cdk.StreamSlice] = None) -> Iterable[Union[Mapping[str, Any], airbyte_cdk.AirbyteMessage]]:
562 def read_records( 563 self, 564 records_schema: Mapping[str, Any], 565 stream_slice: Optional[StreamSlice] = None, 566 ) -> Iterable[StreamData]: 567 """ 568 Fetch a stream's records from an HTTP API source 569 570 :param records_schema: json schema to describe record 571 :param stream_slice: The stream slice to read data for 572 :return: The records read from the API source 573 """ 574 _slice = stream_slice or StreamSlice(partition={}, cursor_slice={}) # None-check 575 yield from self._read_records_or_split_request_window( 576 records_schema, _slice, _slice, depth=0 577 )
Fetch a stream's records from an HTTP API source
Parameters
- records_schema: json schema to describe record
- stream_slice: The stream slice to read data for
Returns
The records read from the API source
Inherited Members
@deprecated('This class is experimental. Use at your own risk.', category=ExperimentalClassWarning)
@dataclass
class
LazySimpleRetriever769@deprecated( 770 "This class is experimental. Use at your own risk.", 771 category=ExperimentalClassWarning, 772) 773@dataclass 774class LazySimpleRetriever(SimpleRetriever): 775 """ 776 A retriever that supports lazy loading from parent streams. 777 """ 778 779 def _read_pages( 780 self, 781 records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]], 782 stream_slice: StreamSlice, 783 ) -> Iterable[Record]: 784 response = stream_slice.extra_fields["child_response"] 785 if response: 786 last_page_size, last_record = 0, None 787 for record in records_generator_fn(response): # type: ignore[call-arg] # only _parse_records expected as a func 788 last_page_size += 1 789 last_record = record 790 yield record 791 792 next_page_token = self._next_page_token( 793 response, 794 last_page_size, 795 last_record, 796 None, 797 **stream_slice_kwargs(self._next_page_token, stream_slice), 798 ) 799 if next_page_token: 800 yield from self._paginate( 801 next_page_token, 802 records_generator_fn, 803 stream_slice, 804 ) 805 806 yield from [] 807 else: 808 # coderabbit detected an interesting bug/gap where if we were to not get a child_response, we 809 # might recurse forever. This might not be the case, but it is worth noting that this code path 810 # isn't comprehensively tested. 811 yield from self._read_pages(records_generator_fn, stream_slice) 812 813 def _paginate( 814 self, 815 next_page_token: Any, 816 records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]], 817 stream_slice: StreamSlice, 818 ) -> Iterable[Record]: 819 """Handle pagination by fetching subsequent pages.""" 820 pagination_complete = False 821 822 while not pagination_complete: 823 response = self._fetch_next_page(stream_slice, next_page_token) 824 last_page_size, last_record = 0, None 825 826 for record in records_generator_fn(response): # type: ignore[call-arg] # only _parse_records expected as a func 827 last_page_size += 1 828 last_record = record 829 yield record 830 831 if not response: 832 pagination_complete = True 833 else: 834 last_page_token_value = ( 835 next_page_token.get("next_page_token") if next_page_token else None 836 ) 837 next_page_token = self._next_page_token( 838 response, 839 last_page_size, 840 last_record, 841 last_page_token_value, 842 **stream_slice_kwargs(self._next_page_token, stream_slice), 843 ) 844 845 if not next_page_token: 846 pagination_complete = True
A retriever that supports lazy loading from parent streams.
LazySimpleRetriever( requester: airbyte_cdk.Requester, record_selector: airbyte_cdk.sources.declarative.extractors.HttpSelector, config: Mapping[str, Any], parameters: dataclasses.InitVar[typing.Mapping[str, typing.Any]], name: str = <property object>, primary_key: Union[str, List[str], List[List[str]], NoneType] = <property object>, paginator: Optional[airbyte_cdk.sources.declarative.requesters.paginators.Paginator] = None, stream_slicer: airbyte_cdk.sources.declarative.stream_slicers.StreamSlicer = <factory>, request_option_provider: airbyte_cdk.sources.declarative.requesters.request_options.RequestOptionsProvider = <factory>, ignore_stream_slicer_parameters_on_paginated_requests: bool = False, additional_query_properties: Optional[airbyte_cdk.sources.declarative.requesters.query_properties.QueryProperties] = None, log_formatter: Optional[Callable[[requests.models.Response], Any]] = None, pagination_tracker_factory: Callable[[], airbyte_cdk.sources.declarative.retrievers.pagination_tracker.PaginationTracker] = <factory>, post_pagination_filter: Optional[airbyte_cdk.sources.declarative.extractors.record_filter.ClientSideIncrementalRecordFilterDecorator] = None, page_size_reduction: Optional[airbyte_cdk.sources.declarative.retrievers.page_size_reducer.PageSizeReduction] = None, request_window_splitting: Optional[airbyte_cdk.sources.declarative.retrievers.request_window_splitting.RequestWindowSplitting] = None, request_window_splitter: Optional[Callable[[airbyte_cdk.StreamSlice, Optional[datetime.timedelta]], Optional[List[airbyte_cdk.StreamSlice]]]] = None)
Inherited Members
- SimpleRetriever
- requester
- record_selector
- config
- parameters
- name
- primary_key
- paginator
- stream_slicer
- request_option_provider
- ignore_stream_slicer_parameters_on_paginated_requests
- additional_query_properties
- log_formatter
- pagination_tracker_factory
- post_pagination_filter
- page_size_reduction
- request_window_splitting
- request_window_splitter
- read_records
- must_deduplicate_query_params