import json import queue import threading from datetime import datetime from ssavva.artists.backfill_spotify_and_apple.spotify_api import \ SpotifyAPI from ssavva.artists.backfill_spotify_and_apple.db_connections import get_releases, dump_responses API_QUEUE_MAX_SIZE = 1000 DB_QUEUE_MAX_SIZE = 1000 API_WORKERS_NUMBER = 2 DB_WORKERS_NUMBER = 1 def gather_data(): api_queue = queue.Queue(maxsize=API_QUEUE_MAX_SIZE) db_queue = queue.Queue(maxsize=DB_QUEUE_MAX_SIZE) spotify_api = SpotifyAPI( [ {'client_id': '', 'secret_id': ''}, ], ) api_worker_pool = WorkerPool( API_WORKERS_NUMBER, _api_worker, (api_queue, db_queue, spotify_api)) db_worker_pool = WorkerPool( DB_WORKERS_NUMBER, _db_worker, (db_queue,)) i = 0 for r in get_releases(): i += 1 api_queue.put(r) if i % 100000 == 0: print(f'loaded to queue {i}') api_queue.join() db_queue.join() api_worker_pool.add_stoppers(api_queue) api_worker_pool.join_workers() db_worker_pool.add_stoppers(db_queue) db_worker_pool.join_workers() def _api_worker(api_queue, db_queue, spotify_api, number): counter = 0 while True: counter += 1 item = api_queue.get() if item is None: break display_upc = item['display_upc'] resp, sign = spotify_api.api_call( 'https://api.spotify.com/v1/search', {'q': f'upc:{display_upc}', 'type': 'album'} ) response = json.loads(resp) if response: item['response'] = resp db_queue.put(item) api_queue.task_done() else: api_queue.task_done() # api_queue.put(item) if counter % 10000 == 0: print(f'api worker {number}', datetime.now()) def _db_worker(q, number): items = [] counter = 0 while True: counter += 1 item = q.get() if item is None: dump_responses(items) break items.append( (item['release_id'], item['display_upc'], item['upc'], item['response'])) if counter % 100 == 0: print('{} buffer {}'.format(number, len(items))) if counter % 1000 == 0: print(f'db dump in {number}', datetime.now()) dump_responses(items) items = [] q.task_done() class WorkerPool: def __init__(self, workers_number, target, worker_args): self._workers_number = workers_number self._target = target self._worker_args = worker_args self._workers = [] for i in range(self._workers_number): worker = threading.Thread(target=self._target, args=worker_args + (i,)) worker.start() self._workers.append(worker) def add_stoppers(self, q): for i in range(self._workers_number): q.put(None) def join_workers(self): for worker in self._workers: worker.join() if __name__ == '__main__': gather_data()