using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Globalization; using System.Linq; using System.Threading.Tasks; using MoreLinq; using MySql.Data.MySqlClient; using NLog; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Contracts.Entities; using Sony.Filtr.Database; using Sony.Filtr.Playlists; using Sony.Filtr.SpotifyAnalytics; using Sony.Filtr.SpotifyAnalytics.Data; using Sony.Filtr.Utility; using Sony.Filtr.Utility.Extensions; using Sony.Filtr.Functional; namespace Sony.Filtr.Tasks.Tasks.Spotify.Analytics { public class PlaylistStreamsSummaryCalculator { private readonly SpotifyAnalyticsManager _spotifyAnalyticsManager; private readonly Logger _logger; private const string TempFolder = "spotifyanalytics"; public PlaylistStreamsSummaryCalculator(SpotifyAnalyticsManager spotifyAnalyticsManager) { _spotifyAnalyticsManager = spotifyAnalyticsManager; _logger = LogManager.GetLogger("SpotifyStreamingAggregatedReportApi"); } public async Task SetPlaylistStreamsSummaryAsync() { var latestDayWithData = await _spotifyAnalyticsManager.GetLatestDayWithDataAsyncNew(); await SetPlaylistStreamsSummaryAsync(latestDayWithData); } private static int GetInitialPlaylistsDBQueryRetries() { return Maybe.GetAppSettingsIntOrDefault("PlaylistStreamsSummaryCalculator_InitialPlaylistsDBQueryRetries", 2); } public async Task SetPlaylistStreamsSummaryAsync(DateTime latestDayWithData) { _logger.Debug("Begin SetPlaylistStreamsSummaryAsync"); await this.GetPlaylistIds() .Bind(ToUriOrId) .TapAsync(_spotifyAnalyticsManager.ClearPlaylistStreamSummaryCacheAsync, list => list.Select(p => p.Id)) .BindAsync(ids => this.CalculatePlaylistStreamsSummaryAsync(ids, latestDayWithData)) .TapAsync(summaries => _logger.Debug($"Done calculating stream summary. Total: {summaries.Length}")) .TapAsync(BulkLoad) .TapAsync(summaries => _logger.Debug($"Bulk loaded stream summaries. Total: {summaries.Length}")) .TapAsync(CalculateGlobalPlaylistStreamSummaryAsync) ; _logger.Debug("Done SetPlaylistStreamsSummaryAsync"); } private async Task CalculatePlaylistStreamsSummaryAsync(IEnumerable items, DateTime latestDayWithData) { ConcurrentBag allPlayliststreamsSummaries = new ConcurrentBag(); var accounts = new[] { SpotifyAnalyticsAccount.Sony, SpotifyAnalyticsAccount.Orchard, SpotifyAnalyticsAccount.SonyMusicEntertainmentJapan, SpotifyAnalyticsAccount.SonyMusicEntertainmentJapanInternational }; await items.ItemIndex().ForEachAsync(10, async playlist => { Console.WriteLine($"Calculating for batch {playlist.Index} / {items.Count()}"); Dictionary> mergedData = new Dictionary>(); foreach (var account in accounts) { List data = await _spotifyAnalyticsManager.CalculatePlaylistStreamSummaryAsync((int)account, playlist.Item.Id, playlist.Item.Uri, latestDayWithData); mergedData.Add((int)account, data); } var groupedPlaylistStreamSummaries = mergedData.SelectMany(p => p.Value, (pair, summary) => new { AccountId = pair.Key, Summary = summary, }) .GroupBy(p => new { p.Summary.PlaylistId, p.Summary.Market }) .Select(playlistMarketData => BuildPlaylistStreamsSummary(playlistMarketData.Key.PlaylistId, playlistMarketData.Key.Market, playlistMarketData.Select(p => p.Summary).ToList(), playlistMarketData.FirstOrDefault(p => p.AccountId == 1)?.Summary)); groupedPlaylistStreamSummaries.ForEach(p => allPlayliststreamsSummaries.Add(p)); }); return allPlayliststreamsSummaries.ToArray(); } private async Task> GetPlaylistIds() { Func> getPlaylists = async b => await GetPlaylistWithRecentStreamsDataAsync(); return await getPlaylists .Timeout(TimeSpan.FromMinutes(25)) .Retry(GetInitialPlaylistsDBQueryRetries()) .TryCatch() .OnFailure((b, result) => _logger.Error(result.Exception, "Could not read playlists with recent streams")) (false); } private async Task BulkLoad(IEnumerable streams) { Func FormatDateForImportFile = (currentDate) => currentDate.Year + "-" + currentDate.Month.ToString("D2") + "-" + currentDate.Day.ToString("D2"); var defaultNumberFormat = new NumberFormatInfo() { NumberDecimalSeparator = "." }; await BulkLoader .ToExtensionFunc( "tblSpotifyPlaylistStreamSummary", new Func[] { s => s.PlaylistId, s => s.Market, s => s.Streams56days.ToString(), s => s.Streams28days.ToString(), s => s.Streams14days.ToString(), s => s.Streams7days.ToString(), s => s.Listeners56days.ToString(), s => s.Listeners28days.ToString(), s => s.Listeners14days.ToString(), s => s.Listeners7days.ToString(), s => s.StreamsPerListener56days.ToString(defaultNumberFormat), s => s.StreamsPerListener28days.ToString(defaultNumberFormat), s => s.StreamsPerListener14days.ToString(defaultNumberFormat), s => s.StreamsPerListener7days.ToString(defaultNumberFormat), s => s.StreamsLatest.ToString(), s => s.ListenersLatest.ToString(), s => s.StreamsPerListenerLatest.ToString(defaultNumberFormat), s => s.StreamDays56Days.ToString(), s => s.StreamDays28Days.ToString(), s => s.StreamDays14Days.ToString(), s => s.StreamDays7Days.ToString(), s => s.TotalStreamDays.ToString(), s => FormatDateForImportFile(s.LatestDate), s => FormatDateForImportFile(DateTime.UtcNow) }, new string[] { "playlistId", "market", "Streams56days", "Streams28days", "Streams14days", "Streams7days", "Listeners56days", "Listeners28days", "Listeners14days", "Listeners7days", "StreamsPerListener56days", "StreamsPerListener28days", "StreamsPerListener14days", "StreamsPerListener7days", "StreamsLatest", "ListenersLatest", "StreamsPerListenerLatest", "StreamDays56days", "StreamDays28days", "StreamDays14days", "StreamDays7days", "TotalStreamDays", "LatestDate", "Updated" }, streams) .Timeout(TimeSpan.FromMinutes(20)) .Retry(2) .TryCatch() .OnFailure((fileName, result) => _logger.Error(result.Exception, $"Could not bulk load PlaylistStreamsSummary. Total: '{streams.Count()}'")) (MySqlBulkLoaderConflictOption.Replace); } private async Task CalculateGlobalPlaylistStreamSummaryAsync() { Func calculateGlobalSummary = async b => { using (var conn = await DatabaseHandler.GetOpenConnectionAsync()) { const string sql = @" INSERT IGNORE INTO tblSpotifyPlaylistStreamSummaryGlobal SELECT PlaylistId, SUM(Streams56days), SUM(Streams28days), SUM(Streams14days), SUM(Streams7days), SUM(Listeners56days), SUM(Listeners28days), SUM(Listeners14days), SUM(Listeners7days), SUM(StreamsLatest), SUM(ListenersLatest), SUM(StreamsPerListenerLatest), SUM(StreamsPerListener7days), SUM(StreamsPerListener14days), SUM(StreamsPerListener28days), SUM(StreamsPerListener56days), MAX(StreamDays56days), MAX(StreamDays28days), MAX(StreamDays14days), MAX(StreamDays7days), MAX(TotalStreamDays), SUM(Streams56Days / StreamDays56Days) as avgStreams56, SUM(Streams28Days / StreamDays28Days) as avgStreams28, SUM(Streams14Days / StreamDays14Days) as avgStreams14, SUM(Streams7Days / StreamDays7Days) as avgStreams7, SUM(Listeners56Days / StreamDays56Days) as avgListeners56, SUM(Listeners28Days / StreamDays28Days) as avgListeners28, SUM(Listeners14Days / StreamDays14Days) as avgListeners14, SUM(Listeners7Days / StreamDays7Days) as avgListeners7, MAX(LatestDate), MAX(Updated) FROM tblSpotifyPlaylistStreamSummary GROUP BY PlaylistId ON DUPLICATE KEY UPDATE Streams56days = VALUES(Streams56days), Streams28days = VALUES(Streams28days), Streams14days = VALUES(Streams14days), Streams7days = VALUES(Streams7days), Listeners56days = VALUES(Listeners56days), Listeners28days = VALUES(Listeners28days), Listeners14days = VALUES(Listeners14days), Listeners7days = VALUES(Listeners7days), StreamsLatest = VALUES(StreamsLatest), ListenersLatest = VALUES(ListenersLatest), StreamsPerListenerLatest = VALUES(StreamsPerListenerLatest), StreamsPerListener7days = VALUES(StreamsPerListener7days), StreamsPerListener14days = VALUES(StreamsPerListener14days), StreamsPerListener28days = VALUES(StreamsPerListener28days), StreamsPerListener56days = VALUES(StreamsPerListener56days), StreamDays56days = VALUES(StreamDays56days), StreamDays28days = VALUES(StreamDays28days), StreamDays14days = VALUES(StreamDays14days), StreamDays7days = VALUES(StreamDays7days), TotalStreamDays = VALUES(TotalStreamDays), avgStreams56 = VALUES(avgStreams56), avgStreams28 = VALUES(avgStreams28), avgStreams14 = VALUES(avgStreams14), avgStreams7 = VALUES(avgStreams7), avgListeners56 = VALUES(avgListeners56), avgListeners28 = VALUES(avgListeners28), avgListeners14 = VALUES(avgListeners14), avgListeners7 = VALUES(avgListeners7), LatestDate = VALUES(LatestDate), Updated = VALUES(Updated)"; MySqlCommand cmd = new MySqlCommand(sql, conn); await cmd.ExecuteNonQueryAsync(); } }; await calculateGlobalSummary .ToUnit() .Timeout(TimeSpan.FromMinutes(20)) .Retry(2) .TryCatch() .OnFailure((b, result) => _logger.Error(result.Exception, $"Could not calculate global playlist stream summary")) .Partial(false) (); } private async Task GetPlaylistWithRecentStreamsDataAsync() { DateTime startDate = DateTime.Today.AddDays(-56); //The longest time period we aggregate for is 56 days. List playlists = new List(); using (var conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { //TODO: Try and optimize. Use actual dates? MySqlCommand cmd = new MySqlCommand("SELECT DISTINCT PlaylistUri FROM tblSpotifyAnalyticsAccountPlaylistStreamInfo WHERE Date >= @startDate", conn); cmd.Parameters.AddWithValue("@startDate", startDate); var reader = await cmd.ExecuteReaderAsync(); while (await reader.ReadAsync()) { playlists.Add(reader.GetString(0)); } } return playlists.ToArray(); } private class SpotifyPlaylistUriOrId { public string Id { get; set; } public string Uri { get; set; } } private SpotifyPlaylistUriOrId[] ToUriOrId(IEnumerable ids) { var playlistUris = ids.Where(p => p.StartsWith("spotify:user:")).ToList(); var plalistUriLookup = playlistUris.Where(p => !string.IsNullOrWhiteSpace(p)).Select(p => new { PlaylistId = new SpotifyLink(p).ExtractPlaylistID(), Uri = p }).Where(p => p.PlaylistId != null).ToDictionary(k => k.PlaylistId, v => v.Uri); var playlistIds = ids.Except(playlistUris).ToList(); var allPlaylistIds = plalistUriLookup.Keys.Union(playlistIds).Distinct().ToList(); List uriOrIdList = new List(); foreach (var playlistId in allPlaylistIds) { uriOrIdList.Add(new SpotifyPlaylistUriOrId() { Id = playlistId, Uri = plalistUriLookup.GetValueOrDefault(playlistId), }); } return uriOrIdList.ToArray(); } private PlaylistStreamsSummary BuildPlaylistStreamsSummary(string playlistId, string market, List allAccounts, PlaylistStreamsSummary sonyData) { if (sonyData == null) { return new PlaylistStreamsSummary() { Market = market, PlaylistId = playlistId, LatestDate = allAccounts.Max(p => p.LatestDate), ListenersLatest = 0, Listeners7days = 0, Listeners14days = 0, Listeners28days = 0, Listeners56days = 0, TotalStreamDays = 0, StreamDays56Days = 0, StreamDays28Days = 0, StreamDays14Days = 0, StreamDays7Days = 0, StreamsPerListenerLatest = 0, StreamsPerListener7days = 0, StreamsPerListener14days = 0, StreamsPerListener28days = 0, StreamsPerListener56days = 0, StreamsLatest = allAccounts.Sum(p => p.StreamsLatest), Streams7days = allAccounts.Sum(p => p.Streams7days), Streams14days = allAccounts.Sum(p => p.Streams14days), Streams28days = allAccounts.Sum(p => p.Streams28days), Streams56days = allAccounts.Sum(p => p.Streams56days), }; } Func SafePercentage = (sonyDataStreamsLatest, sonyDataListenersLatest) => sonyDataListenersLatest == 0 ? 0 : (double)sonyDataStreamsLatest / (double)sonyDataListenersLatest; return new PlaylistStreamsSummary() { Market = market, PlaylistId = playlistId, LatestDate = allAccounts.Max(p => p.LatestDate), ListenersLatest = sonyData.ListenersLatest, Listeners7days = sonyData.Listeners7days, Listeners14days = sonyData.Listeners14days, Listeners28days = sonyData.Listeners28days, Listeners56days = sonyData.Listeners56days, StreamsPerListenerLatest = SafePercentage(sonyData.StreamsLatest, sonyData.ListenersLatest), StreamsPerListener7days = SafePercentage(sonyData.Streams7days, sonyData.Listeners7days), StreamsPerListener14days = SafePercentage(sonyData.Streams14days, sonyData.Listeners14days), StreamsPerListener28days = SafePercentage(sonyData.Streams28days, sonyData.Listeners28days), StreamsPerListener56days = SafePercentage(sonyData.Streams56days, sonyData.Listeners56days), StreamsLatest = allAccounts.Sum(p => p.StreamsLatest), Streams7days = allAccounts.Sum(p => p.Streams7days), Streams14days = allAccounts.Sum(p => p.Streams14days), Streams28days = allAccounts.Sum(p => p.Streams28days), Streams56days = allAccounts.Sum(p => p.Streams56days), TotalStreamDays = sonyData.TotalStreamDays, StreamDays56Days = sonyData.StreamDays56Days, StreamDays28Days = sonyData.StreamDays28Days, StreamDays14Days = sonyData.StreamDays14Days, StreamDays7Days = sonyData.StreamDays7Days, }; } } }