/* appSetting FetchPlaylistStats_MaxThreads_For_Spotify */ using MoreLinq; using Sony.Filtr.ApolloAPI; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Contracts.Entities; using Sony.Filtr.Core.Buzz; using Sony.Filtr.Core.EditorialPlaylists; using Sony.Filtr.Core.LastModified; using Sony.Filtr.Database; using Sony.Filtr.Playlists; using Sony.Filtr.Playlists.Spotify; using Sony.Filtr.Tasks.Helpers; using Sony.Filtr.Utility; using Sony.Filtr.Utility.Extensions; using Sony.Filtr.Functional; using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Threading; using System.Threading.Tasks; using System.Threading.Tasks.Dataflow; namespace Sony.Filtr.Tasks.Tasks.Spotify { public class FetchPlaylistStats : IScheduledTask { private readonly LastModifiedManager _lastModifiedManager; private readonly EditorialPlaylistManager _editorialPlaylistManager; private readonly EditorialPlaylistStatsManager _editorialPlaylistStatsManager; private readonly SpotifyPlaylistManager _spotifyPlaylistManager; private readonly BuzzManager _buzzManager; private readonly LoggerAdapter _loggerFull; private readonly LoggerAdapter _loggerError; private readonly LoggerAdapter _loggerDuration; private readonly ActionTracker blockCountTracker; private readonly IApolloWebApi _vendorApi; public FetchPlaylistStats(LastModifiedManager lastModifiedManager, EditorialPlaylistManager editorialPlaylistManager, EditorialPlaylistStatsManager editorialPlaylistStatsManager, SpotifyPlaylistManager spotifyPlaylistManager, BuzzManager buzzManager, IApolloWebApi vendorApi) { _lastModifiedManager = lastModifiedManager; _editorialPlaylistManager = editorialPlaylistManager; _editorialPlaylistStatsManager = editorialPlaylistStatsManager; _spotifyPlaylistManager = spotifyPlaylistManager; _buzzManager = buzzManager; _loggerFull = LoggerAdapter.GetLogger("FetchPlaylistStats"); _loggerDuration = LoggerAdapter.GetLogger("FetchPlaylistStats_Duration"); _loggerError = LoggerAdapter.GetLogger("FetchPlaylistStats_Errors"); blockCountTracker = new ActionTracker("logs\\FetchPlaylistStats_BlocksCount.txt", 10); this._vendorApi = vendorApi; } private static DateTime GetProcessingDate() { return Maybe.GetAppSettingsDateTimeOrDefault("FetchPlaylistStats_ProcessingDate", DateTime.Today).Value; } public async Task ExecuteAsync(Guid scheduledTaskLogId) { DateTime processingDate = GetProcessingDate(); await processingDate.ToResult() .Tap(dt => _loggerFull.Information(() => $"Starting processing for date {dt}")) .BindAsync(GetPlaylistIdsToProcessAsync) .TapAsync(ids => _loggerFull.Information(() => $"Fetching for {ids.Length} playlists and date '{processingDate.ToString("yyyy-MM-dd")}'.")) .TapAsync(((Func, Task>)FetchStatsAsync) .ToUnit() .Duration((ids, unit, duration) => _loggerFull.Information(() => $"Collecting and saving stats took {duration.TotalSeconds}s"))) .TapAsync(SetLastModifiedAsync) .TapAsync(ClearCacheForBuzzImportUsers) .TapAsync(() => _loggerFull.Information(() => "Begin calculate trend")) .TapAsync(CalculateRecentChangesAsync) //.ToTask() .TapAsync(CalculateRecentChangesNewAsync) .TapAsync(() => _loggerFull.Information(() => "Done calculating trend")) .TapAsync(() => _loggerFull.Information(() => "Done with all!")) ; return null; } private async Task> GetPlaylistIdsToProcessAsync(DateTime date) { //return new string[] { "1H6NwhJTyicXHBvEK9yIsp" }; Func> getPlaylistIdsForDate = async dt => { List playlistIds = new List(); using (var conn = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { var cmd = conn.CreateCommand(); cmd.CommandText = @" SELECT p.playlistId FROM tblSpotifyPlaylist AS p LEFT JOIN tblPlaylistIgnoredPlaylist AS ip ON ip.playlistId=p.playlistUri AND ip.musicServiceId=@spotifyMusicServiceId LEFT JOIN tblSpotifyPlaylistFollowerHistory AS f ON f.playlistId = p.playlistId AND f.date = @date WHERE p.Removed=0 AND ip.playlistId IS NULL AND p.SaveTracklist = 1 AND f.playlistId IS NULL ORDER BY p.TrackLatestAdded DESC ;"; //We might want subscribers for these as well but right now we keep the existing behavior; cmd.Parameters.AddWithValue("@spotifyMusicServiceId", MusicService.Spotify); cmd.Parameters.AddWithValue("@date", date); using (var sqlReader = await cmd.ExecuteReaderAsync()) { while (await sqlReader.ReadAsync()) { playlistIds.Add(sqlReader.GetString(0)); } } } return playlistIds.ToArray(); }; return await getPlaylistIdsForDate .Duration((dt, list, duration) => _loggerFull.Trace(() => $"Loaded sony playlists from db. Took {duration.TotalMilliseconds} ms")) .TryCatch() (date); } private async Task CalculateRecentChangesNewAsync() { var processingDate = GetProcessingDate(); const int limit = 500_000; Func>> getFollowers = async offset => { string selectSql = $@" SELECT p.PlaylistUri, p.playlistId, l1.Followers, l2.Followers as Followers1dayAgo, l3.Followers as Followers7dayAgo, l4.Followers as Followers14daysAgo, l5.Followers as Followers28daysAgo, l6.Followers as Followers56daysAgo FROM tblSpotifyPlaylist AS p LEFT JOIN tblPlaylistIgnoredPlaylist AS ip ON ip.playlistId=p.playlistUri AND ip.musicServiceId=@spotifyMusicServiceId LEFT JOIN tblSpotifyPlaylistFollowerHistory As l1 ON p.PlaylistId = l1.PlaylistId AND l1.Date = @today LEFT JOIN tblSpotifyPlaylistFollowerHistory As l2 ON p.PlaylistId = l2.PlaylistId AND l2.Date = @dayAgo1 LEFT JOIN tblSpotifyPlaylistFollowerHistory As l3 ON p.PlaylistId = l3.PlaylistId AND l3.Date = @dayAgo7 LEFT JOIN tblSpotifyPlaylistFollowerHistory As l4 ON p.PlaylistId = l4.PlaylistId AND l4.Date = @dayAgo14 LEFT JOIN tblSpotifyPlaylistFollowerHistory As l5 ON p.PlaylistId = l5.PlaylistId AND l5.Date = @dayAgo28 LEFT JOIN tblSpotifyPlaylistFollowerHistory As l6 ON p.PlaylistId = l6.PlaylistId AND l6.Date = @dayAgo56 WHERE p.removed = 0 AND ip.playlistID IS NULL GROUP BY p.PlaylistUri ORDER BY p.PlaylistUri LIMIT {limit} OFFSET {offset};"; using (var readonlyConnection = await DatabaseHandler.GetOpenReadOnlyConnectionAsync()) { var selectCommand = readonlyConnection.CreateCommand(); selectCommand.CommandText = selectSql; selectCommand.Parameters.AddWithValue("@spotifyMusicServiceId", MusicService.Spotify); selectCommand.Parameters.AddWithValue("@today", processingDate); selectCommand.Parameters.AddWithValue("@dayAgo1", processingDate.AddDays(-1)); selectCommand.Parameters.AddWithValue("@dayAgo7", processingDate.AddDays(-7)); selectCommand.Parameters.AddWithValue("@dayAgo14", processingDate.AddDays(-14)); selectCommand.Parameters.AddWithValue("@dayAgo28", processingDate.AddDays(-28)); selectCommand.Parameters.AddWithValue("@dayAgo56", processingDate.AddDays(-56)); var dtos = new List(); using (var sqlReader = await selectCommand.ExecuteReaderAsync()) { while (await sqlReader.ReadAsync()) { dtos.Add(new FollowersHistoryDto() { PlaylistId = sqlReader.GetString("playlistId"), PlaylistUri = sqlReader.GetString("PlaylistUri"), Followers = sqlReader.GetIntOrDefault("Followers"), Followers1dayAgo = sqlReader.GetIntOrDefault("Followers1dayAgo"), Followers7daysAgo = sqlReader.GetIntOrDefault("Followers7dayAgo"), Followers14daysAgo = sqlReader.GetIntOrDefault("Followers14daysAgo"), Followers28daysAgo = sqlReader.GetIntOrDefault("Followers28daysAgo"), Followers56daysAgo = sqlReader.GetIntOrDefault("Followers56daysAgo"), }); } } return dtos; } }; Func, Task> bulkLoadFollowers = async dtos => { await BulkLoader.LoadAsync( "tblSpotifyPlaylistFollowers", new Func[] { dto => dto.PlaylistUri, dto => dto.PlaylistId, dto => dto.Followers.ToString(), dto => dto.Followers1dayAgo.ToString(), dto => dto.Followers7daysAgo.ToString(), dto => dto.Followers14daysAgo.ToString(), dto => dto.Followers28daysAgo.ToString(), dto => dto.Followers56daysAgo.ToString(), dto => DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss") }, new string[] { "PlaylistUri", "PlaylistId", "Followers", "Followers1dayAgo", "Followers7daysAgo", "Followers14daysAgo", "Followers28daysAgo", "Followers56daysAgo", "UpdateDate" }, dtos, MySql.Data.MySqlClient.MySqlBulkLoaderConflictOption.Replace); }; bool readAll = false; int currentOffset = 0; var followers = new List(); while (!readAll) { await getFollowers .Retry(1) .Duration((offset, ids, duration) => _loggerFull.Information(() => $"Read {ids.Count()} playlist with offset {offset} ids in {duration.TotalSeconds}secs")) .TryCatch() .OnSuccess((offset, result) => readAll = result.Value.Count() < limit) (currentOffset) .TapAsync(f => followers.AddRange(f)) .TapAsync(f => currentOffset += limit); } await bulkLoadFollowers .ToUnit() .Retry(1) .Duration((ids, unit, duration) => _loggerFull.Information(() => $"Inserted {ids.Count()} in {duration.TotalSeconds}secs")) .TryCatch() .OnFailure((dtos, result) => LogError(result.Exception, () => $"Could not bulk load recent changes to tblSpotifyPlaylistFollowers")) (followers); } private class FollowersHistoryDto { public string PlaylistUri { get; set; } public string PlaylistId { get; set; } public int? Followers { get; set; } public int? Followers1dayAgo { get; set; } public int? Followers7daysAgo { get; set; } public int? Followers14daysAgo { get; set; } public int? Followers28daysAgo { get; set; } public int? Followers56daysAgo { get; set; } } private static int GetMaxDegreeOfParallelismForSpotifyCommunication() { return Maybe.GetAppSettingsIntOrDefault("FetchPlaylistStats_MaxThreads_For_Spotify", 2); } private async Task FetchStatsAsync(IEnumerable playlistIds) { var totalPlaylists = playlistIds.Count(); Func> getPlaylistStatsFromSpotify = this.GetPlaylistStatsFromSpotifyByIdAsync; TransformBlock, Result> getSpotifyDataBlock = getPlaylistStatsFromSpotify .Timeout(TimeSpan.FromSeconds(7)) .Retry(2, TimeSpan.FromSeconds(5)) .Duration((id, stats, duration) => _loggerDuration.Trace(() => $"Received stats for playlist '{id}' in {duration.TotalMilliseconds}ms")) .Map, PlaylistStat, string, PlaylistStat>( inputTransform: idIndex => idIndex.Item, outputTransform: (id, stats, idIndex) => stats ) .TryCatch() .OnSuccess((idIndex, result) => { if (result.Value != null) { _loggerFull.Debug(() => $"Received stats for '{idIndex.Item}' {idIndex.Index}/{totalPlaylists}. Followers {result.Value.Subscribers}"); result.Value.Date = GetProcessingDate(); } }) .OnFailure((id, result) => { this.LogError(result.Exception, () => "Could not get playlist data from Spotify."); }) .AsTransformBlock(new ExecutionDataflowBlockOptions() { MaxDegreeOfParallelism = GetMaxDegreeOfParallelismForSpotifyCommunication(), EnsureOrdered = false, SingleProducerConstrained = true }); var updateBatchBlock = new BatchBlock>(100); Func>, Task> updateFollowers = async subscriberData => { await Task.WhenAll( _spotifyPlaylistManager.SetFollowersAsync(subscriberData.Select(d => new KeyValuePair(d.Value.PlaylistId, d.Value.Subscribers))), _editorialPlaylistStatsManager.AddEditorialPlaylistSubscriptionStatisticAsync(subscriberData.Select(r => r.Value).ToList())); }; var updateFollowersBlock = updateFollowers .ToUnit() .Timeout(TimeSpan.FromSeconds(30)) .Retry(1) .Duration((stats, unit, duration) => _loggerDuration.Trace(() => $"Updated followers for batch in {duration.TotalMilliseconds}ms")) .TryCatch() .OnSuccess((statsResult, result) => _loggerFull.Information(() => $"Updated followers for (playlistId;subscribers) {String.Join(",", statsResult.Select(d => $"({d.Value.PlaylistId};{d.Value.Subscribers})"))}")) .OnFailure((statsResult, result) => this.LogError(result.Exception, () => $"Could not update followers for playlist id: {String.Join(",", statsResult.Select(d => $"({d.Value.PlaylistId};{d.Value.Subscribers})"))}")) .AsActionBlock(new ExecutionDataflowBlockOptions() { MaxDegreeOfParallelism = 2 }); var linkOptions = new DataflowLinkOptions() { PropagateCompletion = true }; getSpotifyDataBlock.LinkTo(updateBatchBlock, linkOptions, dto => dto.IsOk && dto.Value != null); getSpotifyDataBlock.LinkTo(DataflowBlock.NullTarget>(), linkOptions); updateBatchBlock.LinkTo(updateFollowersBlock, linkOptions); CancellationTokenSource cts = new CancellationTokenSource(); var blockTrackingTask = TaskHelper.CreateLongRunningLoggingThread( cts.Token, () => $@" Get block: {getSpotifyDataBlock.InputCount}/{getSpotifyDataBlock.OutputCount} batch block: {updateBatchBlock.OutputCount} update followers block: {updateFollowersBlock.InputCount}", this.blockCountTracker, TimeSpan.FromMinutes(1)); foreach (var playlistId in playlistIds.Where(id => !(string.IsNullOrWhiteSpace(id) || id.Equals("starred"))).ItemIndex()) { try { await getSpotifyDataBlock.SendAsync(playlistId); } catch (Exception ex) { this.LogError(ex, () => $"Could not send playlst id '{playlistId}' to pipeline. get block: {getSpotifyDataBlock.InputCount}/{getSpotifyDataBlock.OutputCount} batch block: {updateBatchBlock.OutputCount} update followers block: {updateFollowersBlock.InputCount}"); } } getSpotifyDataBlock.Complete(); await updateFollowersBlock.Completion; cts.Cancel(); await this.blockCountTracker.CompleteAsync(); } private static bool LogRawHeaders() { return Maybe.GetAppSettingsBooleanOrDefault("FetchPlaylistStats_LogRawHeaders", false); } private async Task GetPlaylistStatsFromSpotifyByIdAsync(string playlistId) { var playlist = await this._vendorApi.GetSpotifyPlaylistByIdAsync(playlistId, TimeSpan.FromHours(23), fields: Sony.Filtr.Tasks.Helpers.Constants.Spotify.VendorPlaylistWithTracksFullFieldsList); if (playlist == null) { _loggerFull.Warn(() => $"Could not find playlist by id. Id: {playlistId}"); return null; } if (!playlist.followers.total.HasValue) { _loggerFull.Warn(() => $"Playlist {playlistId} has no follower info."); return null; } if (playlist.followers.total.Value < 0) { _loggerFull.Warn(() => $"Playlist {playlistId} has negative number of followers."); return null; } Func getRawResponseHeaders = pl => { StringBuilder sb = new StringBuilder(); foreach (var requestPair in pl.RequestHeaders) { sb.Append($"Request URI: '{requestPair.Key}'->"); foreach (var header in requestPair.Value) { sb.Append($"{header.Key}:{header.Value};"); } } return sb.ToString(); }; string logMessage = $"Execution time summary '{playlistId}': { playlist.GetSummary()}"; if (LogRawHeaders()) { logMessage += $" raw headers: {getRawResponseHeaders(playlist)}"; } _loggerFull.Information(() => logMessage); return new PlaylistStat() { PlaylistId = playlistId, Subscribers = (uint)playlist.followers.total.Value }; } private async Task CalculateRecentChangesAsync(IEnumerable playlistIds) { var batchs = playlistIds.Batch(100).ToList(); DateTime processingDate = GetProcessingDate(); await batchs.ItemIndex().ForEachAsync(3, async playlistBatch => { try { _loggerFull.Debug(() => $"Get playlist stats for batch: {playlistBatch.Index} out of {batchs.Count}"); var stats = await _editorialPlaylistStatsManager.GetPlaylistStatsAsync(playlistBatch.Item.ToList(), processingDate.AddDays(-56), processingDate); var grouped = stats.GroupBy(k => k.PlaylistId).Select(g => new KeyValuePair>(g.Key, g.ToList())).ToList(); foreach (var pair in grouped) { try { await _editorialPlaylistManager.SaveAttributesAsync(CalculateAttributeFollowerChange(pair.Key, pair.Value)); } catch (Exception ex) { this.LogError(ex, () => "Could not save attributes"); } } } catch (Exception ex) { this.LogError(ex, () => "Could not process playlist batch"); } }); } private EditorialPlaylistAttributes CalculateAttributeFollowerChange(string playlistId, IEnumerable playlistStat) { var processingDate = GetProcessingDate(); var daysAgo15 = playlistStat.FirstOrDefault(s => s.Date.Date == processingDate.AddDays(-15)); var daysAgo30 = playlistStat.FirstOrDefault(s => s.Date.Date == processingDate.AddDays(-30)); var daysAgo7 = playlistStat.FirstOrDefault(s => s.Date.Date == processingDate.AddDays(-7)); uint latestSubscribers = 0; if (playlistStat.Any()) { latestSubscribers = playlistStat.OrderByDescending(s => s.Date).First().Subscribers; } uint secondLatestSubscribers = 0; if (playlistStat.Count() > 2) { secondLatestSubscribers = playlistStat.OrderByDescending(s => s.Date).ElementAt(1).Subscribers; } int? followersChange15 = null; if (daysAgo15 != null) followersChange15 = (int)(latestSubscribers - daysAgo15.Subscribers); int? followersChange30 = null; if (daysAgo30 != null) followersChange30 = (int)(latestSubscribers - daysAgo30.Subscribers); int? followersChange7 = null; if (daysAgo7 != null) followersChange7 = (int)(latestSubscribers - daysAgo7.Subscribers); var playlistLink = SpotifyLink.FromPlaylistId(playlistId); var attributes = new EditorialPlaylistAttributes(playlistLink.Uri) { Followers = latestSubscribers, LastFollowers = secondLatestSubscribers, FollowerChange7 = followersChange7, FollowerChange15 = followersChange15, FollowerChange30 = followersChange30, }; return attributes; } private async Task SetLastModifiedAsync() { await _lastModifiedManager.SetLastModifyDateAsync(null, EntityType.PlaylistStats); } private void ClearCacheForBuzzImportUsers() { var users = _buzzManager.GetBuzzImportUsers(); foreach (var buzzImportUser in users) { _buzzManager.ClearBuzzImportUserCache(buzzImportUser); } } private void LogError(Exception ex, Func getMessage) { this._loggerFull.Error(ex, getMessage); this._loggerError.Error(ex, getMessage); } } }