[ВОЛК] вводим ANYIO BACKEND для синхронного потока загрузки картинок
ci/woodpecker/push/woodpecker Pipeline was successful
ci/woodpecker/push/woodpecker Pipeline was successful
This commit is contained in:
+18
@@ -1032,10 +1032,26 @@ def api_publish_post(post_id):
|
|||||||
# Функция для запуска в отдельном потоке
|
# Функция для запуска в отдельном потоке
|
||||||
def run_async_func():
|
def run_async_func():
|
||||||
try:
|
try:
|
||||||
|
# Устанавливаем правильную политику event loop для работы с anyio
|
||||||
|
# Это критично для правильной работы httpx в отдельном потоке
|
||||||
|
if sys.platform == 'win32':
|
||||||
|
# На Windows используем ProactorEventLoop
|
||||||
|
policy = asyncio.WindowsProactorEventLoopPolicy()
|
||||||
|
else:
|
||||||
|
# На Linux используем DefaultEventLoopPolicy
|
||||||
|
policy = asyncio.DefaultEventLoopPolicy()
|
||||||
|
|
||||||
|
asyncio.set_event_loop_policy(policy)
|
||||||
|
|
||||||
# Создаем новый event loop для этого потока
|
# Создаем новый event loop для этого потока
|
||||||
# Это важно при работе с gunicorn, который может иметь свой event loop
|
# Это важно при работе с gunicorn, который может иметь свой event loop
|
||||||
loop = asyncio.new_event_loop()
|
loop = asyncio.new_event_loop()
|
||||||
asyncio.set_event_loop(loop)
|
asyncio.set_event_loop(loop)
|
||||||
|
|
||||||
|
# Устанавливаем переменную окружения для anyio backend
|
||||||
|
import os
|
||||||
|
os.environ['ANYIO_BACKEND'] = 'asyncio'
|
||||||
|
|
||||||
try:
|
try:
|
||||||
result = loop.run_until_complete(publish_to_tg(post['vk_post_id']))
|
result = loop.run_until_complete(publish_to_tg(post['vk_post_id']))
|
||||||
return result
|
return result
|
||||||
@@ -1054,6 +1070,8 @@ def api_publish_post(post_id):
|
|||||||
pass
|
pass
|
||||||
finally:
|
finally:
|
||||||
loop.close()
|
loop.close()
|
||||||
|
# Сбрасываем политику event loop
|
||||||
|
asyncio.set_event_loop_policy(None)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
import traceback
|
import traceback
|
||||||
log_event(f"Ошибка в run_async_func: {str(e)}\n{traceback.format_exc()}")
|
log_event(f"Ошибка в run_async_func: {str(e)}\n{traceback.format_exc()}")
|
||||||
|
|||||||
+56
-6
@@ -9,10 +9,22 @@ import re
|
|||||||
import httpx
|
import httpx
|
||||||
import traceback
|
import traceback
|
||||||
import io
|
import io
|
||||||
|
import sys
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from dotenv import load_dotenv
|
from dotenv import load_dotenv
|
||||||
from formatter import get_event_text
|
from formatter import get_event_text
|
||||||
|
|
||||||
|
# Настройка anyio для правильной работы в отдельном потоке
|
||||||
|
# Устанавливаем правильный backend для anyio (используем asyncio)
|
||||||
|
try:
|
||||||
|
import anyio
|
||||||
|
# Устанавливаем backend для anyio явно
|
||||||
|
if not hasattr(anyio, '_backend'):
|
||||||
|
# Используем asyncio backend
|
||||||
|
os.environ.setdefault('ANYIO_BACKEND', 'asyncio')
|
||||||
|
except ImportError:
|
||||||
|
pass
|
||||||
|
|
||||||
# Загрузка переменных окружения
|
# Загрузка переменных окружения
|
||||||
load_dotenv()
|
load_dotenv()
|
||||||
|
|
||||||
@@ -248,6 +260,46 @@ async def publish_to_tg(vk_post_id):
|
|||||||
write_timeout = max(120.0, 60.0 + (file_size_mb * 10.0))
|
write_timeout = max(120.0, 60.0 + (file_size_mb * 10.0))
|
||||||
logger.info(f"Размер файла: {file_size_mb:.2f} МБ, таймаут на запись: {write_timeout:.1f} секунд")
|
logger.info(f"Размер файла: {file_size_mb:.2f} МБ, таймаут на запись: {write_timeout:.1f} секунд")
|
||||||
|
|
||||||
|
# Для больших файлов используем синхронный клиент, чтобы избежать проблем с event loop
|
||||||
|
# Синхронный клиент более стабилен при работе в отдельном потоке с gunicorn
|
||||||
|
use_sync_client = file_size_mb > 0.5 # Для файлов больше 0.5 МБ используем синхронный клиент
|
||||||
|
|
||||||
|
# Логируем время начала запроса
|
||||||
|
request_start_time = time.time()
|
||||||
|
logger.info(f"Начало запроса отправки фото в {datetime.now(timezone.utc).isoformat()}")
|
||||||
|
|
||||||
|
if use_sync_client:
|
||||||
|
logger.info(f"Используем синхронный httpx клиент для файла размером {file_size_mb:.2f} МБ")
|
||||||
|
|
||||||
|
# Используем синхронный httpx клиент для больших файлов
|
||||||
|
import httpx as httpx_sync
|
||||||
|
upload_timeout_config = httpx_sync.Timeout(
|
||||||
|
connect=10.0,
|
||||||
|
read=120.0,
|
||||||
|
write=write_timeout, # Динамический таймаут на запись
|
||||||
|
pool=10.0
|
||||||
|
)
|
||||||
|
|
||||||
|
with httpx_sync.Client(
|
||||||
|
http2=False,
|
||||||
|
timeout=upload_timeout_config,
|
||||||
|
follow_redirects=True
|
||||||
|
) as sync_client:
|
||||||
|
# Отправляем изображение как файл через multipart/form-data
|
||||||
|
image_file = io.BytesIO(image_data)
|
||||||
|
files = {
|
||||||
|
'photo': ('image.jpg', image_file, 'image/jpeg')
|
||||||
|
}
|
||||||
|
data = {
|
||||||
|
'chat_id': CHANNEL_ID,
|
||||||
|
'caption': text,
|
||||||
|
'parse_mode': 'HTML',
|
||||||
|
'disable_notification': str(PUBLISH_SILENTLY).lower()
|
||||||
|
}
|
||||||
|
|
||||||
|
photo_response = sync_client.post(send_photo_url, files=files, data=data)
|
||||||
|
else:
|
||||||
|
# Для маленьких файлов используем асинхронный клиент
|
||||||
# Создаем клиент с увеличенным таймаутом на запись для больших файлов
|
# Создаем клиент с увеличенным таймаутом на запись для больших файлов
|
||||||
upload_timeout_config = httpx.Timeout(
|
upload_timeout_config = httpx.Timeout(
|
||||||
connect=10.0,
|
connect=10.0,
|
||||||
@@ -262,10 +314,6 @@ async def publish_to_tg(vk_post_id):
|
|||||||
)
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# Логируем время начала запроса
|
|
||||||
request_start_time = time.time()
|
|
||||||
logger.info(f"Начало запроса отправки фото в {datetime.now(timezone.utc).isoformat()}")
|
|
||||||
|
|
||||||
# Отправляем изображение как файл через multipart/form-data
|
# Отправляем изображение как файл через multipart/form-data
|
||||||
# Создаем BytesIO объект для файла
|
# Создаем BytesIO объект для файла
|
||||||
image_file = io.BytesIO(image_data)
|
image_file = io.BytesIO(image_data)
|
||||||
@@ -280,9 +328,13 @@ async def publish_to_tg(vk_post_id):
|
|||||||
}
|
}
|
||||||
|
|
||||||
photo_response = await upload_client.post(send_photo_url, files=files, data=data)
|
photo_response = await upload_client.post(send_photo_url, files=files, data=data)
|
||||||
|
finally:
|
||||||
|
await upload_client.aclose()
|
||||||
|
|
||||||
request_duration = time.time() - request_start_time
|
request_duration = time.time() - request_start_time
|
||||||
logger.info(f"Запрос отправки фото завершен за {request_duration:.2f} секунд")
|
logger.info(f"Запрос отправки фото завершен за {request_duration:.2f} секунд")
|
||||||
|
|
||||||
|
# Обработка ответа (общая для синхронного и асинхронного клиентов)
|
||||||
if photo_response.status_code == 200:
|
if photo_response.status_code == 200:
|
||||||
photo_result = photo_response.json()
|
photo_result = photo_response.json()
|
||||||
if photo_result.get('ok'):
|
if photo_result.get('ok'):
|
||||||
@@ -308,8 +360,6 @@ async def publish_to_tg(vk_post_id):
|
|||||||
pass
|
pass
|
||||||
raise RetryAfterException(retry_after)
|
raise RetryAfterException(retry_after)
|
||||||
raise Exception(f"HTTP {photo_response.status_code}: {photo_response.text}")
|
raise Exception(f"HTTP {photo_response.status_code}: {photo_response.text}")
|
||||||
finally:
|
|
||||||
await upload_client.aclose()
|
|
||||||
except RetryAfterException:
|
except RetryAfterException:
|
||||||
raise
|
raise
|
||||||
except (httpx.TimeoutException, httpx.WriteTimeout, httpx.ReadTimeout) as e:
|
except (httpx.TimeoutException, httpx.WriteTimeout, httpx.ReadTimeout) as e:
|
||||||
|
|||||||
Reference in New Issue
Block a user