using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Globalization; using System.IO; using System.IO.Compression; using System.Linq; using System.Threading.Tasks; using Amazon.S3; using Amazon.S3.Model; using Amazon.S3.Transfer; using MoreLinq; using MySql.Data.MySqlClient; using Newtonsoft.Json; using NLog; using PetaPoco.Business; using Sony.Filtr.AppleMusic; using Sony.Filtr.AppleMusic.Data.Streams; using Sony.Filtr.AppleMusic.Playlists; using Sony.Filtr.AppleMusic.Streams; using Sony.Filtr.Contracts.Abstractions; using Sony.Filtr.Contracts.Definitions; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.Tasks.Tasks.AppleMusic { public class UpdateAppleMusicAnalyticsAggregatedTask : IScheduledTask { private readonly AppleMusicStreamingAggregatedReportApi _aggregatedReportApi; private readonly AppleMusicPlaylistManager _appleMusicPlaylistManager; private readonly AppleMusicStreamsManager _appleMusicStreamsManager; private readonly Logger _logger; private const string dataFileFolderName = "appleMusicAnalytics"; public UpdateAppleMusicAnalyticsAggregatedTask(AppleMusicStreamingAggregatedReportApi aggregatedReportApi, AppleMusicPlaylistManager appleMusicPlaylistManager, AppleMusicStreamsManager appleMusicStreamsManager) { _aggregatedReportApi = aggregatedReportApi; _appleMusicPlaylistManager = appleMusicPlaylistManager; _appleMusicStreamsManager = appleMusicStreamsManager; _logger = LogManager.GetLogger("UpdateAppleMusicAnalyticsAggregatedTask"); } public async Task ExecuteAsync(Guid scheduledTaskLogId) { var taskLog = new ScheduledTaskLog(); if (!Directory.Exists(dataFileFolderName)) { Directory.CreateDirectory(dataFileFolderName); } var appleMusicStartDate = new DateTime(2016, 04, 27); var endDate = DateTime.Today; List supportedFileTypes = new List() { //AppleS3FileType.Demographics, //AppleS3FileType.StreamsTracks, //AppleS3FileType.StreamsContainerTracks, AppleS3FileType.StreamsContainers, }; var dates = appleMusicStartDate.GetDateRangeTo(endDate).Reverse(); var shouldCalculate = false; await dates.ForEachAsync(1, async date => { var files = await _aggregatedReportApi.GetFilesAsync(date); var alreadyImportedFiles = PetaPocoRepository.Instance.Fetch("WHERE date=@0", date); var trackFiles = files.Where(p => supportedFileTypes.Contains(p.FileType)).Where(p => !alreadyImportedFiles.Any(a => a.CountryCode == p.Country && a.DistributerId == p.DistributorId && a.VendorId == p.VendorId && a.FileType == p.FileType)).ToList(); _logger.Debug($"Found {trackFiles.Count} for {date.ToShortDateString()}"); await trackFiles.ForEachAsync(1, async trackFile => { try { _logger.Debug($"Downloading {trackFile.FilePath}"); var dateFilenamePart = trackFile.Date.ToString("yyyy-dd-M--HH-mm-ss"); var localPath = $"{dataFileFolderName}/track-{dateFilenamePart}-{trackFile.DistributorId}-{trackFile.VendorId}-{trackFile.Country}-{trackFile.Version}-{trackFile.FileType}.ndjson"; await _aggregatedReportApi.DownloadToFileAsync(localPath, trackFile.FilePath); if (trackFile.FileType == AppleS3FileType.Demographics) { _logger.Debug($"Reading demographics from {trackFile.FilePath}"); var trackDemographics = ReadNJsonData(localPath, trackFile.DistributorId); _logger.Debug($"Importing demographics from {trackFile.FilePath}"); PetaPocoRepository.Instance.ImportBulkFileLoader(trackDemographics); _logger.Debug($"Done importing demographicsfrom {trackFile.FilePath}"); MarkAsProcessed(date, trackFile); } else if (trackFile.FileType == AppleS3FileType.StreamsTracks) { _logger.Debug($"Reading tracks from {trackFile.FilePath}"); var trackStreams = ReadNJsonData(localPath, trackFile.DistributorId); _logger.Debug($"Importing tracks from {trackFile.FilePath}"); int batchIndex = 1; foreach(var trackBatch in trackStreams.Batch(100000)) { PetaPocoRepository.Instance.ImportBulkFileLoader(trackBatch); _logger.Debug($"Done import track batch {batchIndex} from {trackFile.FilePath}"); batchIndex++; } MarkAsProcessed(date, trackFile); } else if (trackFile.FileType == AppleS3FileType.StreamsContainerTracks) { _logger.Debug($"Reading track-container streams from {trackFile.FilePath}"); var trackStreams = ReadNJsonData(localPath, trackFile.DistributorId); _logger.Debug($"Importing track-container streams from {trackFile.FilePath}"); PetaPocoRepository.Instance.ImportBulkFileLoader(trackStreams); _logger.Debug($"Done importing track-container streams from {trackFile.FilePath}"); MarkAsProcessed(date, trackFile); } else if (trackFile.FileType == AppleS3FileType.StreamsContainers) { _logger.Debug($"Reading container streams from {trackFile.FilePath}"); var containerStreams = ReadNJsonData(localPath, trackFile.DistributorId); _logger.Debug($"Importing container streams from {trackFile.FilePath}"); PetaPocoRepository.Instance.ImportBulkFileLoader(containerStreams); _logger.Debug($"Done importing container streams from {trackFile.FilePath}"); if (containerStreams != null && containerStreams.Any()) { shouldCalculate = true; } MarkAsProcessed(date, trackFile); } File.Delete(localPath); } catch(Exception ex) { _logger.Error(ex); } }); }); if (shouldCalculate) { await SetContainerStreamsSummaryAsync(); } return taskLog; } private static void MarkAsProcessed(DateTime date, AppleMusicAggregatedS3File trackFile) { PetaPocoRepository.Instance.Insert(new AppleMusicAnalyticsAggregatedDate() { Date = date, CountryCode = trackFile.Country, VendorId = trackFile.VendorId, DistributerId = trackFile.DistributorId, FileType = trackFile.FileType, }); } private List ReadNJsonData(string localPath, int trackFileDistributorId) where T : class, IAggregatedStreamFormat { List objects = new List(); var jsonSerializer = new JsonSerializer(); using (var fileStream = File.OpenRead(localPath)) { using (var jsonReader = new JsonTextReader(new StreamReader(fileStream)) { SupportMultipleContent = true, }) { while (jsonReader.Read()) { try { var obj = jsonSerializer.Deserialize(jsonReader); obj.DistributerId = trackFileDistributorId; objects.Add(obj); } catch (Exception e) { _logger.Error(e); } } } } return objects; } private async Task SetContainerStreamsSummaryAsync() { _logger.Debug("Begin SetContainerStreamsSummaryAsync"); var playlists = await _appleMusicPlaylistManager.GetPlaylistsAsync(); var containerIds = playlists.Items.Select(p => p.Id).ToList(); _logger.Debug($"Will calculate container summary for {containerIds.Count()} containers"); ConcurrentBag allContainerStreamsSummaries = new ConcurrentBag(); var batches = containerIds.Batch(10).ToList(); var latestDate = await _appleMusicStreamsManager.GetLatestDayWithAggregatedStreamsDataAsync(AppleS3FileType.StreamsContainers); await batches.ItemIndex().ForEachAsync(5, async playlistBatch => { Console.WriteLine($"Calculating for batch {playlistBatch.Index} / {batches.Count}"); var summaries = await _appleMusicStreamsManager.CalculateAppleMusicContainerStreamsSummariesAsync(playlistBatch.Item, latestDate); summaries.ForEach(s => allContainerStreamsSummaries.Add(s)); }); _logger.Debug($"Done calculating container summaries"); _logger.Debug($"Bulk loading to db..."); PetaPocoRepository.Instance.ImportBulkFileLoader(allContainerStreamsSummaries, MySqlBulkLoaderConflictOption.Replace); _logger.Debug("Done SetContainerStreamsSummaryAsync"); } } public class AppleMusicStreamingAggregatedReportApi { private readonly AmazonS3Client _s3Client; private const string _Bucket = "sme-aggregated-stream-reports"; public AppleMusicStreamingAggregatedReportApi(AmazonS3Client s3Client) { _s3Client = s3Client; } public async Task> GetFilesAsync(DateTime date) { List filenames = new List(); var dateString = date.ToString("yyyy-MM-dd", CultureInfo.InvariantCulture); var streamObjects = await _s3Client.ListObjectsAsync(new ListObjectsRequest() { BucketName = _Bucket, Prefix = $"apple/{dateString}/", }); filenames.AddRange(streamObjects.S3Objects.Select(p=> p.Key).ToList()); var nextMarker = streamObjects.NextMarker; while (!string.IsNullOrWhiteSpace(nextMarker)) { var moreResults = await _s3Client.ListObjectsAsync(new ListObjectsRequest() { BucketName = _Bucket, Prefix = $"apple/{dateString}/", Marker = nextMarker }); nextMarker = moreResults.NextMarker; filenames.AddRange(moreResults.S3Objects.Select(p => p.Key).ToList()); } var fileNames = filenames.Select(p => ParseAppleMusicS3FilePath(p)).Where(p => p != null).ToList(); return fileNames; } public async Task DownloadToFileAsync(string destinationFilepath, string fileKey) { using (var downloadStream = await DownloadS3FileAsync(_Bucket, fileKey)) { using (var fileStream = File.Open(destinationFilepath, FileMode.OpenOrCreate)) { using(var gunzippedStream = new GZipStream(downloadStream, CompressionMode.Decompress)) { await gunzippedStream.CopyToAsync(fileStream); } } } } private async Task DownloadS3FileAsync(string bucket, string key) { var transferUtility = new TransferUtility(_s3Client); var contentStream = await transferUtility.OpenStreamAsync(new TransferUtilityOpenStreamRequest() { BucketName = bucket, Key = key, }); return contentStream; } private AppleMusicAggregatedS3File ParseAppleMusicS3FilePath(string filename) { //apple_2017-12-05_theorchard_v1_85420853_streams_tracks.ndjson try { var parsedFilename = Path.GetFileNameWithoutExtension(filename); if (string.IsNullOrWhiteSpace(parsedFilename)) return null; var fileParts = parsedFilename.Split('_'); if (!fileParts.Any()) return null; if (fileParts.ElementAt(0) != "apple") return null; var date = fileParts.ElementAt(1); var distributorName = fileParts.ElementAt(2); var version = fileParts.ElementAt(3); var vendorId = int.Parse(fileParts.ElementAt(4)); string country = "global"; AppleS3FileType fileType = AppleS3FileType.Unknown; var type = fileParts.ElementAt(5); if (type == "streams") { var type2 = fileParts.ElementAt(6); if (type2 == "tracks") { fileType = AppleS3FileType.StreamsTracks; } else if (type2 == "summary") { fileType = AppleS3FileType.StreamsSummary; } else if (type2 == "containers") { fileType = AppleS3FileType.StreamsContainers; } else if (type2 == "container") { fileType = AppleS3FileType.StreamsContainerTracks; } } else if (type == "demographics") { fileType = AppleS3FileType.Demographics; } var distributorId = _distributerMapping[distributorName]; return new AppleMusicAggregatedS3File() { Date = DateTime.Parse(date), DistributorId = distributorId, Country = country, VendorId = vendorId, Version = version, FilePath = filename, FileType = fileType }; } catch(Exception e) { return null; } } private readonly Dictionary _distributerMapping = new Dictionary() { { "sony", (int)SpotifyAnalyticsAccount.Sony }, { "theorchard", (int)SpotifyAnalyticsAccount.Orchard }, { "smej", (int)SpotifyAnalyticsAccount.SonyMusicEntertainmentJapan }, { "smejintl", (int)SpotifyAnalyticsAccount.SonyMusicEntertainmentJapanInternational }, }; } }