from config import config import os from kafka.kafka import send_kafka_message from queries import queries from connection import mysql, snowflake from stores import stores import sentry_sdk sentry_sdk.init(config.SENTRY_DSN) BATCH_SIZE = 5000 LIMIT = 100000 def backfill(log): offset = os.environ.get("CMM_ID") if not offset: raise Exception("no offset") log.info(f"start backfill with offset {offset}") query = queries.get_contribution_ids(LIMIT, int(offset)) try: cs = snowflake.get_cursor() cs.execute(query) contributions = cs.fetchmany(BATCH_SIZE) contribution_counter = 0 while len(contributions) > 0: for contribution in contributions: id = contribution["ID"] send_kafka_message(id) contribution_counter += 1 if contribution_counter % 100 == 0: log.info(f"contribution - {id} - count - {contribution_counter}") contributions = cs.fetchmany(BATCH_SIZE) except Exception as e: log.error(e) raise e finally: log.info("finished contributions") snowflake.close() def backfillOne(log): log.info("backfill one") try: send_kafka_message("00000cdf-a978-462f-8e32-07c37164ade1") log.info(f"sent 00000cdf-a978-462f-8e32-07c37164ade1") except Exception as e: log.error(e) raise e finally: log.info("finished contributions")