Repository navigation
feat(low-code): support request window splitting on partitioned streams - #1192
Conversation
ConcurrentPerPartitionCursor gains split_request_window, delegating to one state-free ConcurrentCursor built from the stream's cursor factory on first use. The factory registers RequestWindowSplitting so a CustomRetriever can declare it, and create_default_stream passes the cursor's splitter to every retriever. A CustomRetriever's request_window_splitting is not validated at load time, consistent with other custom components.
👋 Greetings, Airbyte Team Member!Here are some helpful tips and reminders for your convenience. 💡 Show Tips and TricksTesting This CDK VersionYou can test this version of the CDK using the following: # Run the CLI from this branch:
uvx 'git+https://github.com/airbytehq/airbyte-python-cdk.git@kaizerch/feat/17179-17176-cdk-request-window-splitting-support-google-ads-tiktok#egg=airbyte-python-cdk[dev]' --help
# Update a connector to use the CDK from this branch ref:
cd airbyte-integrations/connectors/source-example
poe use-cdk-branch kaizerch/feat/17179-17176-cdk-request-window-splitting-support-google-ads-tiktokPR Slash CommandsAirbyte Maintainers can execute the following slash commands on your PR:
|
|
/prerelease |
PyTest Results (Fast)4 930 tests +8 4 918 ✅ +8 10m 19s ⏱️ -32s Results for commit 1d62b0f. ± Comparison against base commit d5536bc. This pull request removes 1 and adds 9 tests. Note that renamed tests count towards both.♻️ This comment has been updated with latest results. |
PyTest Results (Full)4 933 tests 4 921 ✅ 15m 47s ⏱️ Results for commit 1d62b0f. ♻️ This comment has been updated with latest results. |
|
The description still covers the first commit, and it becomes the release notes. Could you update it?
|
Thanks, updated. The description now matches e0371aa only: the custom-retriever wiring and tests are gone, and it says a CustomRetriever declaring the block still fails at load, with Google Ads wiring the splitter in Python. Test plan is corrected: all 10 cases fail on main; with only the cursor change reverted, 9 fail, and the tenth asserts the reworded error message. Only the success case runs in global-cursor mode. Spotlight links are pinned to e0371aa. |
|
/prerelease
|
Overview
👉 TL;DR: Streams that sync each account separately, like Google Ads customers or TikTok advertisers, can now automatically split a date window the API rejects instead of failing the sync. Only the account whose request was rejected gets split; other accounts, and the progress saved for them, are untouched.
Specifically,
ConcurrentPerPartitionCursorgainssplit_request_window, so the factory no longer rejectsrequest_window_splittingon a stream with a partition router.Builds on #1171 (released in
7.32.0), which addedSPLIT_REQUEST_WINDOWfor unpartitioned streams only:Unblocks https://github.com/airbytehq/airbyte-internal-issues/issues/17176 (TikTok Marketing) and https://github.com/airbytehq/airbyte-internal-issues/issues/17179 (Google Ads), part of epic https://github.com/airbytehq/airbyte-internal-issues/issues/17173.
Changes
ConcurrentPerPartitionCursor.split_request_window(stream_slice, min_split_window): delegates to one state-freeConcurrentCursor, built from the stream's own_cursor_factoryon first use and cached under the existing_lock. The children keep the slice'spartitionandextra_fields.hasattr(cursor, "split_request_window")check now passes for partitioned streams, so no factory wiring changes. The "unsupported cursor" error no longer says partition routers are unsupported. Every other load-time check applies to partitioned streams unchanged.ConcurrentPerPartitionCursor: halves keeppartitionandextra_fields, a 1-day window returnsNone,min_split_windowis honored, and the splitting cursor is built once and only on first use.http_codes: [400]filter. That account is read as two halves, the other once, and each keeps its own cursor in the final per-partition state. It runs withstart_time_option/end_time_option, and with a list-formpartition_routerplusstream_intervalinterpolation.SWITCH_TO_GLOBAL_LIMIT: it splits and checkpoints one global cursor. Only this success case runs in global-cursor mode; the still-rejected case below runs per partition.min_split_window. The stream fails withtransient_error, that account's cursor stays at its start instead of moving to the record already read, and the other account is still checkpointed.Review Spotlight
Reviewers with limited time, please review first:
concurrent_partition_cursor.py#L682-L708: the newsplit_request_windowand why one state-free cursor is enoughmodel_to_component_factory.py#L4368-L4372: the reworded error, the only factory changetest_concurrent_declarative_source.py#L5396-L5519: the end-to-end per-partition, global-cursor and still-rejected testsNot in scope: custom retrievers
This PR doesn't change custom retrievers. A
CustomRetrieverthat declaresrequest_window_splittingin YAML still fails at load ("Subcomponent creation has not been implemented for 'RequestWindowSplitting'"), and the factory doesn't hand it a splitter. Google Ads, the only custom-retriever consumer, does the wiring in Python, not YAML:GoogleAdsRetrievergets acursorfield (the factory already passescursor=to every retriever) and setsrequest_window_splitterandrequest_window_splittingin__post_init__. I checked that pattern against this branch with a scratch custom retriever on a partitioned read: only the rejected account split, and both accounts' cursors were correct.Risks
ConcurrentCursorpath.7.32.0: split children aren't clamped again to the cursor's start or end, and records already emitted before a rejection are re-emitted when the window is split and re-read (logged as a warning). Both behave exactly as they do for unpartitioned streams.Post-merge actions
7.33.0(minor: new capability, no breaking change). Publishing the tag also publishes thesource-declarative-manifestimage; confirmairbyte/source-declarative-manifest:7.33.0appears on DockerHub.Test plan
e0371aacand fail onmain. With only theconcurrent_partition_cursor.pychange reverted, 9 of them fail. The tenth,test_given_no_incremental_sync_and_request_window_splitting_then_raise, asserts the reworded error message, so it fails only when the message is reverted too.poetry run pytest unit_tests -m "not slow and not requires_creds": 4,917 passed, 2 failed. Both failures are Docker image-build tests (test_docker_image_build_and_spec,test_docker_image_build_and_check) that failed locally because Docker wasn't running; unrelated to this change.ruff check,ruff format --check, andmypyon the two changed source files pass.Design notes: why one state-free cursor instead of each partition's own
Show/Hide Content
ConcurrentCursor.split_request_windowis a pure function of the slice's boundaries and the cursor's granularity and date format. It never reads cursor state. Every per-partition cursor comes from the same_cursor_factory, so they would all return the same split.SWITCH_TO_GLOBAL_LIMIT), when_ensure_partition_limitmay already have evicted it._cursor_factory.createa fixedside_effectlist keep working. This uses the samecreate(stream_state={}, runtime_lookback_window=None)call asget_cursor_datetime_from_state.observeandclose_partitionneed no change:SimpleRetrieveralready re-associates records read from split children with the original slice, and a partition only closes after all of its children succeed.split_request_windowis not added to theCursorABC. A default there would letFinalStateCursorpass the factory'shasattrcheck.