"""Sync manufactured upcs to snowflake processor.""" from csv import DictReader from string import Template import typing from snowflake_connector.etl_connector import SnowflakeSQLExecutor from config import app_logger as logger, SNOWFLAKE_CONFIG from src.connectors.ows_royalties import get_current_statement_period from src.models import StatementPeriod from src.processors.base import Processor from src.processors.exceptions import ProcessingError from .queries import CLEANUP_MANUFACTURED_UPCS, INSERT_MANUFACTURED_UPCS class ManufacturedUpcsProcessor(Processor): """ Processor to sync manufactured upcs to snowflake. Expects the CSV file with the following columns: - upc (str) The first row is a header. """ _data: typing.List[str] _statement_period: StatementPeriod def process(self) -> None: """Run main logic.""" self._load_csv_data() self._load_ows_data() self._cleanup_data() self._insert_data() def _load_csv_data(self) -> None: """Load data from file.""" logger.info('Loading data from csv file') try: raw_data = list(self.csv_dict_reader) except Exception as e: raise ProcessingError(f'Unable to read the file: {e}') try: upc_data = [] for row in raw_data: upc_data.append(row['UPC']) self._data = upc_data except Exception as e: raise ProcessingError(f'Invalid file content: {e}') if not self._data: raise ProcessingError('No items in file to post') logger.info(f'Loaded {len(self._data)} items.') def _load_ows_data(self) -> None: """Fetch any supplemental royalty accounting data.""" self._statement_period = get_current_statement_period() if not self._statement_period: raise ProcessingError('No current statement period') def _cleanup_data(self) -> None: """Soft delete any manufactured upcs from this statement period.""" sql_query = Template(CLEANUP_MANUFACTURED_UPCS).substitute( **{ 'db': SNOWFLAKE_CONFIG.get('db'), 'schema': SNOWFLAKE_CONFIG.get('schema'), 'statement_period_id': self._statement_period.statement_period_id, } ) with SnowflakeSQLExecutor(SNOWFLAKE_CONFIG) as sf_executor: sf_executor.execute(sql_query) def _insert_data(self) -> None: """Insert upc data into snowflake table.""" sql_data = [] period_id = self._statement_period.statement_period_id lambda_user = 'lambda-documents-load-from-s3' for upc in self._data: sql_data.append(f"({period_id}, {upc}, CURRENT_DATE(), '{lambda_user}')") sql_query = Template(INSERT_MANUFACTURED_UPCS).substitute( **{ 'db': SNOWFLAKE_CONFIG.get('db'), 'schema': SNOWFLAKE_CONFIG.get('schema'), 'values': ', '.join(sql_data), } ) with SnowflakeSQLExecutor(SNOWFLAKE_CONFIG) as sf_executor: sf_executor.execute(sql_query)