-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathagents.py
More file actions
68 lines (48 loc) · 2.11 KB
/
Copy pathagents.py
File metadata and controls
68 lines (48 loc) · 2.11 KB
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
from typing import AsyncIterable
import aiohttp
from faust import StreamT
from loguru import logger
from horton.alphavantage import AlphaVantageClient
from horton.app import get_app
from horton.config import API_KEY
from horton.database.cruds.security import SecurityCRUD
from horton.records import CollectSecurityOverview
app = get_app()
collect_securities_topic = app.topic("collect_securities", internal=True, partitions=2)
collect_security_overview_topic = app.topic(
"collect_security_overview", internal=True, value_type=CollectSecurityOverview
)
@app.agent(collect_security_overview_topic)
async def collect_security_overview(
stream: StreamT[CollectSecurityOverview],
) -> AsyncIterable[bool]:
async with aiohttp.ClientSession() as session:
async for event in stream:
logger.info(
"Start collect security [{symbol}] overview", symbol=event.symbol
)
client = AlphaVantageClient(session, API_KEY)
security_overview = await client.get_security_overview(event.symbol)
await SecurityCRUD.update_one({"symbol": event.symbol, "exchange": event.exchange}, security_overview)
yield True
@app.agent(collect_securities_topic)
async def collect_securities(stream: StreamT[None]) -> AsyncIterable[bool]:
async with aiohttp.ClientSession() as session:
async for _ in stream:
logger.info("Start collect securities")
client = AlphaVantageClient(session, API_KEY)
securities = await client.get_securities()
for security in securities:
await SecurityCRUD.update_one(
{"symbol": security["symbol"], "exchange": security["exchange"]},
security,
upsert=True,
)
await collect_security_overview.cast(
CollectSecurityOverview(symbol=security["symbol"], exchange=security["exchange"])
)
yield True
@app.command()
async def start_collect_securities():
"""Collect securities and overview."""
await collect_securities.cast()