diff options
| author | Kapil Viren Ahuja <kvahuja@users.noreply.github.com> | 2017-02-15 09:14:00 +0530 |
|---|---|---|
| committer | GitHub <noreply@github.com> | 2017-02-15 09:14:00 +0530 |
| commit | 00ea8ffaa77efac7ae286359f59187808a13232f (patch) | |
| tree | 3477eb28a6515fabdfd91dffdbcd4dc4293db202 /api.py | |
| parent | ca501a8dab6e931591e43157d53a95f6289b8115 (diff) | |
| parent | 6d2859b3dbee5c22ea7c2cafc17b0d9745c22266 (diff) | |
| download | apigrabber-00ea8ffaa77efac7ae286359f59187808a13232f.tar.gz apigrabber-00ea8ffaa77efac7ae286359f59187808a13232f.zip | |
Merge pull request #12 from vainglorygame/worker_bored
workers sleep and don't die when they have no jobs
Diffstat (limited to 'api.py')
| -rw-r--r-- | api.py | 13 |
1 files changed, 10 insertions, 3 deletions
@@ -63,8 +63,8 @@ class Apigrabber(object): async with con.transaction(): await con.execute(""" INSERT INTO crawljobs(start_date, end_date, finished, region) - SELECT '2017-02-01T00:00:00Z'::TIMESTAMP, - '2017-02-01T00:00:00Z'::TIMESTAMP, + SELECT '2017-02-14T00:00:00Z'::TIMESTAMP, + '2017-02-14T00:00:00Z'::TIMESTAMP, TRUE, region FROM JSONB_TO_RECORDSET($1::JSONB) AS jsn(region TEXT) @@ -148,7 +148,7 @@ class Apigrabber(object): async with con.transaction(isolation="serializable"): # select us our job - jobdate, delta_minutes = await con.fetchrow(""" + row_res = await con.fetchrow(""" SELECT start_date, -- new job's end date LEAST(EXTRACT(EPOCH FROM (start_date-previous_end))/60, $2) -- gap in minutes or default if smaller @@ -163,6 +163,12 @@ class Apigrabber(object): WHERE start_date-previous_end>INTERVAL '0' -- get gaps ORDER BY start_date DESC LIMIT 1 """, region, default_diff) + if row_res == None: + logging.warn("%s: no jobs available. idling.", region) + await asyncio.sleep(60) # a minute TODO make this smarter + asyncio.ensure_future(self.crawl_region(region)) + + jobdate, delta_minutes = row_res delta = datetime.timedelta(minutes=delta_minutes) # store our job as pending jobid = await con.fetchval(""" @@ -211,6 +217,7 @@ class Apigrabber(object): # TODO: respawn a worker if it dies because of connection issues # TODO: insert API version (force update if changed) # TODO: create database indices (id & shardId & type) + # TODO: make workers switch regions flexibly to meet demand? for region in self.regions: await self.request_update(region) |
