/
denver2990
/
async_python_final_assignment
Обзор
Документация
Войти
/
denver2990
/
async_python_final_assignment
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
async_moex_collector.py
217 строк
8 KB
Eugene Popov
потоки отменены, вывод статистики по количесву записей
27 ноя 2025, 20:57
27 ноя 2025, 20:57
47ae39c
Код
Авторство
О чём код?
import asyncio import aiohttp from pathlib import Path import csv from datetime import datetime, timedelta import time from concurrent.futures import ProcessPoolExecutor process_pool = ProcessPoolExecutor(max_workers=16) async def ticker_generator(path: str): for line in Path(path).read_text(encoding="utf-8").splitlines(): ticker = line.strip() if ticker: await asyncio.sleep(0) yield ticker async def fetch_json(session: aiohttp.ClientSession, url: str, params: dict | None = None) -> dict: """ Асинхронно выполняет GET запрос и возвращает JSON как dict. """ async with session.get(url, params=params) as response: response.raise_for_status() return await response.json() async def div_request(session: aiohttp.ClientSession, ticker: str): div_url = f"http://iss.moex.com/iss/securities/{ticker}/dividends.json" data = await fetch_json(session, div_url) div_section = data.get("dividends", {}) columns = div_section.get("columns", []) rows = div_section.get("data", []) if not columns: print(f"Не найдены данные о дивидендах для {ticker}") return [] try: idx_date = columns.index("registryclosedate") idx_value = columns.index("value") idx_currency = columns.index("currencyid") except ValueError as e: print(f"Не удалось найти необходимые колонки дивидендов для {ticker}: {e}") return [] result = [] for row in rows: if row[idx_value] is not None: # Пропускаем записи без значения дивиденда result.append({ "date": row[idx_date], "value": float(row[idx_value]), "currency": row[idx_currency] }) return result def save_dividends_to_csv(ticker: str, records: list[dict], filename: str | None = None) -> None: """ Сохраняет дивиденды в CSV-файл. Формат: date,value,currency """ if filename is None: filename = f"{ticker}_dividends.csv" with open(filename, "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow(["date", "value", "currency"]) # заголовок for item in records: writer.writerow([item["date"], item["value"], item["currency"]]) print(f"Для тикера {ticker} сохранено {len(records)} заисеей дивидендов") def save_prices_to_csv(ticker: str, records: list[dict], filename: str | None = None) -> None: """ Сохраняет цены закрытия в CSV-файл. Формат: date,close """ if filename is None: filename = f"{ticker}_prices.csv" with open(filename, "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow(["date", "close"]) for item in records: writer.writerow([item["date"], item["close"]]) print(f"Для тикера {ticker} сохранено {len(records)} заисеей котировок") async def get_date_bounds(session: aiohttp.ClientSession, ticker: str) -> tuple[str, str]: """Получение граничных дат для истории котировок""" url = f"http://iss.moex.com/iss/history/engines/stock/markets/shares/boards/TQBR/securities/{ticker}/dates.json" data = await fetch_json(session, url) dates_section = data.get("dates", {}) columns = dates_section.get("columns", []) rows = dates_section.get("data", []) if not dates_section: raise ValueError(f"Не удалось получить граничные даты для {ticker}") idx_from = columns.index("from") if "from" in columns else columns.index("FROM") idx_till = columns.index("till") if "till" in columns else columns.index("TILL") if not rows[0][idx_from] or not rows[0][idx_till]: raise ValueError(f"Нет торгов по тикеру {ticker}") date_from = datetime.strptime(rows[0][idx_from], "%Y-%m-%d") date_till = datetime.strptime(rows[0][idx_till], "%Y-%m-%d") print(f"Для тикера {ticker} диапазон дат {date_from.strftime("%Y-%m-%d")} - {date_till.strftime("%Y-%m-%d")}") return date_from, date_till async def fetch_prices(session: aiohttp.ClientSession, ticker: str): date_start, date_end = await get_date_bounds(session, ticker) prices = [] while date_start <= date_end: date_to = min(date_start + timedelta(100), date_end) url = f"http://iss.moex.com/iss/history/engines/stock/markets/shares/boards/TQBR/securities/{ticker}.json?from={date_start.strftime("%Y-%m-%d")}&till={date_to.strftime("%Y-%m-%d")}" try: data = await fetch_json(session, url) history_section = data.get("history", {}) columns = history_section.get("columns", []) rows = history_section.get("data", []) if not columns or not rows: print(f"Пустой ответ цен для {ticker} url {url}") break try: idx_date = columns.index("TRADEDATE") idx_close = columns.index("CLOSE") except ValueError as e: print(f"Не удалось найти необходимые колонки в данных цен для {ticker}: {e}") break for row in rows: if row[idx_close] is not None: # Пропускаем дни без торгов prices.append({ "date": row[idx_date], "close": float(row[idx_close]) }) except: print("ошибка запроса {ticker} {e}") date_start = date_start + timedelta(days=100) if not prices: print(f"Нет истории для тикера {ticker}") return prices async def process_div(session: aiohttp.ClientSession, ticker: str): print(f"Дивы для тикера {ticker}") try: dividends = await div_request(session, ticker) if dividends: save_dividends_to_csv(ticker, dividends, f"result/{ticker}_dividends.csv") #самый быстрый вариант #await asyncio.to_thread(save_dividends_to_csv, ticker, dividends, f"result/{ticker}_dividends.csv") #asyncio.get_event_loop().run_in_executor(process_pool, save_dividends_to_csv, ticker, dividends, f"result/{ticker}_dividends.csv") else: print(f"{ticker} не платит дивиденды") except Exception as e: print(f"Ошибка при обработке тикера {ticker}: {e}") async def process_prices(session: aiohttp.ClientSession, ticker: str): print(f"Котировки тикера {ticker}") try: prices = await fetch_prices(session, ticker) if prices: save_prices_to_csv(ticker, prices, f"result/{ticker}_prices.csv") #самый быстрый вариант #await asyncio.to_thread(save_prices_to_csv, ticker, prices, f"result/{ticker}_prices.csv") #asyncio.get_event_loop().run_in_executor(process_pool, save_prices_to_csv, ticker, prices, f"result/{ticker}_prices.csv") except Exception as e: print(f"Ошибка при обработке тикера {ticker}: {e}") async def main(): start_time = time.perf_counter() async with aiohttp.ClientSession() as session: tasks = [] async for ticker in ticker_generator("test_data/tickers.txt"): # запрос дивидендов task_div = asyncio.create_task(process_div(session, ticker)) tasks.append(task_div) # запрос тикеров task_prices = asyncio.create_task(process_prices(session, ticker)) tasks.append(task_prices) await asyncio.gather(*tasks, return_exceptions=True) end_time = time.perf_counter() print(f"Время выполнения {(end_time - start_time) : 3f}, тикеров обработано {len(tasks)/2}") if __name__ == "__main__": asyncio.run(main())