diff --git a/.gitignore b/.gitignore index fc5cd0f..bfc2e0b 100644 --- a/.gitignore +++ b/.gitignore @@ -19,3 +19,6 @@ tap_google_sheets/.vscode/settings.json *.ipynb .DS_Store tap_target_commands.sh + +.idea +.env diff --git a/README.md b/README.md index b813b5f..148a519 100644 --- a/README.md +++ b/README.md @@ -75,6 +75,7 @@ The [**Google Sheets Setup & Authentication**](https://drive.google.com/open?id= - client_secret: authenticates your application - refresh_token: generates an access token to authorize your session - spreadsheet_id: unique identifier for each spreadsheet in Google Drive + - sheet_id: optional worksheet gid. When provided, discovery exposes the matching worksheet as stream `gid:` and sync still calls the Google Sheets API by the worksheet title. - start_date: absolute minimum start date to check file modified - user_agent: tap-name and email address; identifies your application in the Remote API server logs @@ -111,6 +112,7 @@ The [**Google Sheets Setup & Authentication**](https://drive.google.com/open?id= "client_secret": "YOUR_CLIENT_SECRET", "refresh_token": "YOUR_REFRESH_TOKEN", "spreadsheet_id": "YOUR_GOOGLE_SPREADSHEET_ID", + "sheet_id": "0", "start_date": "2019-01-01T00:00:00Z", "user_agent": "tap-google-sheets ", "request_timeout": 300 diff --git a/tap_google_sheets/__init__.py b/tap_google_sheets/__init__.py index 535f7c5..03de778 100644 --- a/tap_google_sheets/__init__.py +++ b/tap_google_sheets/__init__.py @@ -8,6 +8,7 @@ from tap_google_sheets.client import GoogleClient from tap_google_sheets.discover import discover from tap_google_sheets.sync import sync +from tap_google_sheets.streams import get_sheet_id_filter LOGGER = singer.get_logger() @@ -20,10 +21,10 @@ 'user_agent' ] -def do_discover(client, spreadsheet_id): +def do_discover(client, config): LOGGER.info('Starting discover') - catalog = discover(client, spreadsheet_id) + catalog = discover(client, config.get('spreadsheet_id'), sheet_id=get_sheet_id_filter(config)) json.dump(catalog.to_dict(), sys.stdout, indent=2) LOGGER.info('Finished discover') @@ -48,11 +49,14 @@ def main(): spreadsheet_id = config.get('spreadsheet_id') if parsed_args.discover: - do_discover(client, spreadsheet_id) + do_discover(client, config) else: sync(client=client, config=config, - catalog=parsed_args.catalog or discover(client, spreadsheet_id), + catalog=parsed_args.catalog or discover( + client, + spreadsheet_id, + sheet_id=get_sheet_id_filter(config)), state=state) if __name__ == '__main__': diff --git a/tap_google_sheets/discover.py b/tap_google_sheets/discover.py index 1440dee..59ece17 100644 --- a/tap_google_sheets/discover.py +++ b/tap_google_sheets/discover.py @@ -2,11 +2,11 @@ from tap_google_sheets.schema import STREAMS -def discover(client, spreadsheet_id): +def discover(client, spreadsheet_id, sheet_id=None): catalog = Catalog([]) for stream, stream_obj in STREAMS.items(): - stream_object = stream_obj(client, spreadsheet_id) + stream_object = stream_obj(client, spreadsheet_id, sheet_id=sheet_id) schemas, field_metadata = stream_object.get_schemas() # loop over the schema and prepare catalog diff --git a/tap_google_sheets/streams.py b/tap_google_sheets/streams.py index f2b803a..b822992 100644 --- a/tap_google_sheets/streams.py +++ b/tap_google_sheets/streams.py @@ -15,6 +15,28 @@ import tap_google_sheets.schema as schema LOGGER = singer.get_logger() +GID_STREAM_PREFIX = 'gid:' + + +def get_sheet_id_filter(config): + sheet_id = config.get('sheet_id') or os.environ.get('GOOGLE_SHEET_GID') + if sheet_id in (None, ''): + return None + return str(sheet_id) + + +def sheet_matches_filter(sheet, sheet_id_filter): + if sheet_id_filter is None: + return True + sheet_id = sheet.get('properties', {}).get('sheetId') + return str(sheet_id) == str(sheet_id_filter) + + +def get_sheet_stream_name(sheet, sheet_id_filter=None): + if sheet_id_filter is not None: + sheet_id = sheet.get('properties', {}).get('sheetId') + return '{}{}'.format(GID_STREAM_PREFIX, sheet_id) + return sheet.get('properties', {}).get('title') def update_currently_syncing(state, stream_name): """ @@ -126,11 +148,12 @@ class GoogleSheets: params = None state = None - def __init__(self, client, spreadsheet_id, start_date=None, batch_rows=200): + def __init__(self, client, spreadsheet_id, start_date=None, batch_rows=200, sheet_id=None): self.client = client self.config_start_date = start_date self.spreadsheet_id = spreadsheet_id self.batch_rows = batch_rows + self.sheet_id = str(sheet_id) if sheet_id not in (None, '') else None def get_path(self, sheet_title_encoded=""): """ @@ -315,13 +338,15 @@ def get_schemas(self): if sheets: # Loop thru each worksheet in spreadsheet for sheet in sheets: + if not sheet_matches_filter(sheet, self.sheet_id): + continue # GET sheet_json_schema for each worksheet (from function above) sheet_json_schema, columns = schema.get_sheet_metadata(sheet, self.spreadsheet_id, self.client) # SKIP empty sheets (where sheet_json_schema and columns are None) if sheet_json_schema and columns: - sheet_title = sheet.get('properties', {}).get('title') - schemas[sheet_title] = sheet_json_schema + sheet_stream_name = get_sheet_stream_name(sheet, self.sheet_id) + schemas[sheet_stream_name] = sheet_json_schema sheet_mdata = metadata.new() sheet_mdata = metadata.get_standard_metadata( schema=sheet_json_schema, @@ -337,7 +362,7 @@ def get_schemas(self): mdata = metadata.to_map(sheet_mdata) sheet_mdata = metadata.write(mdata, ('properties', column.get('columnName')), 'inclusion', 'unsupported') sheet_mdata = metadata.to_list(mdata) - field_metadata[sheet_title] = sheet_mdata + field_metadata[sheet_stream_name] = sheet_mdata return schemas, field_metadata @@ -462,8 +487,11 @@ def load_data(self, catalog, state, selected_streams, sheets, spreadsheet_time_e if sheets: # Loop through sheets (worksheet tabs) in spreadsheet for sheet in sheets: + if not sheet_matches_filter(sheet, self.sheet_id): + continue sheet_title = sheet.get('properties', {}).get('title') sheet_id = sheet.get('properties', {}).get('sheetId') + sheet_stream_name = get_sheet_stream_name(sheet, self.sheet_id) # GET sheet_metadata and columns sheet_schema, columns = schema.get_sheet_metadata(sheet, self.spreadsheet_id, self.client) @@ -480,26 +508,26 @@ def load_data(self, catalog, state, selected_streams, sheets, spreadsheet_time_e # SHEET_DATA # Should this worksheet tab be synced? - if sheet_title in selected_streams: - LOGGER.info('STARTED Syncing Sheet {}'.format(sheet_title)) - update_currently_syncing(self.state, sheet_title) - selected_fields = get_selected_fields(catalog, sheet_title) # -------------------- - LOGGER.info('Stream: {}, selected_fields: {}'.format(sheet_title, selected_fields)) - write_schema(catalog, sheet_title) + if sheet_stream_name in selected_streams: + LOGGER.info('STARTED Syncing Sheet {}'.format(sheet_stream_name)) + update_currently_syncing(self.state, sheet_stream_name) + selected_fields = get_selected_fields(catalog, sheet_stream_name) # -------------------- + LOGGER.info('Stream: {}, selected_fields: {}'.format(sheet_stream_name, selected_fields)) + write_schema(catalog, sheet_stream_name) # Emit a Singer ACTIVATE_VERSION message before initial sync (but not subsequent syncs) # everytime after each sheet sync is complete. # This forces hard deletes on the data downstream if fewer records are sent. # https://github.com/singer-io/singer-python/blob/master/singer/messages.py#L137 - last_integer = int(get_bookmark(self.state, sheet_title, 0)) + last_integer = int(get_bookmark(self.state, sheet_stream_name, 0)) activate_version = int(time.time() * 1000) activate_version_message = singer.ActivateVersionMessage( - stream=sheet_title, + stream=sheet_stream_name, version=activate_version) if last_integer == 0: # initial load, send activate_version before AND after data sync singer.write_message(activate_version_message) - LOGGER.info('INITIAL SYNC, Stream: {}, Activate Version: {}'.format(sheet_title, activate_version)) + LOGGER.info('INITIAL SYNC, Stream: {}, Activate Version: {}'.format(sheet_stream_name, activate_version)) # Determine max range of columns and rows for "paging" through the data sheet_last_col_index = 1 @@ -568,7 +596,7 @@ def load_data(self, catalog, state, selected_streams, sheets, spreadsheet_time_e # Process records, send batch of records to target record_count = self.process_records( catalog=catalog, - stream_name=sheet_title, + stream_name=sheet_stream_name, records=sheet_data_transformed, time_extracted=spreadsheet_time_extracted, version=activate_version) @@ -584,10 +612,10 @@ def load_data(self, catalog, state, selected_streams, sheets, spreadsheet_time_e # End of Stream: Send Activate Version and update State singer.write_message(activate_version_message) - write_bookmark(self.state, sheet_title, activate_version) - LOGGER.info('COMPLETE SYNC, Stream: {}, Activate Version: {}'.format(sheet_title, activate_version)) + write_bookmark(self.state, sheet_stream_name, activate_version) + LOGGER.info('COMPLETE SYNC, Stream: {}, Activate Version: {}'.format(sheet_stream_name, activate_version)) LOGGER.info('FINISHED Syncing Sheet {}, Total Rows: {}'.format( - sheet_title, row_num - 2)) # subtract 1 for header row + sheet_stream_name, row_num - 2)) # subtract 1 for header row update_currently_syncing(self.state, None) # SHEETS_LOADED diff --git a/tap_google_sheets/sync.py b/tap_google_sheets/sync.py index e50295f..855a376 100644 --- a/tap_google_sheets/sync.py +++ b/tap_google_sheets/sync.py @@ -1,5 +1,5 @@ import singer -from tap_google_sheets.streams import STREAMS, SheetsLoadData, write_bookmark, strftime +from tap_google_sheets.streams import STREAMS, SheetsLoadData, get_sheet_id_filter, write_bookmark, strftime LOGGER = singer.get_logger() @@ -26,11 +26,13 @@ def sync(client, config, catalog, state): LOGGER.info("No stream is selected.") return + sheet_id = get_sheet_id_filter(config) + # loop through main streams for stream_name, stream_obj in STREAMS.items(): # get the stream object - stream_obj = stream_obj(client, config.get("spreadsheet_id"), config.get("start_date")) + stream_obj = stream_obj(client, config.get("spreadsheet_id"), config.get("start_date"), sheet_id=sheet_id) # to sync the sheet's data, we need to get "spreadsheet_metadata" if stream_name == "spreadsheet_metadata": @@ -44,7 +46,12 @@ def sync(client, config, catalog, state): # get sheets from the metadata sheets = spreadsheet_metadata.get("sheets") # class to load sheet's data - sheets_load_data = SheetsLoadData(client, config.get("spreadsheet_id"), config.get("start_date"), config.get('batch_rows')) + sheets_load_data = SheetsLoadData( + client, + config.get("spreadsheet_id"), + config.get("start_date"), + config.get('batch_rows'), + sheet_id=sheet_id) # perform sheet's sync and get sheet's metadata and sheet loaded records for "sheet_metadata" and "sheets_loaded" streams sheet_metadata_records, sheets_loaded_records = sheets_load_data.load_data(catalog=catalog, diff --git a/tests/unittests/test_sheet_id_streams.py b/tests/unittests/test_sheet_id_streams.py new file mode 100644 index 0000000..35c52e6 --- /dev/null +++ b/tests/unittests/test_sheet_id_streams.py @@ -0,0 +1,99 @@ +import unittest +from unittest import mock + +from tap_google_sheets.client import GoogleClient +from tap_google_sheets.streams import SheetsLoadData, SpreadSheetMetadata + + +sheet_schema = { + 'type': 'object', + 'additionalProperties': False, + 'properties': { + '__sdc_spreadsheet_id': {'type': ['null', 'string']}, + '__sdc_sheet_id': {'type': ['null', 'integer']}, + '__sdc_row': {'type': ['null', 'integer']}, + 'value': {'type': ['null', 'string']}, + }, +} +columns = [ + { + 'columnIndex': 1, + 'columnLetter': 'A', + 'columnName': 'value', + 'columnType': 'stringValue', + 'columnSkipped': False, + } +] + + +def _sheet(sheet_id, title): + return { + 'properties': { + 'sheetId': sheet_id, + 'title': title, + 'index': 0, + 'sheetType': 'GRID', + 'gridProperties': { + 'rowCount': 10, + 'columnCount': 1, + }, + } + } + + +class FakeClient: + def get(self, **_kwargs): + return {'sheets': [_sheet(0, 'First'), _sheet(123, 'Second')]} + + +class TestSheetIdStreams(unittest.TestCase): + @mock.patch('tap_google_sheets.streams.schema.get_sheet_metadata', return_value=[sheet_schema, columns]) + def test_discovery_can_name_selected_sheet_stream_by_gid(self, _mocked_sheet_metadata): + stream = SpreadSheetMetadata(FakeClient(), 'spreadsheet-id', sheet_id='123') + + schemas, field_metadata = stream.get_schemas() + + self.assertIn('gid:123', schemas) + self.assertIn('gid:123', field_metadata) + self.assertNotIn('Second', schemas) + self.assertNotIn('gid:0', schemas) + + @mock.patch('tap_google_sheets.client.GoogleClient.get') + @mock.patch('tap_google_sheets.streams.schema.get_sheet_metadata', return_value=[sheet_schema, columns]) + @mock.patch('tap_google_sheets.streams.get_selected_fields', return_value=[]) + @mock.patch('tap_google_sheets.streams.write_schema') + @mock.patch('tap_google_sheets.streams.GoogleSheets.process_records') + def test_load_data_uses_gid_stream_name_and_title_for_api( + self, + mock_process_records, + mock_write_schema, + _mocked_get_selected_fields, + _mocked_sheet_metadata, + mocked_get, + ): + client = GoogleClient('dummy_client_id', 'dummy_client_secret', 'dummy_refresh_token', 300) + sheets_load_data = SheetsLoadData( + client, + 'spreadsheet-id', + '2019-01-01T00:00:00Z', + batch_rows=200, + sheet_id='1260142713', + ) + + sheets_load_data.load_data({}, {}, ['gid:1260142713'], [_sheet(1260142713, 'Sheet13')], 'time') + + mock_write_schema.assert_called_with({}, 'gid:1260142713') + self.assertEqual(mock_process_records.call_args.kwargs['stream_name'], 'gid:1260142713') + self.assertEqual( + mock.call( + api='sheets', + endpoint='Sheet13', + params='dateTimeRenderOption=SERIAL_NUMBER&valueRenderOption=FORMATTED_VALUE&majorDimension=ROWS', + path="spreadsheets/spreadsheet-id/values/'Sheet13'!A2:A10", + ), + mocked_get.mock_calls[0], + ) + + +if __name__ == '__main__': + unittest.main()