From 612f1ad9a38c065641259412828a85f333080f89 Mon Sep 17 00:00:00 2001 From: schneefux Date: Mon, 20 Feb 2017 19:49:58 +0100 Subject: specify region and crawl times per env vars --- api.py | 47 ++++++++++++++++++++++++++++++++--------------- 1 file changed, 32 insertions(+), 15 deletions(-) (limited to 'api.py') diff --git a/api.py b/api.py index d178828..db4fa6e 100644 --- a/api.py +++ b/api.py @@ -14,8 +14,8 @@ import crawler # SEMC API is a bit strict about the iso format def date2iso(d): """Convert datetime to iso8601 string.""" - date = d.isoformat() - date = ".".join(date.split(".")[:-1]) # remove microseconds + date = d.replace(microsecond=0) + date = date.isoformat() date = date + "Z" return date @@ -28,9 +28,19 @@ def iso2date(d): class Apigrabber(object): - def __init__(self): - #self.regions = ["na", "eu", "sg", "ea", "sa", "cn"] - self.regions = ["na", "eu"] + def __init__(self, regions, first_fetch=None, last_fetch=None): + self.update_live = False + # TODO refactor this + if first_fetch is None: + first_fetch = "2017-02-14T00:00:00Z" + if last_fetch is None: # update this in 53 years + last_fetch = "2070-01-01T00:00:00Z" + self.update_live = True + + # TODO first_fetch parameter is only valid on fresh databases + self.first_fetch = iso2date(first_fetch) # when to end fetching history + self.last_fetch = iso2date(last_fetch) # when to start fetching + self.regions = regions async def connect(self, **args): self._pool = await asyncpg.create_pool(**args) @@ -47,13 +57,11 @@ class Apigrabber(object): async with con.transaction(): await con.execute(""" INSERT INTO crawljobs(start_date, end_date, finished, region) - SELECT '2017-02-14T00:00:00Z'::TIMESTAMP, - '2017-02-14T00:00:00Z'::TIMESTAMP, - TRUE, - region + SELECT $2, $2, TRUE, region FROM JSONB_TO_RECORDSET($1::JSONB) AS jsn(region TEXT) ON CONFLICT DO NOTHING; - """, json.dumps([{"region": r} for r in self.regions])) + """, json.dumps([{"region": r} for r in self.regions]), + self.first_fetch) async def _db_insert(self, con, objects): """Insert a list of API response objects into respective tables.""" @@ -181,12 +189,16 @@ class Apigrabber(object): try: async with con.transaction(isolation="serializable"): # insert a job from now->now + # (or last_fetch->last_fetch if asked) # so history crawler picks up the time diff between # now and the last query. await con.fetchval(""" - INSERT INTO crawljobs(start_date, end_date, finished, region) - VALUES (NOW(), NOW(), TRUE, $1) - """, region) + INSERT INTO crawljobs + (start_date, end_date, finished, region) + VALUES ( + LEAST(NOW(), $2), + LEAST(NOW(), $2), TRUE, $1) + """, region, self.last_fetch) logging.info("%s: scheduled for live update", region) break except asyncpg.exceptions.SerializationError: @@ -195,7 +207,8 @@ class Apigrabber(object): async def _recall_later(): await asyncio.sleep(300) # wait 5 min and repeat asyncio.ensure_future(self.request_update(region)) - asyncio.ensure_future(_recall_later()) + if self.update_live: + asyncio.ensure_future(_recall_later()) async def start(self): @@ -219,7 +232,11 @@ class Apigrabber(object): async def startup(): - apigrabber = Apigrabber() + apigrabber = Apigrabber( + regions=os.environ["VAINSOCIAL_REGIONS"].split(","), + first_fetch=os.environ.get("VAINSOCIAL_STARTDATE"), + last_fetch=os.environ.get("VAINSOCIAL_ENDDATE") + ) await apigrabber.connect( host=os.environ["POSTGRESQL_HOST"], port=os.environ["POSTGRESQL_PORT"], -- cgit v1.3.1