1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
|
#!/usr/bin/python
import os
import logging
import asyncio
import asyncpg
source_db = {
"host": os.environ.get("POSTGRESQL_SOURCE_HOST") or "localhost",
"port": os.environ.get("POSTGRESQL_SOURCE_PORT") or 5532,
"user": os.environ.get("POSTGRESQL_SOURCE_USER") or "vgstats",
"password": os.environ.get("POSTGRESQL_SOURCE_PASSWORD") or "vgstats",
"database": os.environ.get("POSTGRESQL_SOURCE_DB") or "vgstats"
}
dest_db = {
"host": os.environ.get("POSTGRESQL_DEST_HOST") or "localhost",
"port": os.environ.get("POSTGRESQL_DEST_PORT") or 5432,
"user": os.environ.get("POSTGRESQL_DEST_USER") or "vainsocial",
"password": os.environ.get("POSTGRESQL_DEST_PASSWORD") or "vainsocial",
"database": os.environ.get("POSTGRESQL_DEST_DB") or "vainsocial"
}
class Database(object):
async def connect(self, **connect_kwargs):
"""Connect to database by arguments."""
self._pool = await asyncpg.create_pool(**connect_kwargs)
async def each(self, fetchquery, func, **func_args):
"""Execute a function for each row that is fetched with query."""
# TODO use iterator instead of callback
async with self._pool.acquire() as conn:
async with conn.transaction():
tasks = []
async for record in conn.cursor(fetchquery):
tasks.append(
asyncio.ensure_future(func(record, **func_args))
)
await asyncio.gather(*tasks)
async def into(self, data, table):
"""Insert a named tuple into a database."""
items = list(data.items())
keys, values = [x[0] for x in items], [x[1] for x in items]
placeholders = ["${}".format(i) for i, _ in enumerate(values, 1)]
query = "INSERT INTO {} (\"{}\") VALUES ({})".format(
table, "\", \"".join(keys), ", ".join(placeholders))
async with self._pool.acquire() as conn:
await conn.fetch(query, (*data))
async def process():
sdb = Database()
await sdb.connect(**source_db)
ddb = Database()
await ddb.connect(**dest_db)
selqry = """
SELECT
(data->'attributes'->>'duration')::int AS "duration",
data->'attributes'->>'gameMode' AS "gameMode",
COALESCE(NULLIF(data->'attributes'->>'patchVersion', ''), '0')::int AS "patchVersion",
data->'attributes'->>'shardId' AS "shard",
data->'attributes'->'stats'->>'endGameReason' AS "result",
data->'attributes'->'stats'->>'queue' AS "queue",
false AS "anyAFK",
'' AS "winningTeam",
0 AS "krakenCaptures",
0 AS "laneMinionsSlayed",
0 AS "jungleMinionsSlayed",
0 AS "heroDeaths"
FROM match LIMIT 100
"""
await sdb.each(selqry, ddb.into, table="match")
if __name__ == "__main__":
loop = asyncio.get_event_loop()
loop.run_until_complete(
process()
)
|