Плохо работает многопоточность
Всех приветствую. Я сделал скрипт, который должен отправлять ежесекундно отправлять данные в телеграмм бот. Использую телеграмм API, отправка через библиотеку httpx и для каждого пользователя я выделял отдельный поток. Я для теста вбил в базу данных 10000 пользователей, запустив свой скрипт, производительность упала очень сильно. Раньше когда я запускал одного пользователя, то сообщения от телеграмм бота приходили ежесекундно, но при 10000 пользователях, приходят раз в 5 минут. Прошу помогите, решить не легкую задачку для меня, может вы можете посоветовать какую либо другую альтернативу для решение проблемы
Код:
from read_sql_config import read_db_config
import config
import asyncio
from threading import *
import httpx
import mysql.connector
# Инициализация Базы данных
db_config = read_db_config()
conn = mysql.connector.connect(**db_config)
cursor = conn.cursor(buffered=True)
Data_base = [] # Список данных с БД
def requests_message(id_telegram, wallet):
while True:
# Получение данных с сервера с помощью API
try:
response = httpx.get('https://apilist.tronscan.org/api/transaction',
params={
'address': wallet,
'limit': 3,
'sort': '-timestamp',
})
print(id_telegram)
# Отбор нужных данных полученных с сервера
in_address = response.json()['data'][0]['ownerAddress']
to_address = response.json()['data'][0]['toAddress']
token = response.json()['data'][0]['tokenInfo']['tokenName'].upper()
# Отправление данных с помощью телеграм API
httpx.post(f'https://api.telegram.org/bot{config.api_token}/sendMessage?',
params={
'chat_id': id_telegram,
'text': (f"Полученные данные с кошелька {wallet}\n{in_address, to_address, token}")
})
except:
pass
class AsyncIterator:
def __init__(self, seq):
self.iter = iter(seq)
def __aiter__(self):
return self
async def __anext__(self):
try:
return next(self.iter)
except StopIteration:
raise StopAsyncIteration
# Загрузка базы данных в список
async def iter_row():
cursor.execute("SELECT * FROM users")
while True:
rows = cursor.fetchmany(300)
if not rows:
break
async for row in AsyncIterator(rows):
Data_base.append(row)
# Запуск на получение и отправку данных
async def brute_force():
async for row in AsyncIterator(Data_base):
Thread(target=requests_message, args=(row[1], row[2])).start()
if __name__ == "__main__":
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
loop.run_until_complete(iter_row())
loop.run_until_complete(brute_force())
loop.close()