-
Notifications
You must be signed in to change notification settings - Fork 10
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
TDL-18879: Clear offset when max_skip
error is encountered and return 0 if skip
is greater than 250k in the current state.
#29
Open
hpatel41
wants to merge
7
commits into
master
Choose a base branch
from
TDL-18879-update-default-window-size-for-activities-stream-to-5-days
base: master
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from 4 commits
Commits
Show all changes
7 commits
Select commit
Hold shift + click to select a range
7848ba0
updated default window size to 5 days for activities stream
2add6aa
addressed review comments
e6ef2be
updated the code to clear offset of max_skip error is encountered
451c435
resolved review comments
a33f38b
added unittests
5ce87fb
updated error message
4a99c74
updated error message
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -12,6 +12,9 @@ | |
|
||
LOGGER = singer.get_logger() | ||
|
||
# default date date window size in days | ||
DATE_WINDOW_SIZE = 15 | ||
|
||
PATHS = { | ||
IDS.CUSTOM_FIELDS: "/custom_fields/lead/", | ||
IDS.LEADS: "/lead/", | ||
|
@@ -115,15 +118,24 @@ def paginated_sync(tap_stream_id, ctx, request, start_date): | |
# There may be streams other than `leads` that will run into | ||
# `max_skip` errors but YAGNI. We can make the tap more | ||
# complicated once we have an extant need for it. | ||
if 'max_skip = ' in str(e) and tap_stream_id == IDS.LEADS: | ||
LOGGER.info(("Hit max_skip error. " | ||
"Setting bookmark to `{}` and restarting pagination.".format( | ||
max_bookmark))) | ||
skip = 0 | ||
ctx.clear_offsets(tap_stream_id) | ||
ctx.set_bookmark(bookmark(tap_stream_id), max_bookmark) | ||
_request = create_leads_request(ctx) | ||
ctx.write_state() | ||
if 'max_skip = ' in str(e): | ||
if tap_stream_id == IDS.ACTIVITIES: | ||
LOGGER.warning("Hit max_skip error so clearing skip offset, please reduce the date window size and try again.") | ||
# clear offset | ||
ctx.clear_offsets(tap_stream_id) | ||
ctx.write_state() | ||
raise | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Need to add the message to reduce the date window. There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Updated error message |
||
elif tap_stream_id == IDS.LEADS: | ||
LOGGER.info(("Hit max_skip error. " | ||
"Setting bookmark to `{}` and restarting pagination.".format( | ||
max_bookmark))) | ||
skip = 0 | ||
ctx.clear_offsets(tap_stream_id) | ||
ctx.set_bookmark(bookmark(tap_stream_id), max_bookmark) | ||
_request = create_leads_request(ctx) | ||
ctx.write_state() | ||
else: | ||
raise | ||
else: | ||
raise | ||
ctx.clear_offsets(tap_stream_id) | ||
|
@@ -168,15 +180,15 @@ def sync_activities(ctx): | |
|
||
try: | ||
# get date window from config | ||
date_window = int(ctx.config.get("date_window", 15)) | ||
# if date_window is 0, '0' or None, then set default window size of 15 days | ||
date_window = int(ctx.config.get("date_window", DATE_WINDOW_SIZE)) | ||
# if date_window is 0, '0' or None, then set the default window size to DATE_WINDOW_SIZE (15 days) | ||
if not date_window: | ||
LOGGER.warning("Invalid value of date window is passed: \'{}\', using default window size of 15 days.".format(ctx.config.get("date_window"))) | ||
date_window = 15 | ||
LOGGER.warning("Invalid value of date window is passed: \'{}\', using default window size of {} days.".format(ctx.config.get("date_window"), DATE_WINDOW_SIZE)) | ||
date_window = DATE_WINDOW_SIZE | ||
except ValueError: | ||
LOGGER.warning("Invalid value of date window is passed: \'{}\', using default window size of 15 days.".format(ctx.config.get("date_window"))) | ||
LOGGER.warning("Invalid value of date window is passed: \'{}\', using default window size of {} days.".format(ctx.config.get("date_window"), DATE_WINDOW_SIZE)) | ||
# In case of empty string(''), use default window | ||
date_window = 15 | ||
date_window = DATE_WINDOW_SIZE | ||
|
||
LOGGER.info("Using offset seconds {}".format(offset_secs)) | ||
start_date -= timedelta(seconds=offset_secs) | ||
|
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,41 @@ | ||
from tap_closeio.schemas import IDS | ||
import unittest | ||
from tap_closeio.context import Context | ||
from tap_closeio.streams import paginated_sync, create_request | ||
|
||
class TestExistingStateOffset(unittest.TestCase): | ||
def test_existing_state_existing_state_offset_greater_than_250K(self): | ||
config = { | ||
"start_date": "2022-01-01", | ||
"api_key": "test_API_key" | ||
} | ||
state = { | ||
"currently_syncing": "activities", | ||
"bookmarks": { | ||
"activities": { | ||
"date_created": "2022-04-01T00:00:00", | ||
"offset": {"skip": 259000} | ||
} | ||
} | ||
} | ||
context = Context(config, state) | ||
offset = context.get_offset(["activities", "skip"]) | ||
self.assertEqual(offset, 0) | ||
|
||
def test_existing_state_existing_state_offset_lesser_than_250K(self): | ||
config = { | ||
"start_date": "2022-01-01", | ||
"api_key": "test_API_key" | ||
} | ||
state = { | ||
"currently_syncing": "activities", | ||
"bookmarks": { | ||
"activities": { | ||
"date_created": "2022-04-01T00:00:00", | ||
"offset": {"skip": 1000} | ||
} | ||
} | ||
} | ||
context = Context(config, state) | ||
offset = context.get_offset(["activities", "skip"]) | ||
self.assertEqual(offset, 1000) |
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
It should be an error and must be raised with this message.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Added