summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-03-27 20:58:21 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-03-27 20:58:21 +0200
commit3c7b19d77476d774762b4a6c97bff18bcd0b249a (patch)
tree32fbc9f93eacbc46b073920801b5ee79ba05fa0e
parent60a8aac5f7acc61432fc6a59843e2c91d5df6a35 (diff)
downloadapigrabber-3c7b19d77476d774762b4a6c97bff18bcd0b249a.tar.gz
apigrabber-3c7b19d77476d774762b4a6c97bff18bcd0b249a.zip
use 2.0 joblib
-rw-r--r--Dockerfile5
-rw-r--r--api.py92
-rw-r--r--cli.py40
-rw-r--r--crawler.py72
-rw-r--r--insert.sql74
-rw-r--r--requirements.txt12
-rw-r--r--worker.py55
7 files changed, 136 insertions, 214 deletions
diff --git a/Dockerfile b/Dockerfile
new file mode 100644
index 0000000..fb20bea
--- /dev/null
+++ b/Dockerfile
@@ -0,0 +1,5 @@
+FROM python:3.6-alpine
+ADD requirements.txt /code/requirements.txt
+WORKDIR /code
+RUN pip install -r requirements.txt
+CMD ["python", "worker.py"]
diff --git a/api.py b/api.py
deleted file mode 100644
index 71ec6a3..0000000
--- a/api.py
+++ /dev/null
@@ -1,92 +0,0 @@
-#!/usr/bin/python
-
-import asyncio
-import os
-import logging
-import json
-import asyncpg
-
-import crawler
-import joblib.worker
-
-
-class Apigrabber(joblib.worker.Worker):
- def __init__(self, apitoken):
- self._apitoken = apitoken
- 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._con = await asyncpg.connect(**args)
-
- async def setup(self):
- """Initialize the database."""
- logging.info("initializing database")
- 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 = await self._con.prepare(
- file.read())
-
- async def _execute_job(self, jobid, payload, priority):
- """Finish a job."""
- api = crawler.Crawler(self._apitoken)
- # if a player is queried, pass that information to processor
- if "filter[playerNames]" in payload["params"]:
- playername = payload["params"]["filter[playerNames]"]
- else:
- playername = ""
-
- 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.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],
- False) # not critical
-
-
-async def startup():
- worker = Apigrabber(
- apitoken=os.environ["VAINSOCIAL_APITOKEN"]
- )
- await worker.connect(
- host=os.environ["POSTGRESQL_HOST"],
- port=os.environ["POSTGRESQL_PORT"],
- user=os.environ["POSTGRESQL_USER"],
- password=os.environ["POSTGRESQL_PASSWORD"],
- database=os.environ["POSTGRESQL_DB"]
- )
- await worker.setup()
- await worker.run(batchlimit=1)
-
-
-if __name__ == "__main__":
- logging.basicConfig(level=logging.DEBUG)
-
- loop = asyncio.get_event_loop()
- loop.run_until_complete(startup())
diff --git a/cli.py b/cli.py
new file mode 100644
index 0000000..f26901a
--- /dev/null
+++ b/cli.py
@@ -0,0 +1,40 @@
+#!/usr/bin/python3
+import os
+import argparse
+
+import joblib.joblib
+
+
+RABBIT = {
+ "host": os.environ.get("RABBITMQ_HOST"),
+ "port": os.environ.get("RABBITMQ_PORT"),
+ "credentials": os.environ.get("RABBITMQ_CREDS")
+}
+
+
+def main(args):
+ queue = joblib.joblib.JobQueue()
+ queue.connect(**RABBIT)
+
+ if args.player:
+ payload = {
+ "region": args.region,
+ "params": {
+ "filter[createdAt-start]": "2017-02-01T00:00:00Z",
+ "filter[playerNames]": args.player
+ }
+ }
+ queue.request("grab", payload)
+
+parser = argparse.ArgumentParser(description="Request a Vainsocial update.")
+parser.add_argument("-n", "--player",
+ help="Player name",
+ type=str)
+parser.add_argument("-r", "--region",
+ help="Specify a region",
+ type=str,
+ choices=["na", "eu", "sg"],
+ default="na")
+
+if __name__ == "__main__":
+ main(parser.parse_args())
diff --git a/crawler.py b/crawler.py
index f8fd279..84e4692 100644
--- a/crawler.py
+++ b/crawler.py
@@ -1,9 +1,9 @@
#!/usr/bin/python
import json
-import asyncio
+import time
import logging
-import aiohttp
+import requests
APIURL = "https://api.dc01.gamelockerapp.com/"
@@ -18,11 +18,9 @@ class Crawler(object):
self._token = token
self._pagelimit = 50
- async def _req(self, session, path, params):
+ def _req(self, path, params):
"""Sends an API request and returns the response dict.
- :param session: aiohttp client session.
- :type session: :class:`aiohttp.ClientSession`
:param path: URL path.
:type path: str
:param params: Request parameters.
@@ -39,23 +37,19 @@ class Crawler(object):
retries = 5
while True:
try:
- async with session.get(self._apiurl + path, headers=headers,
- params=params) as response:
- if response.status == 429:
- logging.warning("rate limited, retrying")
+ response = requests.get(self._apiurl + path,
+ headers=headers,
+ params=params)
+ if response.status_code == 429:
+ logging.warning("rate limited, retrying")
+ else:
+ if response.status_code > 500:
+ logging.error("API server error %s",
+ response.status_code)
+ raise ApiError(response.status_code)
else:
- if response.status > 500:
- logging.error("API server error %s",
- response.status)
- raise ApiError(response.status)
- else:
- return await response.json()
- except (aiohttp.errors.ContentEncodingError,
- aiohttp.errors.ServerDisconnectedError,
- aiohttp.errors.ClientResponseError,
- aiohttp.errors.ClientOSError,
- LookupError,
- json.decoder.JSONDecodeError) as err:
+ return response.json()
+ except (json.decoder.JSONDecodeError) as err:
# API bug?
logging.error("API error '%s', retrying", err)
retries -= 1
@@ -63,9 +57,9 @@ class Crawler(object):
logging.error("Giving up")
raise ApiError(str(err))
- await asyncio.sleep(5)
+ time.sleep(5)
- async def matches(self, params, region="na"):
+ def matches(self, params, region="na"):
"""Queries the API for matches and their related data.
:param region: (optional) Region where the matches were played.
@@ -76,23 +70,21 @@ class Crawler(object):
"""
params["page[offset]"] = 0
params["page[limit]"] = self._pagelimit
- async with aiohttp.ClientSession() as session:
- while True:
- res = await self._req(session,
- "shards/" + region + "/matches",
- params)
+ while True:
+ res = self._req("shards/" + region + "/matches",
+ params)
- if "errors" in res:
- if res["errors"][0].get("title") == "Not Found" \
- and params["page[offset]"] > 0:
- # a query returned exactly 50 matches
- # which is expected, so don't fail.
- return
- raise ApiError(res["errors"])
+ if "errors" in res:
+ if res["errors"][0].get("title") == "Not Found" \
+ and params["page[offset]"] > 0:
+ # a query returned exactly 50 matches
+ # which is expected, so don't fail.
+ return
+ raise ApiError(res["errors"])
- yield res
+ yield res
- if len(res["data"]) < self._pagelimit:
- # asked for 50, got less -> exhausted
- break
- params["page[offset]"] += params["page[limit]"]
+ if len(res["data"]) < self._pagelimit:
+ # asked for 50, got less -> exhausted
+ break
+ params["page[offset]"] += params["page[limit]"]
diff --git a/insert.sql b/insert.sql
deleted file mode 100644
index bc6485d..0000000
--- a/insert.sql
+++ /dev/null
@@ -1,74 +0,0 @@
-WITH
--- parse source string
-srcjson AS (
- SELECT $1::JSONB AS data
-),
--- split into data / included
-matches AS (
- SELECT JSONB_ARRAY_ELEMENTS(srcjson.data->'data') AS data FROM srcjson
-),
-includes AS (
- SELECT JSONB_ARRAY_ELEMENTS(srcjson.data->'included') AS data FROM srcjson
-),
--- filter included by type
-rosters AS (
- SELECT includes.data AS data FROM includes WHERE includes.data->>'type'='roster'
-),
-participants AS (
- SELECT includes.data AS data FROM includes WHERE includes.data->>'type'='participant'
-),
-players AS (
- SELECT includes.data AS data FROM includes WHERE includes.data->>'type'='player'
-),
--- cleanup players
-linked_players AS (
- SELECT
- JSONB_BUILD_OBJECT(
- 'data', players.data-'relationships'
- ) AS data
- FROM players
-),
--- link participant-player
-linked_participants AS (
- SELECT
- JSONB_BUILD_OBJECT(
- 'data', participants.data-'relationships',
- 'relations', TO_JSONB(ARRAY(
- SELECT * FROM linked_players
- WHERE participants.data->'relationships'->'player'->'data' = JSONB_BUILD_OBJECT('id', linked_players.data->'data'->>'id', 'type', linked_players.data->'data'->>'type')
- ))
- ) AS data
- FROM participants
-),
- -- link roster-participants
-linked_rosters AS (
- SELECT
- JSONB_BUILD_OBJECT(
- 'data', rosters.data-'relationships',
- 'relations', TO_JSONB(ARRAY(
- SELECT * FROM linked_participants
- WHERE rosters.data->'relationships'->'participants'->'data' @> JSONB_BUILD_ARRAY(JSONB_BUILD_OBJECT('id', linked_participants.data->'data'->>'id', 'type', linked_participants.data->'data'->>'type'))
- ))
- ) AS data
- FROM rosters
-),
--- link match-rosters
-linked_matches AS (
- SELECT
- matches.data->>'id' AS id,
- JSONB_BUILD_OBJECT(
- 'data', matches.data-'relationships',
- 'relations', TO_JSONB(ARRAY(
- SELECT * FROM linked_rosters
- WHERE matches.data->'relationships'->'rosters'->'data' @> JSONB_BUILD_ARRAY(JSONB_BUILD_OBJECT('id', linked_rosters.data->'data'->>'id', 'type', linked_rosters.data->'data'->>'type'))
- ))
- ) AS data
- FROM matches
-),
--- insert!
-insert_matches AS (
- INSERT INTO match(id, data) SELECT * FROM linked_matches
- ON CONFLICT(id) DO NOTHING
- RETURNING id
-)
-SELECT DISTINCT id FROM insert_matches
diff --git a/requirements.txt b/requirements.txt
index 7d8614e..b01174d 100644
--- a/requirements.txt
+++ b/requirements.txt
@@ -1,10 +1,6 @@
-aiohttp==1.2.0
-appdirs==1.4.0
-async-timeout==1.1.0
-asyncpg==0.8.4
-chardet==2.3.0
-multidict==2.1.4
+appdirs==1.4.3
packaging==16.8
-pyparsing==2.1.10
+pika==0.10.0
+pyparsing==2.2.0
+requests==2.13.0
six==1.10.0
-yarl==0.9.0
diff --git a/worker.py b/worker.py
new file mode 100644
index 0000000..d4c5fab
--- /dev/null
+++ b/worker.py
@@ -0,0 +1,55 @@
+#!/usr/bin/python
+
+import os
+import logging
+
+import crawler
+import joblib.joblib
+
+
+RABBIT = {
+ "host": os.environ.get("RABBITMQ_HOST"),
+ "port": os.environ.get("RABBITMQ_PORT"),
+ "credentials": os.environ.get("RABBITMQ_CREDS")
+}
+
+APITOKEN = os.environ["MADGLORY_TOKEN"]
+
+
+class Apigrabber(joblib.joblib.Worker):
+ def __init__(self, apitoken):
+ super().__init__("grab")
+ self._apitoken = apitoken
+
+ def work(self, payload):
+ """Finish a job."""
+ api = crawler.Crawler(self._apitoken)
+ logging.info("running on %s with parameters '%s'",
+ payload["region"], payload["params"])
+ try:
+ for data in api.matches(region=payload["region"],
+ params=payload["params"]):
+ items = data["data"] + data["included"]
+ for item in items:
+ self.request("process",
+ payload={
+ "id": item["id"],
+ "type": item["type"],
+ "data": item
+ })
+ except crawler.ApiError as error:
+ logging.warning("API returned error '%s'", error.args[0])
+ raise joblib.JobFailed(error.args[0],
+ False) # not critical
+
+
+def startup():
+ worker = Apigrabber(APITOKEN)
+ worker.connect(**RABBIT)
+ worker.setup()
+ worker.run()
+
+
+if __name__ == "__main__":
+ logging.basicConfig(level=logging.INFO)
+ startup()