from service.tasks.data_model import list_caw_schemas from service.utils.aws_connectors import pd_read_sql from service.utils.data_model_utils import get_management_schema from .audience import Audience class AllianceAudience(Audience): """ Class created to add support the Alliance logic Acts like a normal audience if alliance_schema = None Allows to generate data from collection created in Alliance schema. - Only returns the fans that they have in their own workspace(s) - For each fan only returns data points that are available in their own workspace(s) """ def __init__(self, collection_id, schema: str, alliance_schema: str = None): super().__init__(collection_id=collection_id, schema=schema) self.alliance_schema = alliance_schema self.management_schema = get_management_schema(schema) self.related_workspaces = [ x for x in list_caw_schemas( self.management_schema, workspaces=True, alliances=False, with_caw_type=False, ) ] def _generate_available_data(self): if self.alliance_schema is None: return super()._generate_available_data() else: attributes_list = [str(x) for x in self.source_attribute_ids] attributes_dynamic_sql = [] for attribute in attributes_list: # We use min so in case it has female vs unknown, it would take female. query_str = f', min(CASE WHEN fa.attribute_id = {attribute} THEN value END) "{attribute}"' attributes_dynamic_sql.append(query_str) source_dynamic_sql = [] for w_schema in self.related_workspaces: or_str = f"source like '{w_schema}%%'" source_dynamic_sql.append(or_str) params = {"collection_id": self.base_collection_id} sql = f""" WITH owned_collection_ids AS ( SELECT DISTINCT id FROM {self.alliance_schema}.collection WHERE {' OR '.join(source_dynamic_sql)} ), owned_fan_ids AS ( SELECT DISTINCT fan_id FROM {self.alliance_schema}.collection_fan WHERE collection_id IN (SELECT * FROM owned_collection_ids) ), base_fan_ids AS ( SELECT DISTINCT fan_id FROM {self.alliance_schema}.collection_fan WHERE collection_id = %(collection_id)s AND fan_id IN (SELECT * from owned_fan_ids) ), attribute_data AS ( SELECT fa.fan_id as fan_id_2 {''.join(attributes_dynamic_sql)} FROM {self.alliance_schema}.fan_attribute fa WHERE fan_id IN (SELECT * FROM base_fan_ids) AND collection_id IN (SELECT * FROM owned_collection_ids) GROUP BY 1 ) SELECT * FROM base_fan_ids bfi LEFT JOIN attribute_data ad ON bfi.fan_id = ad.fan_id_2; """ data = pd_read_sql(sql, params=params) data.drop(columns=["fan_id", "fan_id_2"], inplace=True) return data def get_general_data(self): """Getting general information about base_collection_id""" params = {"collection_id": self.base_collection_id} sql = f""" SELECT * FROM {self.alliance_schema or self.schema}.collection WHERE id = %(collection_id)s; """ data = pd_read_sql(sql, params=params) return data