""" Tornado App to pull RT data""" import json import os import momoko import pprint import psycopg2 import tornado.ioloop import tornado.web from datetime import date, datetime from psycopg2.extras import RealDictCursor from tornado import httpserver from tornado import gen from tornado.options import define, options import dma class DateTimeEncoder(json.JSONEncoder): def default(self, o): if isinstance(o, datetime) or isinstance(o, date): return o.isoformat() return json.JSONEncoder.default(self, o) class BaseHandler(tornado.web.RequestHandler): @property def db(self): return self.application.db class TrackPlaysHandler(BaseHandler): @gen.coroutine def get(self): isrc = self.get_argument('isrc') dt = self.get_argument('dt', '2015-12-01') #USA561410900 - ted leo me & mia query = ( "SELECT " "os, device_type, user_country, user_region, zipcode, " "CASE WHEN gender = 'female' THEN 'F' WHEN gender = 'male' THEN 'M' ELSE 'U' END as gender, " "DATEDIFF(year, TO_DATE(birth_year, 'YYYY'), GETDATE()) as age, " "track_name, download_date, tmstamp, to_char(tmstamp, 'HH24:MI:SS') as time " "FROM production.staging_raw_spotify_v2 " "WHERE download_date = '{dt}' and isrc = '{isrc}' " "AND user_country = 'US' " "ORDER BY tmstamp asc".format(isrc=isrc, dt=dt)) try: cursor = yield self.db.execute(query) except psycopg2.Error as error: self.write(str(error)) else: results = [] for row in cursor.fetchall(): if row['user_region'] in dma.mapping: row.update(dma.mapping[row['user_region']]) results.append(row) row['age_bucket'] = bucket_age(row['age']) self.write(json.dumps(results, cls=DateTimeEncoder)) self.finish() def set_default_headers(self): self.set_header('Content-Type', 'application/json') class TrackInfoHandler(BaseHandler): @gen.coroutine def get(self): isrc = self.get_argument("isrc") #USA561410900 - ted leo me & mia query = ( "SELECT da.artistname as artist, di.isrcname as track_name " "FROM dim_isrc di " "JOIN dim_release dr on di.upc = dr.display_upc " "JOIN dim_artist da on dr.artistid = da.artistid " "WHERE di.isrc = '{isrc}' ".format(isrc=isrc)) try: cursor = yield self.db.execute(query) except psycopg2.Error as error: self.write(str(error)) else: self.write(json.dumps(cursor.fetchall()[0], cls=DateTimeEncoder)) self.finish() def set_default_headers(self): self.set_header('Content-Type', 'application/json') class IndexHandler(tornado.web.RequestHandler): def get(self): self.render("index.html") class TestHandler(BaseHandler): def get(self): self.write('Some text here!') self.finish() def bucket_age(age): try: if age >= 12 and age <= 17: bucket = 1 elif age >= 18 and age <= 24: bucket = 2 elif age >= 25 and age <= 34: bucket = 3 elif age >= 35 and age <= 44: bucket = 4 elif age >= 45 and age <= 54: bucket = 5 elif age >= 55 and age <= 64: bucket = 6 elif age > 65: bucket = 7 else: bucket = 0 except TypeError: bucket = 0 return bucket if __name__ == "__main__": app = tornado.web.Application( [ (r'/track_plays', TrackPlaysHandler), (r'/track_info', TrackInfoHandler), (r'/', IndexHandler), (r'/test', TestHandler)], template_path=os.path.join(os.path.dirname(__file__), "templates"), static_path=os.path.join(os.path.dirname(__file__), "static"), debug=True) ioloop = tornado.ioloop.IOLoop.instance() db = os.environ.get('REDSHIFT_DB') user = os.environ.get('REDSHIFT_USER') pwd = os.environ.get('REDSHIFT_PASSWORD') host = os.environ.get('REDSHIFT_HOST') port = os.environ.get('REDSHIFT_PORT') dsn = ( 'dbname={db} user={user} password={pwd} ' 'host={host} port={port}').format( db=db, user=user, pwd=pwd, host=host, port=port) app.db = momoko.Pool( dsn=dsn, size=1, ioloop=ioloop, cursor_factory=RealDictCursor ) # this is a one way to run ioloop in sync future = app.db.connect() ioloop.add_future(future, lambda f: ioloop.stop()) ioloop.start() future.result() # raises exception on connection error http_server = httpserver.HTTPServer(app) http_server.listen(8888, '0.0.0.0') #http_server.listen(8889, '10.70.0.215') ioloop.start()