diff options
Diffstat (limited to 'api.py')
| -rw-r--r-- | api.py | 62 |
1 files changed, 30 insertions, 32 deletions
@@ -13,30 +13,29 @@ import joblib.worker class Apigrabber(joblib.worker.Worker): def __init__(self, apitoken): self._apitoken = apitoken - self._pool = None - self._insertquery = "" + self._con = None + self._insertquery = None super().__init__(jobtype="grab") async def connect(self, **args): """Connect to database.""" logging.warning("connecting to database") await super().connect(**args) - self._pool = await asyncpg.create_pool(min_size=1, - **args) + self._con = await asyncpg.connect(**args) async def setup(self): """Initialize the database.""" logging.info("initializing database") - async with self._pool.acquire() as con: - await con.execute(""" - CREATE TABLE IF NOT EXISTS - match (id TEXT PRIMARY KEY, data JSONB) - """) + await self._con.execute(""" + CREATE TABLE IF NOT EXISTS + match (id TEXT PRIMARY KEY, data JSONB) + """) root = os.path.realpath( os.path.join(os.getcwd(), os.path.dirname(__file__))) with open(root + "/insert.sql", "r", encoding="utf-8-sig") as file: - self._insertquery = file.read() + self._insertquery = await self._con.prepare( + file.read()) async def _execute_job(self, jobid, payload, priority): """Finish a job.""" @@ -47,28 +46,27 @@ class Apigrabber(joblib.worker.Worker): else: playername = "" - async with self._pool.acquire() as con: - logging.debug("%s: running on %s with parameters '%s'", - jobid, payload["region"], payload["params"]) - try: - async for data in api.matches(region=payload["region"], - params=payload["params"]): - async with con.transaction(): - matchids = await con.fetch( - self._insertquery, json.dumps(data)) - logging.info("%s: inserted %s matches for player '%s' from API into database", - jobid, len(matchids), playername) - payloads = [{ - "id": mat["id"], - "playername": playername - } for mat in matchids] - await self._queue.request(jobtype="process", - payload=payloads, - priority=priority) - except crawler.ApiError as error: - logging.warning("%s: API returned error '%s'", - jobid, error.args[0]) - raise joblib.worker.JobFailed(error.args[0]) + logging.debug("%s: running on %s with parameters '%s'", + jobid, payload["region"], payload["params"]) + try: + async for data in api.matches(region=payload["region"], + params=payload["params"]): + async with self._con.transaction(): + matchids = await self._insertquery.fetch( + json.dumps(data)) + logging.info("%s: inserted %s matches for player '%s' from API into database", + jobid, len(matchids), playername) + payloads = [{ + "id": mat["id"], + "playername": playername + } for mat in matchids] + await self._queue.request(jobtype="process", + payload=payloads, + priority=priority) + except crawler.ApiError as error: + logging.warning("%s: API returned error '%s'", + jobid, error.args[0]) + raise joblib.worker.JobFailed(error.args[0]) async def startup(): |
