summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--api.py62
1 files changed, 30 insertions, 32 deletions
diff --git a/api.py b/api.py
index 692c931..19ea541 100644
--- a/api.py
+++ b/api.py
@@ -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():