using System; using System.Collections.Generic; using System.IO; using System.Linq; using System.Runtime.Caching; using System.Threading.Tasks; using MoreLinq; using MySql.Data.MySqlClient; using NLog; using PetaPoco.Business; using Sony.Filtr.Contracts.Abstractions; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.SpotifyAnalytics; using Sony.Filtr.SpotifyAnalytics.Data; using Sony.Filtr.SpotifyAnalytics.Models; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.Tasks.Tasks.Spotify.Analytics { public class AggregatedGeneralSpotifyImportTask : AggregatedImportTaskBase { private readonly SpotifyAnalyticsManager _spotifyAnalyticsManager; private readonly Logger _logger; public AggregatedGeneralSpotifyImportTask(SpotifyStreamingAggregatedReportApi aggregatedReportApi, SpotifyAnalyticsManager spotifyAnalyticsManager) : base(aggregatedReportApi) { _spotifyAnalyticsManager = spotifyAnalyticsManager; _logger = LogManager.GetLogger("SpotifyStreamingAggregatedReportApi"); } public override async Task ExecuteAsync(Guid scheduledTaskLogId) { var taskLog = new ScheduledTaskLog(); var fromDate = new DateTime(2018, 03, 21); var endDate = DateTime.Today; var dates = fromDate.GetDateRangeTo(endDate).Reverse().ToList(); if (!Directory.Exists(TempFolder)) { Directory.CreateDirectory(TempFolder); } List supportedFileTypes = new List() { //SpotifyS3FileType.Demographics, //SpotifyS3FileType.SaveSkips, SpotifyS3FileType.Summary }; await dates.ForEachAsync(1, async date => { var files = await GetFilesToImportFromS3(date, supportedFileTypes); _logger.Debug($"Found {files.Count} files for {date.ToShortDateString()}"); await files.Where(p => p.FileType != SpotifyS3FileType.Playlists).ForEachAsync(1, async trackFile => { try { var localPath = await DownloadFileFromS3Async(trackFile); if (trackFile.FileType == SpotifyS3FileType.Tracks) { _logger.Debug($"Reading tracks from {trackFile.FilePath}"); var tracks = ReadNJsonData(localPath, trackFile.DistributorId); var orderedTracks = tracks.OrderBy(p => p.country_code).ThenBy(p => p.report_country_code).ThenBy(p => p.date).ThenBy(p => p.isrc); int batchIndex = 1; foreach (var batch in orderedTracks.Batch(100000)) { _logger.Debug($"Importing batch {batchIndex} from {trackFile.FilePath}"); await PetaPocoRepository.Instance.ImportBulkFileLoaderAsync(batch); _logger.Debug($"Done batch {batchIndex} from {trackFile.FilePath}"); batchIndex++; } MarkAsProcessed(date, trackFile); } //else if (trackFile.FileType == SpotifyS3FileType.Demographics) //{ // await ImportAsync(trackFile, localPath, date, "Demographics"); //} //else if (trackFile.FileType == SpotifyS3FileType.SaveSkips) //{ // await ImportAsync(trackFile, localPath, date, "tracks save/skips"); //} else if (trackFile.FileType == SpotifyS3FileType.Summary) { _logger.Debug($"Reading tracks country summary from {trackFile.FilePath}"); var summary = ReadNJsonData(localPath, trackFile.DistributorId); var oldStructure = await MapToOldStreamSummaryStructureAsync(trackFile.Date, trackFile.DistributorId, trackFile.Country, summary.ToList()); await PetaPocoRepository.Instance.ImportBulkFileLoaderAsync(new[] { oldStructure }, MySqlBulkLoaderConflictOption.Replace); MarkAsProcessed(date, trackFile); } else if (trackFile.FileType == SpotifyS3FileType.PlaylistTracks) { _logger.Debug($"Reading playlist-track streams from {trackFile.FilePath}"); var playlistTrackStreams = ReadNJsonData(localPath, trackFile.DistributorId); _logger.Debug($"Importing playlist-track streams from {trackFile.FilePath}"); PetaPocoRepository.Instance.ImportBulkFileLoader(playlistTrackStreams); _logger.Debug($"Done importing playlist-track streams from {trackFile.FilePath}"); MarkAsProcessed(date, trackFile); } else if (trackFile.FileType == SpotifyS3FileType.ListenerSummaryTracks) { _logger.Debug($"Reading listener summary for track streams from {trackFile.FilePath}"); var trackListeners = ReadNJsonData(localPath, trackFile.DistributorId); _logger.Debug($"Importing listener summary for track streams from {trackFile.FilePath}"); PetaPocoRepository.Instance.ImportBulkFileLoader(trackListeners); _logger.Debug($"Done importing listener summary for track streams from {trackFile.FilePath}"); MarkAsProcessed(date, trackFile); } else if (trackFile.FileType == SpotifyS3FileType.ListenerSummaryPlaylistTracks) { _logger.Debug($"Reading listener summary for playlist-track streams from {trackFile.FilePath}"); var playlistTrackStreams = ReadNJsonData(localPath, trackFile.DistributorId); _logger.Debug($"Importing listener summary for playlist-track streams from {trackFile.FilePath}"); PetaPocoRepository.Instance.ImportBulkFileLoader(playlistTrackStreams); _logger.Debug($"Done importing listener summary for playlist-track streams from {trackFile.FilePath}"); MarkAsProcessed(date, trackFile); } else if (trackFile.FileType == SpotifyS3FileType.ListenerSummaryPlaylists) { _logger.Debug($"Reading listener summary for playlist streams from {trackFile.FilePath}"); var playlistListeners = ReadNJsonData(localPath, trackFile.DistributorId); _logger.Debug($"Importinglistener summary for playlist streams from {trackFile.FilePath}"); PetaPocoRepository.Instance.ImportBulkFileLoader(playlistListeners); _logger.Debug($"Done importing listener summary for playlist streams from {trackFile.FilePath}"); MarkAsProcessed(date, trackFile); } File.Delete(localPath); } catch (Exception ex) { _logger.Error(ex); } }); }); return taskLog; } private async Task MapToOldStreamSummaryStructureAsync(DateTime trackFileDate, int trackFileDistributorId, string trackFileCountry, List summary) { var playlistData = await GetLocalPlaylistStreamsDataAsync(trackFileDate); var localData = playlistData.FirstOrDefault(p => p.AccountId == trackFileDistributorId && p.Market == trackFileCountry); return new StreamDataSummary() { AccountId = trackFileDistributorId, Date = trackFileDate, Market = trackFileCountry, DeviceDesktopStreams = summary.Sum(s => s.device_type_personal_computer), DeviceMobileStreams = summary.Sum(s => s.device_type_cell_phone), DeviceTabletStreams = summary.Sum(s => s.device_type_tablet), OSAndroidStreams = summary.Sum(s => s.device_os_android), OSMacOSStreams = summary.Sum(s => s.device_os_mac), OSOtherStreams = summary.Sum(s => s.device_os_other), OSWindowsStreams = summary.Sum(s => s.device_os_windows), OSiOSStreams = summary.Sum(s => s.device_os_ios), SourceAlbum = summary.Sum(s => s.source_album), SourceArtist = summary.Sum(s => s.source_artist), SourceCollection = summary.Sum(s => s.source_collection), SourceOther = summary.Sum(s => s.source_other), SourcePlaylist = summary.Sum(s => s.source_others_playlist), SourceSearch = summary.Sum(s => s.source_search), StreamsFromLocalSonyPlaylists = localData?.LocalListStreams ?? 0, StreamsFromSonyPlaylists = localData?.TotalListStreams ?? 0, TotalStreams = summary.Sum(s => s.streams), UniqueUsers = summary.Sum(s => s.listeners), UniqueUsersDeviceDesktop = -1, //TODO: UniqueUsersDeviceMobile = -1, //TODO: UniqueUsersDeviceTablet = -1, //TODO: UniqueUsersSourceAlbum = -1, //TODO: UniqueUsersSourceArtist = -1, //TODO: UniqueUsersSourceCollection = -1, //TODO: UniqueUsersSourceOther = -1, //TODO: UniqueUsersSourceSearch = -1, //TODO: UniqueUsersSourcePlaylist = summary.Sum(s => s.listeners_source_others_playlist), }; } private async Task> GetLocalPlaylistStreamsDataAsync(DateTime summaryDate) { var cacheKey = $"UpdateSpotifyAnalyticsAggregatedTask-{summaryDate.ToString()}"; var localPlaylistStreamsData = MemoryCache.Default.Get(cacheKey) as List; if (localPlaylistStreamsData == null) { localPlaylistStreamsData = await _spotifyAnalyticsManager.GetMarketPlaylistStreamsAsync(summaryDate, (int)StaticBuzzCategory.SonyMusic); MemoryCache.Default.Set(cacheKey, localPlaylistStreamsData, new CacheItemPolicy() { Priority = CacheItemPriority.Default }); } return localPlaylistStreamsData; } private async Task ImportAsync(SpotifyS3File trackFile, string localPath, DateTime date, string typeName) where T : IAggregatedStreamFormat { _logger.Debug($"Reading {typeName} from {trackFile.FilePath}"); var tracks = ReadNJsonData(localPath, trackFile.DistributorId); _logger.Debug($"Importing {typeName} from {trackFile.FilePath}"); const int sizePerBatch = 10000; var batches = tracks.Batch(sizePerBatch).ToList(); await batches.ItemIndex().ForEachAsync(1, async batch => { _logger.Debug($"Importing {typeName} batch {batch.Index} of {batches.Count} from {trackFile.FilePath}"); await PetaPocoRepository.Instance.ImportBulkFileLoaderAsync(batch.Item, MySqlBulkLoaderConflictOption.Replace); _logger.Debug($"Done {typeName} batch {batch.Index} from {trackFile.FilePath}"); }); _logger.Debug($"Done importing {typeName} from {trackFile.FilePath}"); MarkAsProcessed(date, trackFile); } } }