from typing import Any, Dict from ...base import BaseStep from ...mixins import SnowflakeMixin from ...utils import Cleanup from ..constants import NON_SIGNED_TRACK_COLUMNS, OK_STATUS, VALIDATION_STATUS_NON_SIGNED_TRACK __all__ = ["Validation"] class Validation(SnowflakeMixin, BaseStep): depends_on = {Cleanup} def _check_exist_data_non_signed_tracks(self) -> str: """ Do we have any data in table T_NON_SIGNED_TRACK? """ status = OK_STATUS query = """ SELECT COUNT(*) AS COUNTER FROM DNA.DNA_PUBLIC.T_NON_SIGNED_TRACK """ result_query = self.run_raw_query(query=query, is_async=False) if result_query[0].get("counter") == 0: # type: ignore status = VALIDATION_STATUS_NON_SIGNED_TRACK["non_signed_tracks"] return status def _check_columns_non_signed_tracks(self): """ Check if table T_NON_SIGNED_TRACK has correct columns """ status = OK_STATUS query = """ SHOW COLUMNS IN TABLE DNA.DNA_PUBLIC.T_NON_SIGNED_TRACK """ res = self.run_raw_query(query=query, is_async=False) columns = set([i.get("column_name") for i in res]) diff_columns = list(NON_SIGNED_TRACK_COLUMNS ^ columns) if diff_columns: status = VALIDATION_STATUS_NON_SIGNED_TRACK["non_signed_tracks_columns"] status = status + ":" + ",".join(str(x) for x in diff_columns) return status def _check_id_is_chartmetric(self): """ Track ID column - is chartmetric value """ status = OK_STATUS query = """ SELECT T_NON_SIGNED_TRACK.ID AS ID FROM DNA.DNA_PUBLIC.T_NON_SIGNED_TRACK AS T_NON_SIGNED_TRACK WHERE T_NON_SIGNED_TRACK.ID NOT IN ( SELECT V_CM_TRACK.ID AS ID FROM DELPHI_EXPLORATION.CHARTMETRIC.V_CM_TRACK AS V_CM_TRACK ) """ res = self.run_raw_query(query=query, is_async=False) res = [i.get("id") for i in res] # type: ignore if res: status = VALIDATION_STATUS_NON_SIGNED_TRACK["id_is_chartmetric"] status = status + ":" + ",".join([str(x) for x in res]) return status def _check_track_id_is_chartmetric(self): """ ARTIST_ID column - is chartmetric value """ status = OK_STATUS query = """ SELECT ARTIST_ID FROM DNA.DNA_PUBLIC.T_NON_SIGNED_TRACK AS T_NON_SIGNED_TRACK WHERE ARTIST_ID NOT IN ( SELECT ID FROM DELPHI_EXPLORATION.CHARTMETRIC.V_CM_ARTIST AS V_CM_ARTIST ) """ res = self.run_raw_query(query=query, is_async=False) res = [i.get("artist_id") for i in res] # type: ignore if res: status = VALIDATION_STATUS_NON_SIGNED_TRACK["track_id_is_chartmetric"] status = status + ":" + ",".join([str(x) for x in res]) return status def _check_one_track_one_artist(self): """ Regarding current business logic, one track has one artist ( should change in feature) """ status = OK_STATUS query = """ SELECT T_NON_SIGNED_TRACK.ID, T_NON_SIGNED_TRACK.ARTIST_ID, COUNT(T_NON_SIGNED_TRACK.ARTIST_ID) AS COUNTER FROM DNA.DNA_PUBLIC.T_NON_SIGNED_TRACK AS T_NON_SIGNED_TRACK GROUP BY T_NON_SIGNED_TRACK.ID, T_NON_SIGNED_TRACK.ARTIST_ID HAVING COUNTER > 1 """ res = self.run_raw_query(query=query, is_async=False) res = [i.get("id") for i in res] # type: ignore if res: status = VALIDATION_STATUS_NON_SIGNED_TRACK["one_artist_one_track"] status = status + ":" + ",".join([str(x) for x in res]) return status def _check_track_has_own_isrc(self): """ Each track should have own ISRC """ status = OK_STATUS query = """ SELECT T_NON_SIGNED_TRACK.ID, T_NON_SIGNED_TRACK.ISRC, COUNT(T_NON_SIGNED_TRACK.ARTIST_ID) AS COUNTER FROM DNA.DNA_PUBLIC.T_NON_SIGNED_TRACK AS T_NON_SIGNED_TRACK GROUP BY T_NON_SIGNED_TRACK.ID, T_NON_SIGNED_TRACK.ISRC HAVING COUNTER > 1 """ res = self.run_raw_query(query=query, is_async=False) res = [i.get("id") for i in res] # type: ignore if res: status = VALIDATION_STATUS_NON_SIGNED_TRACK["track_own_isrc"] status = status + ":" + ",".join([str(x) for x in res]) return status def _check_not_null_columns(self): """ Columns ID, ARTIST_ID, ISRC, NAME - don't have NULL inside """ status = OK_STATUS query = """ SELECT T_NON_SIGNED_TRACK.ID AS ID, T_NON_SIGNED_TRACK.ISRC AS ISRC, T_NON_SIGNED_TRACK.ARTIST_ID AS ARTIST_ID, T_NON_SIGNED_TRACK.NAME AS NAME FROM DNA.DNA_PUBLIC.T_NON_SIGNED_TRACK AS T_NON_SIGNED_TRACK WHERE T_NON_SIGNED_TRACK.ID IS NULL OR T_NON_SIGNED_TRACK.ISRC IS NULL OR T_NON_SIGNED_TRACK.ARTIST_ID IS NULL OR T_NON_SIGNED_TRACK.NAME IS NULL """ res = self.run_raw_query(query=query, is_async=False) res = [i.get("id") for i in res] # type: ignore if res: status = VALIDATION_STATUS_NON_SIGNED_TRACK["not_null_columns"] status = status + ":" + ",".join([str(x) for x in res]) return status def process(self) -> Dict[str, Any]: result = { "non_signed_tracks": self._check_exist_data_non_signed_tracks(), "non_signed_tracks_columns": self._check_columns_non_signed_tracks(), "id_is_chartmetric": self._check_id_is_chartmetric(), "track_id_is_chartmetric": self._check_track_id_is_chartmetric(), "one_artist_one_track": self._check_one_track_one_artist(), "track_own_isrc": self._check_track_has_own_isrc(), "not_null_columns": self._check_not_null_columns(), } if all(x == OK_STATUS for x in result.values()): self.logger.info(result) else: self.logger.error(result) return result