import datetime import time import multiprocessing import traceback import sys import random import boto3 THREADS = 12 cutoff = datetime.datetime.strptime('2020-09-13 00:00:00', '%Y-%m-%d %H:%M:%S') def delete_logs(arguments): stream = arguments[0] creation_day = datetime.datetime.fromtimestamp(stream['creationTime'] / 1000) if creation_day < cutoff: time.sleep(random.random()) logs.delete_log_stream( logGroupName=log_group_name, logStreamName=stream['logStreamName'] ) return stream['logStreamName'], creation_day < cutoff if __name__ == '__main__': logs = boto3.client('logs') log_group_name = '/aws/batch/job' next_token = None while True: if next_token: log_streams = logs.describe_log_streams(logGroupName=log_group_name, nextToken=next_token, limit=50) else: log_streams = logs.describe_log_streams(logGroupName=log_group_name) next_token = log_streams.get('nextToken', None) pool = multiprocessing.Pool(THREADS) iterator = pool.imap_unordered(delete_logs, list(zip(log_streams['logStreams']))) while True: try: log_stream_name, deleted = next(iterator) except multiprocessing.TimeoutError: continue except StopIteration: break except Exception: # Print traceback because we can't reraise it here traceback.print_exc(file=sys.stdout) else: if deleted: print('deleted {} {}'.format(log_stream_name, deleted)) pool.close() pool.join() time.sleep(3) if not next_token or len(log_streams['logStreams']) == 0: break