import importlib from datetime import datetime from sqlalchemy import and_, exc, or_ from src.config import config as _config from src.constants import DBType from src.logger import logger from src.utils.common import timer from src.utils.db import get_session session = get_session(db_type=DBType[_config.DB_TYPE]) def get_model(model_path: str): path = model_path.split(".") model_name = path.pop() globals().update(importlib.import_module(".".join(path), model_name).__dict__) return globals()[model_name] def clear_table(model_path: str, filter_by: str, date: datetime): """Clear records from specified table before passed date.""" model = get_model(model_path) session.query(model).filter(getattr(model, filter_by) < date).delete(synchronize_session=False) session.commit() @timer(logger) def clear_chunk(model, limit: int, filter_by: str, date: datetime) -> int: """Query and clear chunk of records, returns its count.""" query = session.query(model.id.label("p_id")).filter(getattr(model, filter_by) < date.isoformat()).limit(limit) records = [and_(model.id == r.p_id) for r in query.all()] records_len = len(records) if not records_len: return 0 try: session.query(model).filter(or_(*records)).delete(synchronize_session=False) session.commit() except exc.InternalError: pass return records_len def clear_table_by_chunks(model_path: str, filter_by: str, date: datetime): """Iterable clear records from specified table before passed date.""" model = get_model(model_path) records_len = 1 exec_time = _config.TIME_FRAME time_frame = _config.TIME_FRAME while records_len: limit = int(max((records_len or 1) / (exec_time or 1) * time_frame, 1)) # adaptive chunk size records_len, exec_time = clear_chunk(model, limit, filter_by, date)