from abc import abstractmethod from collections import defaultdict from collections.abc import Iterator from typing import Any, Callable, Dict, List, Tuple from constants import AGG_DATA_STRUCTURE, AGG_ID_KEY, AGG_VALUE_KEY from ...base import BaseStep from ...mixins import CSVMixin, S3Mixin from ..constants import RAW_ARTIST_ID_KEY, RAW_TRACK_ID_KEY __all__ = ["AggregatesTransformBase"] class AggregatesTransformBase(S3Mixin, CSVMixin, BaseStep): """ Get raw data from S3, aggregates it, and writes back """ csv_normalizers: Dict[str, Callable[[Any], Any]] = {RAW_ARTIST_ID_KEY: str, RAW_TRACK_ID_KEY: str} @property @abstractmethod def raw_data_path(self) -> str: pass @property @abstractmethod def agg_data_path(self) -> str: pass @property @abstractmethod def raw_id_key(self) -> str: pass agg_id_key: str = AGG_ID_KEY agg_value_key: str = AGG_VALUE_KEY agg_keys: Tuple[str, ...] = tuple(AGG_DATA_STRUCTURE.keys()) def _read_raw_data(self, prefix: str) -> Iterator[Tuple[str, Dict[str, Any]]]: """ Reads multiple files with raw artist data and combines them togather :param prefix: Path to S3 sub-folder to read :return: """ for key in self.get_s3_keys(prefix): for record in self.read_csv(self.get_s3_object(key)): yield record.pop(self.raw_id_key), record def _write_agg_data_batch(self, data: List[Dict[str, Any]], batch_index: int) -> str: """ Write aggregated artists data to S3 :param data: Aggregated data batch to write :param batch_index: Index of the current batch :return: S3 key to written file """ key = f"{self.agg_data_path}/data_{batch_index}.csv" self.logger.debug(f"Writing '{key}'") self.put_s3_object(key, self.write_csv(data, self.agg_keys)) return key def process(self) -> Dict[str, Any]: result: Dict[str, List[str]] = {"processed_ids": [], "files": []} self.wipe_folder(self.agg_data_path) self.logger.info(f"Processing data from '{self.raw_data_path}'") for prefixes in self.get_batched_s3_keys(f"{self.raw_data_path}/"): for index, prefix in enumerate(prefixes): artists_data = defaultdict(list) for artist_id, track in self._read_raw_data(prefix): artists_data[artist_id].append(track) result["processed_ids"].append(artist_id) upload_data = [] for artist_id, tracks in artists_data.items(): upload_data.append({self.agg_id_key: artist_id, self.agg_value_key: self.format_json(tracks)}) f_key = self._write_agg_data_batch(upload_data, index) result["files"].append(f_key) return result