using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.IO; using System.Linq; using System.Threading; using System.Threading.Tasks; using log4net; using MoreLinq; using Sony.Filtr.SpotifyAnalytics.Data; using Sony.Filtr.SpotifyAnalytics.Models; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.SpotifyAnalytics { public class SpotifyAnalyticsAnalyzer { private readonly ILog _logger; private static readonly object _ConsoleWriterLock = new object(); public SpotifyAnalyticsAnalyzer() { _logger = LogManager.GetLogger("SpotifyAnalyticsUpdate"); } public async Task AnalyzeStreamsMarketDataFileAsync(string streamsFilePath, ConcurrentDictionary tracksDictionary, string market, DateTime fileDate, AnalysisInputData inputData) { var analyticsData = new AnalyticsData() { TrackInformation = new ConcurrentDictionary(), StreamInformation = new StreamDataSummaryForAnalysis() { Market = market, Date = fileDate }, PlaylistStreamInformation = new ConcurrentDictionary(), NewPlaylists = new ConcurrentDictionary(), }; var users = new SpotifyAnalyticsUniqueUsers() { UsersPerPlaylist = new ConcurrentDictionary>(), AllUsers = new ConcurrentDictionary(), AlbumSourceUsers = new ConcurrentDictionary(), ArtistSourceUsers = new ConcurrentDictionary(), CollectionSourceUsers = new ConcurrentDictionary(), SearchSourceUsers = new ConcurrentDictionary(), PlaylistSourceUsers = new ConcurrentDictionary(), OtherSourceUsers = new ConcurrentDictionary(), DesktopDeviceUsers = new ConcurrentDictionary(), MobileDeviceUsers = new ConcurrentDictionary(), TabletDeviceUsers = new ConcurrentDictionary(), }; var streamFileRows = File.ReadLines(streamsFilePath); _logger.Info("\nAnalyzing tracks and streams for market " + market); int row = 0; int missedRow = 0; var startTime = DateTime.Now; Console.WriteLine("Number of rows analyzed:"); Console.Write(row); var allPlaylists = Enumerable.ToHashSet(inputData.AllPlaylistUris); await streamFileRows.ForEachAsync(10, async line => { if (!string.IsNullOrWhiteSpace(line)) { try { var streamsObject = ParseStreamObject(line); AdditionalTrackData additionalTrackData; if (tracksDictionary.TryGetValue(streamsObject.TrackId, out additionalTrackData)) { analyticsData.TrackInformation.TryAdd(additionalTrackData.TrackUri, new AdditionalTrackData() { ISRC = additionalTrackData.ISRC, UPC = additionalTrackData.UPC }); } // Set Playlist stream data analyticsData.StreamInformation = GetDetailedStreamInformation(streamsObject, analyticsData.StreamInformation, market, inputData.SonyPlaylistsByMarket, users); UpdatePlaylistStreamInformation(streamsObject, analyticsData.PlaylistStreamInformation, market, inputData.SonyPlaylistsByMarket, users); if (streamsObject.Source == StreamSource.Others_Playlist) { var playlistUri = streamsObject.SourceUri; if (!allPlaylists.Contains(playlistUri)) { analyticsData.NewPlaylists.TryAdd(playlistUri, new byte()); } } } catch (Exception ex) { _logger.Error(ex); } Interlocked.Increment(ref row); if (row % 100000 == 0) { WritePositionToConsole(row); } } else { missedRow++; } }); Console.Write("\n"); analyticsData.StreamInformation.UniqueUsers = users.AllUsers.Count; analyticsData.StreamInformation.UniqueUsersDeviceDesktop = users.DesktopDeviceUsers.Count; analyticsData.StreamInformation.UniqueUsersDeviceMobile = users.MobileDeviceUsers.Count; analyticsData.StreamInformation.UniqueUsersDeviceTablet = users.TabletDeviceUsers.Count; analyticsData.StreamInformation.UniqueUsersSourceAlbum = users.AlbumSourceUsers.Count; analyticsData.StreamInformation.UniqueUsersSourceArtist = users.ArtistSourceUsers.Count; analyticsData.StreamInformation.UniqueUsersSourceCollection = users.CollectionSourceUsers.Count; analyticsData.StreamInformation.UniqueUsersSourceOther = users.OtherSourceUsers.Count; analyticsData.StreamInformation.UniqueUsersSourcePlaylist = users.PlaylistSourceUsers.Count; analyticsData.StreamInformation.UniqueUsersSourceSearch = users.SearchSourceUsers.Count; analyticsData.PlaylistStreamInformation.ForEach(p => p.Value.UniqueUsers = users.UsersPerPlaylist[p.Key].Count); var endTime = DateTime.Now.Subtract(startTime); _logger.Info("Finished analysis of " + row + " rows in " + endTime.Hours.ToString("D2") + ":" + endTime.Minutes.ToString("D2") + ":" + endTime.Seconds.ToString("D2") + "." + endTime.Milliseconds); return analyticsData; } public async Task AnalyzeNewPlaylistsAsync(string streamsFilePath, string market, AnalysisInputData inputData) { var analyticsData = new AnalyticsData() { NewPlaylists = new ConcurrentDictionary(), }; var streamFileRows = File.ReadLines(streamsFilePath); _logger.Info("\nAnalyzing new playlists for market " + market); int row = 0; int missedRow = 0; var startTime = DateTime.Now; Console.WriteLine("Number of rows analyzed:"); Console.Write(row); var existingPlaylists = Enumerable.ToHashSet(inputData.AllPlaylistUris); Parallel.ForEach(streamFileRows, (line, _, lineNumber) => { if (!string.IsNullOrWhiteSpace(line)) { try { var streamsObject = ParseStreamObject(line); if (streamsObject.Source == StreamSource.Others_Playlist) { var playlistUri = streamsObject.SourceUri; if (!existingPlaylists.Contains(playlistUri)) { analyticsData.NewPlaylists.TryAdd(playlistUri, new byte()); } } } catch (Exception ex) { _logger.Error(ex); } Interlocked.Increment(ref row); if (lineNumber % 10000 == 0) { WritePositionToConsole(lineNumber); } } else { missedRow++; } }); Console.Write("\n"); var endTime = DateTime.Now.Subtract(startTime); _logger.Info("Finished analysis of " + row + " rows in " + endTime.Hours.ToString("D2") + ":" + endTime.Minutes.ToString("D2") + ":" + endTime.Seconds.ToString("D2") + "." + endTime.Milliseconds); return analyticsData; } private void WritePositionToConsole(long lineNumber) { lock (_ConsoleWriterLock) { ClearCurrentConsoleLine(); Console.SetCursorPosition(0, Console.CursorTop); Console.Write(lineNumber); } } private static void ClearCurrentConsoleLine() { var currentLineCursor = Console.CursorTop; Console.SetCursorPosition(0, Console.CursorTop); Console.Write(new string(' ', Console.WindowWidth)); Console.SetCursorPosition(0, currentLineCursor); } private StreamsApiObject ParseStreamObject(string line) { var streamsApiObject = Jil.JSON.Deserialize(line); return streamsApiObject; } private void UpdatePlaylistStreamInformation(StreamsApiObject streamsObject, ConcurrentDictionary playlistStreamInformation, string market, IDictionary sonyPlaylistMarketsByUri, SpotifyAnalyticsUniqueUsers users) { if (streamsObject.Source == StreamSource.Others_Playlist) { var playlistId = streamsObject.SourceUri; playlistStreamInformation.TryAdd(playlistId, new StreamDataSummaryForAnalysis()); var userIds = users.UsersPerPlaylist.GetOrAdd(playlistId, new ConcurrentDictionary()); userIds.GetOrAdd(streamsObject.UserId, new byte()); playlistStreamInformation[playlistId] = GetDetailedStreamInformation(streamsObject, playlistStreamInformation[playlistId], market, sonyPlaylistMarketsByUri, null); } } private StreamDataSummaryForAnalysis GetDetailedStreamInformation(StreamsApiObject streamsObject, StreamDataSummaryForAnalysis streamSummary, string market, IDictionary sonyPlaylistMarketsByUri, SpotifyAnalyticsUniqueUsers users) { Interlocked.Increment(ref streamSummary.TotalStreams); streamSummary = SetStreamOS(streamsObject, streamSummary); streamSummary = SetStreamDevice(streamsObject, streamSummary, users); streamSummary = SetStreamSource(streamsObject, streamSummary, users); streamSummary = SetSonyStream(streamsObject, streamSummary, market, sonyPlaylistMarketsByUri); users?.AllUsers?.TryAdd(streamsObject.UserId, new byte()); return streamSummary; } private StreamDataSummaryForAnalysis SetSonyStream(StreamsApiObject streamsObject, StreamDataSummaryForAnalysis streamSummary, string market, IDictionary sonyPlaylistMarketsByUri) { if (streamsObject.Source == StreamSource.Others_Playlist) { var sonyPlaylistMarket = GetValueOrDefault(sonyPlaylistMarketsByUri,streamsObject.SourceUri); if (sonyPlaylistMarket != null) { Interlocked.Increment(ref streamSummary.StreamsFromSonyPlaylists); if (sonyPlaylistMarket.Equals(market, StringComparison.InvariantCultureIgnoreCase)) { Interlocked.Increment(ref streamSummary.StreamsFromLocalSonyPlaylists); } } } return streamSummary; } private TValue GetValueOrDefault(IDictionary dictionary, TKey key) { TValue value; return dictionary.TryGetValue(key, out value) ? value : default(TValue); } private StreamDataSummaryForAnalysis SetStreamOS(StreamsApiObject streamsObject, StreamDataSummaryForAnalysis streamSummary) { var os = streamsObject.Os.ToLowerInvariant(); if (os == "windows") Interlocked.Increment(ref streamSummary.OSWindowsStreams); else if (os == "ios") Interlocked.Increment(ref streamSummary.OSiOSStreams); else if (os == "android") Interlocked.Increment(ref streamSummary.OSAndroidStreams); else if (os == "mac") Interlocked.Increment(ref streamSummary.OSMacOSStreams); //else if (os == "browser") // Interlocked.Increment(ref streamSummary.OSOtherStreams); //else if (os == "linux") // Interlocked.Increment(ref streamSummary.OSOtherStreams); //else if (os == "windows phone") // Interlocked.Increment(ref streamSummary.OSOtherStreams); //else if (os == "blackberry") // Interlocked.Increment(ref streamSummary.OSOtherStreams); //else if (os == "other") // Interlocked.Increment(ref streamSummary.OSOtherStreams); else Interlocked.Increment(ref streamSummary.OSOtherStreams); return streamSummary; } private StreamDataSummaryForAnalysis SetStreamDevice(StreamsApiObject streamsObject, StreamDataSummaryForAnalysis streamSummary, SpotifyAnalyticsUniqueUsers users) { if (streamsObject.DeviceType == DeviceType.Desktop) { Interlocked.Increment(ref streamSummary.DeviceDesktopStreams); users?.DesktopDeviceUsers?.TryAdd(streamsObject.UserId, new byte()); } else if (streamsObject.DeviceType == DeviceType.Mobile) { Interlocked.Increment(ref streamSummary.DeviceMobileStreams); users?.MobileDeviceUsers?.TryAdd(streamsObject.UserId, new byte()); } else if (streamsObject.DeviceType == DeviceType.Tablet) { Interlocked.Increment(ref streamSummary.DeviceTabletStreams); users?.TabletDeviceUsers?.TryAdd(streamsObject.UserId, new byte()); } return streamSummary; } private StreamDataSummaryForAnalysis SetStreamSource(StreamsApiObject streamsObject, StreamDataSummaryForAnalysis streamSummary, SpotifyAnalyticsUniqueUsers users) { // Available from stream data: “other”, “search”, “artist”, “album”, “collection”, “others_playlist” if (streamsObject.Source == StreamSource.Others_Playlist) { Interlocked.Increment(ref streamSummary.SourcePlaylist); users?.PlaylistSourceUsers?.TryAdd(streamsObject.UserId, new byte()); } else if (streamsObject.Source == StreamSource.Search) { Interlocked.Increment(ref streamSummary.SourceSearch); users?.SearchSourceUsers?.TryAdd(streamsObject.UserId, new byte()); } else if (streamsObject.Source == StreamSource.Artist) { Interlocked.Increment(ref streamSummary.SourceArtist); users?.ArtistSourceUsers?.TryAdd(streamsObject.UserId, new byte()); } else if (streamsObject.Source == StreamSource.Album) { Interlocked.Increment(ref streamSummary.SourceAlbum); users?.AlbumSourceUsers?.TryAdd(streamsObject.UserId, new byte()); } else if (streamsObject.Source == StreamSource.Collection) { Interlocked.Increment(ref streamSummary.SourceCollection); users?.CollectionSourceUsers?.TryAdd(streamsObject.UserId, new byte()); } else { Interlocked.Increment(ref streamSummary.SourceOther); users?.OtherSourceUsers?.TryAdd(streamsObject.UserId, new byte()); } return streamSummary; } private class SpotifyAnalyticsUniqueUsers { public ConcurrentDictionary> UsersPerPlaylist { get; set; } public ConcurrentDictionary AllUsers { get; set; } public ConcurrentDictionary PlaylistSourceUsers { get; set; } public ConcurrentDictionary SearchSourceUsers { get; set; } public ConcurrentDictionary ArtistSourceUsers { get; set; } public ConcurrentDictionary AlbumSourceUsers { get; set; } public ConcurrentDictionary OtherSourceUsers { get; set; } public ConcurrentDictionary CollectionSourceUsers { get; set; } public ConcurrentDictionary DesktopDeviceUsers { get; set; } public ConcurrentDictionary MobileDeviceUsers { get; set; } public ConcurrentDictionary TabletDeviceUsers { get; set; } } } }