import concurrent from concurrent.futures.thread import ThreadPoolExecutor from typing import Dict, List, Tuple, Type from delphi_api.bigtable import BigTableModel from delphi_api.v3.enums import AggByParam from delphi_api.v3.view_models.generic_dsp import GenericDspViewModel from delphi_api.v3.view_models.params import Params from delphi_api.v3.view_models.query_builder import QueryBuilder class BigTableLoader: """ Static helper utility for making sync/async requests to Bigtable tables, and loading ViewModels with data from one or more tables """ @staticmethod def _get_base_key_name(params: Params): """Used by aggregation functions as the highest level in the tree. See :class:`delphi_api.utils.models.Models` """ if params.agg_by == AggByParam.ISRC.value: return 'isrc' elif params.agg_by == AggByParam.ARTIST.value: return 'artist_id' elif params.track_id: return 'track_id' elif params.isrc: return 'isrc' elif params.playlist_id: return 'playlist_id' elif params.chart_id: return 'chart_id' elif params.artist_id: return 'artist_id' @staticmethod def _load_multiple(view_models: List[GenericDspViewModel], fn: str, async_: bool, **kwargs): if async_: return BigTableLoader._load_multiple_async(view_models=view_models, fn=fn, **kwargs) return BigTableLoader._load_multiple_sync(view_models=view_models, fn=fn, **kwargs) @staticmethod def _load_multiple_sync(view_models: List[GenericDspViewModel], fn: callable, **kwargs): """Make iterative calls to BigTable and merge the results""" results = [] for model in view_models: method = getattr(model, fn) results.extend(method(**kwargs).results) return results @staticmethod def _load_multiple_async(view_models: List[GenericDspViewModel], fn: callable, **kwargs): """Make concurrent calls to BigTable and merge the results""" results = [] if not view_models: return results with ThreadPoolExecutor(max_workers=len(view_models)) as executor: futures = [executor.submit(getattr(model, fn), **kwargs) for model in view_models] for future in concurrent.futures.as_completed(futures): model = future.result() results.extend(model.results) return results @staticmethod def get_multiple(params: Params, data_models: Dict[str, Type[BigTableModel]], item_keys: List[Tuple[bytes, bytes]], async_: bool = True) -> List[dict]: """ Query multiple database tables and consolidate results into a single generic view model. This method calls :meth:`GenericDspViewModel.bigtable_load_results` to query to Bigtable. """ if params.dsp: # query only requested DSPs tables dsps = QueryBuilder.param_as_list(params.dsp) models = [v for k, v in data_models.items() if k in dsps] else: # query all DSP tables models = list(set(data_models.values())) base_key_name = BigTableLoader._get_base_key_name(params) view_models = [ GenericDspViewModel(params, base_key_name=base_key_name, data_model=model) for model in models ] results = BigTableLoader._load_multiple( view_models, fn=GenericDspViewModel.bigtable_load_results.__name__, async_=async_, item_keys=frozenset(item_keys) ) if not results: return [] return GenericDspViewModel(params, base_key_name=base_key_name, results=results).items