using System; using System.Collections.Generic; using System.Data; using System.Data.Common; using System.Data.Odbc; using System.Linq; using System.Threading.Tasks; using NLog; using Npgsql; using Sony.Filtr.Applications; using Sony.Filtr.Core.SpotifyRegion; using Sony.Filtr.SpotifyWebAPI.Model; using Sony.Filtr.Utility.Extensions; namespace Sony.Filtr.Worker.Onetime { public class CopyS3ToRedshift { private readonly ApplicationInstanceManager _applicationInstanceManager; private readonly SpotifyRegionManager _spotifyRegionManager; private Logger _logger; private const string _endpointAddress = "filtr-analytics.cpe1o4eyrhse.eu-west-1.redshift.amazonaws.com"; private const string _db = "spotify"; private const string _userName = "filtr"; private const string _password = "***REMOVED***"; private const string _port = "5439"; public CopyS3ToRedshift(ApplicationInstanceManager applicationInstanceManager, SpotifyRegionManager spotifyRegionManager) { _applicationInstanceManager = applicationInstanceManager; _spotifyRegionManager = spotifyRegionManager; _logger = NLog.LogManager.GetLogger("CopyS3ToRedshift"); } public async Task ExecuteAsync() { _logger.Info("."); var allRegions = await GetAllRegionsAsync(); await allRegions.ForEachAsync(1, async region => { await CopyMarketStreamsAsync(region); }); } private async Task CopyMarketStreamsAsync(string market) { _logger.Info($"Starting market {market}"); string odbcConnectionString = $"Driver={{PostgreSQL Unicode(x64)}}; Server={_endpointAddress}; Database={_db}; UID={_userName}; PWD={_password}; Port={_port}"; var connectionString = $"Host={_endpointAddress};Username={_userName};Password={_password};Database={_db}"; using (var conn = new OdbcConnection(odbcConnectionString)) { await conn.OpenAsync(); try { using(var command = BuildCopyCommand(market, 4, conn)) { command.CommandTimeout = Int32.MaxValue; await command.ExecuteNonQueryAsync(); } } catch(Exception ex) { _logger.Error(ex); } } _logger.Info($"Done with {market}"); } private async Task> GetAllRegionsAsync() { var applicationRegions = _applicationInstanceManager.GetApplications().Select(a => a.SpotifyRegionCode).Where(r => r.Length == 2).ToList(); var spotifyRegions = await _spotifyRegionManager.GetAvailableRegionsAsync(); //return new List() {"ad", "pt"}; return applicationRegions.Union(spotifyRegions).Select(r => r.ToLowerInvariant()).Distinct().OrderBy(r => r).ToList(); } private DbCommand BuildCopyCommand(string market, int month, OdbcConnection conn) { var tempTableName = $"stream_preload_{market}_{month}"; var query = $"CREATE TABLE {tempTableName} (LIKE stream_raw_format); " + $"COPY {tempTableName} from 's3://filtr-spotify-analytics/streams/{market.ToUpper()}/2017-{month:D2}' " + "access_key_id '***REMOVED***' " + "secret_access_key '***REMOVED***' " + "timeformat 'auto' " + "format as json 'auto' " + "gzip " + "COMPUPDATE off; " + "INSERT INTO stream(user_id, track_id, timestamp, length, source, source_uri, device_type, os, stream_date, market) " + $"SELECT user_id, track_id, (timestamp), length, source, source_uri, device_type, os, DATE_TRUNC('d', timestamp), '{market}' " + $"FROM {tempTableName}; " + $"DROP TABLE {tempTableName}; " + "ANALYZE stream; " + "VACUUM stream; "; var command = new OdbcCommand(query, conn); return command; } } }