from abc import abstractmethod from typing import Any, Dict, List from app_types import S3Path from constants import AGG_DATA_STRUCTURE from ... import BaseStep from ...mixins import PGMixin, S3Mixin __all__ = ["SearchLoadDataBase"] class SearchLoadDataBase(S3Mixin, PGMixin, BaseStep): priority = 4 @property @abstractmethod def data_path(self) -> S3Path: pass @property @abstractmethod def table_name(self) -> str: pass @property def alt_table_name(self) -> str: # TODO: use self.timestamp instead of "alt" # but first find a way to delete tables by wildcard like "DROP TABLE TABLE_NAME_*" return f"{self.table_name}_alt" table_definition: Dict[str, str] = AGG_DATA_STRUCTURE def process(self) -> Dict[str, Any]: result: Dict[str, List[str]] = {"files": []} self.logger.info(f"Recreating '{self.alt_table_name}'") self.drop_table(self.alt_table_name) self.create_table(self.alt_table_name, self.table_definition) self.logger.info(f"Importing data from '{self.data_path}'") for key in self.get_s3_keys(self.data_path): self.logger.debug(f"Importing '{key}'") self.import_from_s3(self.alt_table_name, self.s3_artifacts_bucket, key) result["files"].append(key) return result