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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -19,3 +19,6 @@ tap_google_sheets/.vscode/settings.json
*.ipynb
.DS_Store
tap_target_commands.sh

.idea
.env
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:<sheet_id>` 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

Expand Down Expand Up @@ -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 <api_user_email@example.com>",
"request_timeout": 300
Expand Down
12 changes: 8 additions & 4 deletions tap_google_sheets/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand All @@ -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')

Expand All @@ -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__':
Expand Down
4 changes: 2 additions & 2 deletions tap_google_sheets/discover.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
62 changes: 45 additions & 17 deletions tap_google_sheets/streams.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
"""
Expand Down Expand Up @@ -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=""):
"""
Expand Down Expand Up @@ -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,
Expand All @@ -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

Expand Down Expand Up @@ -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)
Expand All @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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
Expand Down
13 changes: 10 additions & 3 deletions tap_google_sheets/sync.py
Original file line number Diff line number Diff line change
@@ -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()

Expand All @@ -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":
Expand All @@ -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,
Expand Down
99 changes: 99 additions & 0 deletions tests/unittests/test_sheet_id_streams.py
Original file line number Diff line number Diff line change
@@ -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()