import collections import datetime import json import os import click import connection import consts_2 as consts progress_dir = os.environ.get('PROGRESS_DIR', 'progress') def main(): for query_order, (query_name, query_config) in enumerate(consts.MOVE_PROJECT_QUERIES.items()): print(f"Running query {query_name}") processed_ids = read_progress(query_order, query_name) unprocessed_ids = sorted(list(set(consts.project_ids) - set(processed_ids))) print(f"Unprocessed project ids left: {len(unprocessed_ids)}") ids = chunk_list(unprocessed_ids, query_config['batch_size']) with click.progressbar(ids) as project_ids_with_bar: with connection.driver.session() as session: for i, project_ids in enumerate(project_ids_with_bar, start=1): start_time = datetime.datetime.now() session.run(query_config['query'], project_ids=project_ids) processed_ids.extend(project_ids) # if i % 5 == 0: print('\n\n') print(f"{datetime.datetime.now()} Processed {len(processed_ids)} project ids. Batch {query_config['batch_size']} took {(datetime.datetime.now() - start_time).total_seconds()}.") print('\n\n') with open(f'{progress_dir}/{query_order}_{query_name}_progress.json', 'w') as f: json.dump(processed_ids, f) def chunk_list(lst, chunk_size): """Yield successive chunk_size chunks from lst.""" result = [] for i in range(0, len(lst), chunk_size): # yield lst[i:i + chunk_size] result.append(lst[i:i + chunk_size]) return result def read_progress(query_order, query_name): filename = f'{progress_dir}/{query_order}_{query_name}_progress.json' if os.path.exists(filename): with open(filename, 'r') as f: return json.load(f) return [] if __name__ == '__main__': main()