import functools from concurrent.futures.thread import ThreadPoolExecutor from datetime import datetime from queue import Queue import requests from ssavva.artists.backfill_spotify_and_apple.db_connections import \ get_valid_releases, dump_itunes_responses from ssavva.artists.backfill_spotify_and_apple.utils import WorkerPool def _db_worker(q, number): items = [] counter = 0 while True: counter += 1 item = q.get() if item is None: dump_itunes_responses(items) break items.append( (item['release_id'], item['display_upc'], item['upc'], item['response'])) if counter % 10 == 0: print('{} buffer {}'.format(number, len(items))) if counter % 100 == 0: print(f'db dump in {number}', datetime.now()) dump_itunes_responses(items) items = [] q.task_done() def request_itunes_by_upc(item, queue): params = {'upc': item['display_upc']} try: response = requests.get('https://itunes.apple.com/lookup', params=params) resposne_data = response.text if response.text else None item['response'] = resposne_data queue.put(item) except Exception as e: print(e) def gather_data(): releases = get_valid_releases() q = Queue(maxsize=500) db_worker_pool = WorkerPool(4, _db_worker, (q,)) with ThreadPoolExecutor(max_workers=40) as executor: executor.map( functools.partial(request_itunes_by_upc, queue=q), releases) q.join() q.put(None) db_worker_pool.add_stoppers(q) db_worker_pool.join_workers() if __name__ == '__main__': gather_data()