using MySql.Data.MySqlClient; using System; using System.Collections.Generic; using System.Globalization; using System.Linq; using System.Threading.Tasks; using PetaPoco.Business; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Contracts.Entities; using Sony.Filtr.Database; using Sony.Filtr.DistributedCaching; using Sony.Filtr.SpotifyAnalytics.Data; using Sony.Filtr.SpotifyAnalytics.Models; using Sony.Filtr.Utility.Extensions; using Dapper; using System.Text; using Sony.Filtr.Utility; using System.Collections.Concurrent; using TaskExtensions = Sony.Filtr.Utility.Extensions.TaskExtensions; namespace Sony.Filtr.SpotifyAnalytics { public class SpotifyAnalyticsManager { private readonly DistributedCacheHandler _distributedCacheHandler; private readonly SpotifyAnalyticsFactory _factory; public SpotifyAnalyticsManager(DistributedCacheHandler distributedCacheHandler, SpotifyAnalyticsFactory spotifyAnalyticsFactory) { _distributedCacheHandler = distributedCacheHandler; _factory = spotifyAnalyticsFactory; } public async Task GetLatestDayWithDataAsyncOld() { DateTime latestDate = DateTime.MinValue; using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { MySqlCommand cmd = new MySqlCommand("SELECT Max(Date) FROM tblSpotifyAnalyticsAccountDate WHERE LogType = @logtype", conn); cmd.Parameters.AddWithValue("@logtype", AnalyticsLogType.Finished); var reader = await cmd.ExecuteReaderAsync(); if (await reader.ReadAsync()) { latestDate = reader.GetDateTime(0); } } return latestDate; } public async Task> GetFilesProcessedCountForDates(IEnumerable dates, IEnumerable supportedFileTypes, int version) { Dictionary result = new Dictionary(); using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { MySqlCommand cmd = new MySqlCommand($@" SELECT Date, count(*) as Count FROM tblSpotifyAnalyticsDate WHERE Version = @version AND Date IN ({Maybe.ToCommaSeparated(dates.Select(d => d.ToString("yyyy-MM-dd")))}) AND FileType IN({String.Join(", ", supportedFileTypes.Select(i => (int)i))}) GROUP BY Date;", conn); cmd.Parameters.AddWithValue("@version", version); using (var reader = await cmd.ExecuteReaderAsync()) { while (await reader.ReadAsync()) { result.Add(reader.GetDateTime(0), reader.GetInt32(1)); } } } return result; } public async Task GetLatestDayWithDataAsyncNew() { DateTime latestDate = DateTime.MinValue; using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { MySqlCommand cmd = new MySqlCommand("SELECT MAX(Date) FROM tblSpotifyAnalyticsDate WHERE FileType = @fileType AND Version = 2", conn); cmd.Parameters.AddWithValue("@fileType", SpotifyS3FileType.Playlists); var reader = await cmd.ExecuteReaderAsync(); if (await reader.ReadAsync()) { latestDate = reader.GetDateTime(0); } } return latestDate; } public async Task GetLatestTimestampWithDataAsync() { DateTime latestDate = DateTime.MinValue; using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { MySqlCommand cmd = new MySqlCommand("SELECT MAX(TimeStamp) FROM tblSpotifyAnalyticsDate WHERE FileType = @fileType AND Version = 2", conn); cmd.Parameters.AddWithValue("@fileType", SpotifyS3FileType.Playlists); var reader = await cmd.ExecuteReaderAsync(); if (await reader.ReadAsync()) { latestDate = reader.GetDateTime(0); } } return latestDate; } public List GetStreamSummary(SpotifyAnalyticsAccount account, DateTime startDate, DateTime endDate) { return _factory.GetStreamSummary((int)account, startDate, endDate); } public List GetStreamSummary(DateTime startDate, DateTime endDate) { return _factory.GetStreamSummary(startDate, endDate); } public async Task SaveStreamSummaryAsync(int accountId, StreamDataSummaryForAnalysis summary) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { const string sqlString = "INSERT INTO tblSpotifyAnalyticsAccountStreamInfo VALUES (@AccountId, @Market, @Date, @TotalStreams, @DeviceDesktopStreams, @DeviceTabletStreams, @DeviceMobileStreams, " + "@OSAndroidStreams, @OSiOSStreams, @OSWindowsStreams, @OSMacoSStreams, @OSOtherStreams, @SourcePlaylist, @SourceSearch, @SourceArtist, @SourceAlbum, " + "@SourceCollection, @SourceOther, @StreamsFromSonyPlaylists, @StreamsFromLocalSonyPlaylists, @UniqueUsers, " + "@UniqueUsersSourcePlaylist, @UniqueUsersSourceSearch, @UniqueUsersSourceArtist, @UniqueUsersSourceAlbum, @UniqueUsersSourceCollection, @UniqueUsersSourceOther, " + "@UniqueUsersDeviceDesktop, @UniqueUsersDeviceMobile, @UniqueUsersDeviceTablet) " + "ON DUPLICATE KEY UPDATE TotalStreams = @TotalStreams, DeviceDesktopStreams = @DeviceDesktopStreams, DeviceTabletStreams = @DeviceTabletStreams, " + "DeviceMobileStreams = @DeviceMobileStreams, OSAndroidStreams = @OSAndroidStreams, OSiOSStreams = @OSiOSStreams, OSWindowsStreams = @OSWindowsStreams, " + "OSMacoSStreams = @OSMacoSStreams, OSOtherStreams = @OSOtherStreams, SourcePlaylist = @SourcePlaylist, SourceSearch = @SourceSearch, SourceArtist = @SourceArtist, " + "SourceAlbum = @SourceAlbum, SourceCollection = @SourceCollection, SourceOther = @SourceOther, StreamsFromSonyPlaylists = @StreamsFromSonyPlaylists, " + "StreamsFromLocalSonyPlaylists = @StreamsFromLocalSonyPlaylists, UniqueUsers = @UniqueUsers, " + "UniqueUsersSourcePlaylist = @UniqueUsersSourcePlaylist, UniqueUsersSourceSearch = @UniqueUsersSourceSearch, UniqueUsersSourceArtist = @UniqueUsersSourceArtist, " + "UniqueUsersSourceAlbum = @UniqueUsersSourceAlbum, UniqueUsersSourceCollection = @UniqueUsersSourceCollection, UniqueUsersSourceOther = @UniqueUsersSourceOther, " + "UniqueUsersDeviceDesktop = @UniqueUsersDeviceDesktop, UniqueUsersDeviceMobile = @UniqueUsersDeviceMobile, UniqueUsersDeviceTablet = @UniqueUsersDeviceTablet"; using(var command = new MySqlCommand(sqlString, conn)) { command.Parameters.AddWithValue("@AccountId", accountId); command.Parameters.AddWithValue("@Market", summary.Market); command.Parameters.AddWithValue("@Date", summary.Date); command.Parameters.AddWithValue("@TotalStreams", summary.TotalStreams); command.Parameters.AddWithValue("@DeviceDesktopStreams", summary.DeviceDesktopStreams); command.Parameters.AddWithValue("@DeviceTabletStreams", summary.DeviceTabletStreams); command.Parameters.AddWithValue("@DeviceMobileStreams", summary.DeviceMobileStreams); command.Parameters.AddWithValue("@OSAndroidStreams", summary.OSAndroidStreams); command.Parameters.AddWithValue("@OSiOSStreams", summary.OSiOSStreams); command.Parameters.AddWithValue("@OSWindowsStreams", summary.OSWindowsStreams); command.Parameters.AddWithValue("@OSMacoSStreams", summary.OSMacOSStreams); command.Parameters.AddWithValue("@OSOtherStreams", summary.OSOtherStreams); command.Parameters.AddWithValue("@SourcePlaylist", summary.SourcePlaylist); command.Parameters.AddWithValue("@SourceSearch", summary.SourceSearch); command.Parameters.AddWithValue("@SourceArtist", summary.SourceArtist); command.Parameters.AddWithValue("@SourceAlbum", summary.SourceAlbum); command.Parameters.AddWithValue("@SourceCollection", summary.SourceCollection); command.Parameters.AddWithValue("@SourceOther", summary.SourceOther); command.Parameters.AddWithValue("@StreamsFromSonyPlaylists", summary.StreamsFromSonyPlaylists); command.Parameters.AddWithValue("@StreamsFromLocalSonyPlaylists", summary.StreamsFromLocalSonyPlaylists); command.Parameters.AddWithValue("@UniqueUsers", summary.UniqueUsers); command.Parameters.AddWithValue("@UniqueUsersSourcePlaylist", summary.UniqueUsersSourcePlaylist); command.Parameters.AddWithValue("@UniqueUsersSourceSearch", summary.UniqueUsersSourceSearch); command.Parameters.AddWithValue("@UniqueUsersSourceArtist", summary.UniqueUsersSourceArtist); command.Parameters.AddWithValue("@UniqueUsersSourceAlbum", summary.UniqueUsersSourceAlbum); command.Parameters.AddWithValue("@UniqueUsersSourceCollection", summary.UniqueUsersSourceCollection); command.Parameters.AddWithValue("@UniqueUsersSourceOther", summary.UniqueUsersSourceOther); command.Parameters.AddWithValue("@UniqueUsersDeviceDesktop", summary.UniqueUsersDeviceDesktop); command.Parameters.AddWithValue("@UniqueUsersDeviceMobile", summary.UniqueUsersDeviceMobile); command.Parameters.AddWithValue("@UniqueUsersDeviceTablet", summary.UniqueUsersDeviceTablet); await command.ExecuteNonQueryAsync(); } } } public async Task SetSpotifyAnalyticsLogAsync(int accountId, string market, AnalyticsLogType logType, DateTime analyticsDate, int version) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { const string sqlString = "INSERT INTO tblSpotifyAnalyticsAccountDate (AccountId, Market, Date, LogType, Version) VALUES(@accountId, @market, @analyticsDate, @logType, @version) ON DUPLICATE KEY UPDATE Timestamp = CURRENT_TIMESTAMP"; var command = new MySqlCommand(sqlString, conn); command.Parameters.AddWithValue("@accountId", accountId); command.Parameters.AddWithValue("@market", market); command.Parameters.AddWithValue("@analyticsDate", analyticsDate); command.Parameters.AddWithValue("@logType", logType); command.Parameters.AddWithValue("@version", version); await command.ExecuteNonQueryAsync(); } } public async Task> GetSpotifyAnalyticsLogAsync(int accountId, DateTime analyticsDate, AnalyticsLogType logType, int version) { var dates = new Dictionary(); using (var conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { const string sqlString = "SELECT Market, Date FROM tblSpotifyAnalyticsAccountDate WHERE AccountId=@accountId AND Date=@date AND LogType=@logType AND Version=@version ORDER BY Timestamp DESC"; var command = new MySqlCommand(sqlString, conn); command.Parameters.AddWithValue("@accountId", accountId); command.Parameters.AddWithValue("@date", analyticsDate); command.Parameters.AddWithValue("@logType", logType); command.Parameters.AddWithValue("@version", version); var reader = await command.ExecuteReaderAsync(); while (await reader.ReadAsync()) { dates.Add(reader.GetString(0), reader.GetDateTime(1)); } reader.Close(); } return dates; } public async Task> GetSpotifyAnalyticsLogAsync(AnalyticsLogType logType, int version) { return PetaPocoRepository.ReadOnlyInstance.Fetch("WHERE logType=@0 AND Version = @1", logType, version); } public async Task> GetPlaylistStreamDataAsync(string playlistId, string legacyPlaylistUri, DateTime startDate, DateTime endDate) { return await _factory.GetPlaylistStreamDataAsync(playlistId, legacyPlaylistUri, null, startDate, endDate); } public async Task> GetPlaylistStreamDataAsync(string playlistId, string legacyPlaylistUri, SpotifyAnalyticsAccount account, DateTime startDate, DateTime endDate) { return await _factory.GetPlaylistStreamDataAsync(playlistId, legacyPlaylistUri, (int)account, startDate, endDate); } public async Task> GetPlaylistsStreamSummaryAsync(List playlistIds) { var cachedStreamSummaries = await GetCachedPlaylistStreamSummariesAsync(playlistIds); var nonCachedStreamSummaries = playlistIds.Except(cachedStreamSummaries.Select(p => p.PlaylistId)).ToList(); if (nonCachedStreamSummaries.Any()) { var dbResults = await _factory.GetPlaylistsStreamSummaryAsync(nonCachedStreamSummaries); await CachePlaylistStreamSummariesAsync(dbResults); return cachedStreamSummaries.Union(dbResults).ToList(); } return cachedStreamSummaries; } private async Task CachePlaylistStreamSummariesAsync(List dbResults) { foreach(var playlistStreamSummary in dbResults) { var cacheKey = GetPlaylistStreamSummaryCacheKey(playlistStreamSummary.PlaylistId); _distributedCacheHandler.Put(cacheKey, playlistStreamSummary); } } private async Task> GetCachedPlaylistStreamSummariesAsync(List playlistIds) { List playlistStreamsSummaries = new List(); foreach(var playlist in playlistIds) { var cacheKey = GetPlaylistStreamSummaryCacheKey(playlist); var playlistResult = _distributedCacheHandler.Get(cacheKey) as List; if (playlistResult != null) playlistStreamsSummaries.AddRange(playlistResult); } return playlistStreamsSummaries; } public async Task> GetMarketPlaylistStreamsAsync(DateTime date, int buzzCategoryId) { return await _factory.GetMarketStreamDataAsync(date, buzzCategoryId); } public async Task> GetMarketPlaylistStreamsAsync(string market, DateTime startDate, DateTime endDate, int buzzCategoryId) { return await _factory.GetMarketStreamDataAsync(market, startDate, endDate, buzzCategoryId); } public async Task> GetMarketPlaylistStreamsAsync(SpotifyAnalyticsAccount spotifyAnalyticsAccount, string market, DateTime startDate, DateTime endDate, int buzzCategoryId) { return await _factory.GetMarketStreamDataAsync((int)spotifyAnalyticsAccount, market, startDate, endDate, buzzCategoryId); } public async Task> CalculatePlaylistStreamSummaryAsync(int spotifyAnalyticsAccountId, string playlistId, string playlistUri, DateTime latestDate) { return await _factory.CalculatePlaylistStreamSummaryAsync(spotifyAnalyticsAccountId, playlistId, playlistUri, latestDate); } public async Task> GetTrackDemographicsAsync(string isrc, DateTime startDate, DateTime endDate, string countryCode) { using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { // This also prevent from sql-injections isrc = Maybe.ToCommaSeparated(isrc); var sbSql = new StringBuilder(); // This select statement returns only required fields and do grouping, sorting // on SQL server side to prevent memory leaks and CPU consumption. sbSql.Append("SELECT " + " d.Date, " + " d.ReportCountryCode as CountryCode," + " d.Isrc," + " SUM(d.GenderFemaleAge_0_17) as GenderFemaleAge0_17," + " SUM(d.GenderFemaleAge_18_22) as GenderFemaleAge18_22, " + " SUM(d.GenderFemaleAge_23_27) as GenderFemaleAge23_27, " + " SUM(d.GenderFemaleAge_28_34) as GenderFemaleAge28_34, " + " SUM(d.GenderFemaleAge_35_44) as GenderFemaleAge35_44, " + " SUM(d.GenderFemaleAge_45_59) as GenderFemaleAge45_59, " + " SUM(d.GenderFemaleAge_60_150) as GenderFemaleAge60_150," + " SUM(d.GenderFemaleAge_unknown) as GenderFemaleAgeUnknown," + " SUM(d.GenderMaleAge_0_17) as GenderMaleAge0_17," + " SUM(d.GenderMaleAge_18_22) as GenderMaleAge18_22," + " SUM(d.GenderMaleAge_23_27) as GenderMaleAge23_27," + " SUM(d.GenderMaleAge_28_34) as GenderMaleAge28_34," + " SUM(d.GenderMaleAge_35_44) as GenderMaleAge35_44," + " SUM(d.GenderMaleAge_45_59) as GenderMaleAge45_59," + " SUM(d.GenderMaleAge_60_150) as GenderMaleAge60_150," + " SUM(d.GenderMaleAge_unknown) as GenderMaleAgeUnknown," + " SUM(d.GenderNeutralAge_0_17) as GenderNeutralAge0_17," + " SUM(d.GenderNeutralAge_18_22) as GenderNeutralAge18_22," + " SUM(d.GenderNeutralAge_23_27) as GenderNeutralAge23_27," + " SUM(d.GenderNeutralAge_28_34) as GenderNeutralAge28_34," + " SUM(d.GenderNeutralAge_35_44) as GenderNeutralAge35_44," + " SUM(d.GenderNeutralAge_45_59) as GenderNeutralAge45_59," + " SUM(d.GenderNeutralAge_60_150) as GenderNeutralAge60_150," + " SUM(d.GenderNeutralAge_unknown) as GenderNeutralAgeUnknown," + " SUM(d.GenderUnknownAge_0_17) as GenderUnknownAge0_17," + " SUM(d.GenderUnknownAge_18_22) as GenderUnknownAge18_22," + " SUM(d.GenderUnknownAge_23_27) as GenderUnknownAge23_27," + " SUM(d.GenderUnknownAge_28_34) as GenderUnknownAge28_34," + " SUM(d.GenderUnknownAge_35_44) as GenderUnknownAge35_44," + " SUM(d.GenderUnknownAge_45_59) as GenderUnknownAge45_59," + " SUM(d.GenderUnknownAge_60_150) as GenderUnknownAge60_150," + " SUM(d.GenderUnknownAge_unknown) as GenderUnknownAgeUnknown " + "FROM tblSpotifyTrackStreamDemographics AS d "); sbSql.Append($"WHERE isrc IN ({isrc}) AND Date >= 'startDate' AND Date <= 'endDate' "); // In most cases we do not send country code so we do not need to inline this in sql query if (!string.IsNullOrWhiteSpace(countryCode)) sbSql.Append($"AND countryCode = '{Maybe.ToSingleString(countryCode)}' "); sbSql.Append("GROUP BY Date, CountryCode, Isrc;"); var days = endDate.Subtract(startDate).TotalDays; // 14 best fit to performance if (days <= 14) return await ExeRequestForDateRangeAsync(startDate, endDate, sbSql); var tasks = new List(); var result = new ConcurrentDictionary>(); var dateRanges = Maybe.ToDateRange(startDate, endDate); foreach (var dateRange in dateRanges) { var date = dateRange; tasks.Add(Task.Run(async () => { IEnumerable items = await ExeRequestForDateRangeAsync(dateRange.From, dateRange.To, sbSql); result.TryAdd(date.From.Ticks, items); })); } await TaskExtensions.EnsureAllSuccessAsync(tasks); return result.OrderBy(x => x.Key).SelectMany(x => x.Value).ToArray(); } } public async Task> GetTrackStreamsAsync(List isrcs, DateTime startDate, DateTime endDate, string countryCode, SpotifyAnalyticsAccount? licensor = null) { var licensorId = (int?)licensor; var sql = "SELECT DistributerId, Date, CountryCode, Isrc, ReportDate, " + "SUM(ReportCountryCode) AS ReportCountryCode, " + "SUM(DeviceTypePersonalComputer) AS DeviceTypePersonalComputer, " + "SUM(DeviceTypeCellPhone) AS DeviceTypeCellPhone, " + "SUM(DeviceTypeTablet) AS DeviceTypeTablet, " + "SUM(DeviceTypeGamingConsole) AS DeviceTypeGamingConsole, " + "SUM(DeviceTypeSmartTvDevice) AS DeviceTypeSmartTvDevice, " + "SUM(DeviceTypeConnectedAudioDevice) AS DeviceTypeConnectedAudioDevice, " + "SUM(DeviceTypeBuiltinCarApplication) AS DeviceTypeBuiltinCarApplication, " + "SUM(DeviceTypeOther) AS DeviceTypeOther, " + "SUM(DeviceOsAndroid) AS DeviceOsAndroid, " + "SUM(DeviceOsBlackberry) AS DeviceOsBlackberry, " + "SUM(DeviceOsBrowser) AS DeviceOsBrowser, " + "SUM(DeviceOsIos) AS DeviceOsIos, " + "SUM(DeviceOsLinux) AS DeviceOsLinux, " + "SUM(DeviceOsMac) AS DeviceOsMac, " + "SUM(DeviceOsOther) AS DeviceOsOther, " + "SUM(DeviceOsWindows) AS DeviceOsWindows, " + "SUM(RepeatPlay) AS RepeatPlay, " + "SUM(ShufflePlay) AS ShufflePlay, " + "SUM(SourceAlbum) AS SourceAlbum, " + "SUM(SourceArtist) AS SourceArtist, " + "SUM(SourceChart) AS SourceChart, " + "SUM(SourceCollection) AS SourceCollection, " + "SUM(SourceOther) AS SourceOther, " + "SUM(SourceOthersPlaylist) AS SourceOthersPlaylist, " + "SUM(SourceRadio) AS SourceRadio, " + "SUM(SourceSearch) AS SourceSearch, " + "SUM(Streams) AS Streams " + "FROM tblSpotifyTrackStreams " + "WHERE (@countryCode IS NULL OR CountryCode = @countryCode) " + "AND isrc IN(@isrcs) " + "AND (@licensorId IS NULL OR DistributerId = @licensorId) " + "AND Date >= @startDate AND Date <= @endDate " + "GROUP BY isrc, countryCode"; return PetaPocoRepository.ReadOnlyInstance.Fetch(sql, new { countryCode = countryCode, isrcs = isrcs, licensorId = licensorId, startDate = startDate, endDate = endDate }); } public async Task CalculatePlaylistCategoryStreams(int spotifyAnalyticsAccountId, DateTime streamsDate, Application app, List appRegions, StaticBuzzCategory buzzCategory) { var categoryStreams = new PlaylistCategoryStreams() { AccountId = spotifyAnalyticsAccountId, Market = app.SpotifyRegionCode, BuzzCategoryId = (int)buzzCategory, Date = streamsDate, TotalListStreams = 0, LocalListStreams = 0, TotalListListeners = 0, LocalListListeners = 0, }; foreach (var region in appRegions) { using (var conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { const string sql = @" SELECT p.BuzzCategoryId AS StreamCategory, s.date AS StreamDate, SUM(s.TotalStreams) AS TotalListStreams, SUM(IF(p.CountryCode = @market, s.TotalStreams, 0)) AS LocalListStreams, SUM(s.UniqueUsers) AS TotalListListeners, SUM(IF(p.CountryCode = @market, s.UniqueUsers, 0)) AS LocalListListeners FROM tblSpotifyPlaylist AS p INNER JOIN tblSpotifyAnalyticsAccountPlaylistStreamInfo AS s ON p.playlistId = s.playlistUri AND s.Date = @date AND s.accountId = @accountId LEFT JOIN tblPlaylistIgnoredPlaylist AS ip ON ip.playlistId=p.playlistUri AND ip.musicServiceId=@spotifyMusicServiceId WHERE p.removed = 0 AND ip.playlistID IS NULL AND p.BuzzCategoryId = @buzzCategoryId AND (@market = '_gl' OR s.Market = @region) GROUP BY s.date"; MySqlCommand cmd = new MySqlCommand(sql, conn); cmd.Parameters.AddWithValue("@market", app.SpotifyRegionCode); cmd.Parameters.AddWithValue("@accountId", spotifyAnalyticsAccountId); cmd.Parameters.AddWithValue("@region", region); cmd.Parameters.AddWithValue("@date", streamsDate); cmd.Parameters.AddWithValue("@buzzCategoryId", buzzCategory); cmd.Parameters.AddWithValue("@spotifyMusicServiceId", MusicService.Spotify); var reader = await cmd.ExecuteReaderAsync(); while (await reader.ReadAsync()) { categoryStreams.TotalListStreams += reader.GetIntOrFallback("TotalListStreams", 0); categoryStreams.LocalListStreams += reader.GetIntOrFallback("LocalListStreams", 0); categoryStreams.TotalListListeners += reader.GetIntOrFallback("TotalListListeners", 0); categoryStreams.LocalListListeners += reader.GetIntOrFallback("LocalListListeners", 0); } } } return categoryStreams; } private string GetPlaylistStreamSummaryCacheKey(string playlistId) { return "SpotifyAnalytics_PlaylistStreamSummary_2.0_" + playlistId; } public async Task ClearPlaylistStreamSummaryCacheAsync(IEnumerable playlistIds) { foreach(var playlistId in playlistIds) { var cacheKey = GetPlaylistStreamSummaryCacheKey(playlistId); _distributedCacheHandler.Remove(cacheKey); } } public async Task> GetTrackListenersAsync(List isrcs, DateTime startDate, DateTime endDate, string countryCode) { var distributerId = (int)SpotifyAnalyticsAccount.Sony; //We only get data this distributer so we specify it so no behaviour changes. return PetaPocoRepository.ReadOnlyInstance.Fetch("WHERE isrc IN (@isrcs) AND DistributerId = @distributerId AND date >= @startDate AND date <= @endDate AND (@countryCode IS NULL OR countryCode = @countryCode)", new { isrcs = isrcs, distributerId = distributerId, startDate = startDate, endDate = endDate, countryCode = countryCode }); } public List GetPlaylistProcessedDates(DateTime startDate, DateTime endDate) { return PetaPocoRepository.ReadOnlyInstance.Fetch("WHERE Date > @startDate AND Date < @endDate", new { startDate, endDate }); } public DateTime? GetLatestProccessedDateWithFullData() { return PetaPocoRepository.ReadOnlyInstance.FirstOrDefaultWithSql("WHERE AggregatedTimestamp IS NOT NULL AND CategoryStreamsSummary IS NOT NULL ORDER BY Date DESC")?.Date; } public List GetImportedFiles(DateTime date) { return PetaPocoRepository.ReadOnlyInstance.Fetch("WHERE date=@0", date); } public async Task SetMarketsAsCalculatedCategoryStreamsAsync(DateTime date) { string sql = @" INSERT INTO tblSpotifyAnalyticsDatePlaylistProcessed (CountryCode, DistributerId, Date, CategoryStreamsSummary) SELECT CountryCode, DistributerId, @date, Now() FROM tblSpotifyAnalyticsDate WHERE Date = @date ON DUPLICATE KEY UPDATE CategoryStreamsSummary = VALUES(CategoryStreamsSummary);"; using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { MySqlCommand cmd = conn.CreateCommand(); cmd.CommandText = sql; cmd.Parameters.AddWithValue("@date", date); await cmd.ExecuteNonQueryAsync(); } } public async Task SetMarketsAsCalculatedCategoryStreamsAsync(List importedFilesForToday, DateTime date) { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var timestamp = DateTime.UtcNow; var sqlDateFormatCulture = CultureInfo.GetCultureInfo("sv-SE"); var sql = "INSERT INTO tblSpotifyAnalyticsDatePlaylistProcessed (CountryCode, DistributerId, Date, CategoryStreamsSummary) " + "VALUES " + string.Join(",", importedFilesForToday.Select(p => $"('{MySqlHelper.EscapeString(p.CountryCode)}', {p.DistributerId}, '{date.ToString(sqlDateFormatCulture)}', '{timestamp.ToString(sqlDateFormatCulture)}')")) + "ON DUPLICATE KEY UPDATE CategoryStreamsSummary = VALUES(CategoryStreamsSummary);"; MySqlCommand cmd = new MySqlCommand(sql, conn); await cmd.ExecuteNonQueryAsync(); } } public async Task SetMarketsAsAggregatedAsync(List importedFiles, DateTime date) { if (!importedFiles.Any()) { return; } using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { var timestamp = DateTime.UtcNow; var sqlDateFormatCulture = CultureInfo.GetCultureInfo("sv-SE"); var sql = "INSERT INTO tblSpotifyAnalyticsDatePlaylistProcessed (CountryCode, DistributerId, Date, AggregatedTimestamp) " + "VALUES " + string.Join(",", importedFiles.Select(p => $"('{MySqlHelper.EscapeString(p.CountryCode)}', {p.DistributerId}, '{date.ToString(sqlDateFormatCulture)}', '{timestamp.ToString(sqlDateFormatCulture)}')")) + "ON DUPLICATE KEY UPDATE aggregatedTimestamp = VALUES(AggregatedTimestamp);"; MySqlCommand cmd = new MySqlCommand(sql, conn); await cmd.ExecuteNonQueryAsync(); } } public async Task> GetPlaylistsWithStreamsAsync(DateTime date) { List playlistsToAggregate = new List(); using (var conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { const string sql = "SELECT DISTINCT PlaylistUri FROM tblSpotifyPlaylistStream WHERE date = @date"; MySqlCommand cmd = new MySqlCommand(sql, conn); cmd.Parameters.AddWithValue("@date", date); var reader = await cmd.ExecuteReaderAsync(); while (await reader.ReadAsync()) { playlistsToAggregate.Add(reader.GetString(0)); } } return playlistsToAggregate; } private static async Task> ExeRequestForDateRangeAsync(DateTime startDate, DateTime endDate, StringBuilder sbSql) { string sqlBase = sbSql.ToString(); string execSql = sqlBase.Replace("startDate", startDate.ToNeutralShortDateString()); execSql = execSql.Replace("endDate", endDate.ToNeutralShortDateString()); using (MySqlConnection conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { IEnumerable items = await conn.QueryAsync(execSql, commandType: System.Data.CommandType.Text); return items; } } } public enum AnalyticsLogType { Started, Finished } }