Введение
Всем привет. Я не являюсь разработчиком или DevOps специалистом. Для личных целей набросал скрипт для периодической публикации контента в несколько групп в ВК. Ранее рассказывал об этом в этих статьях:
Поначалу скрипт отрабатывал без проблем, и необходимости в системе мониторинга я не увидел. Но спустя несколько месяцев стал чаще сталкиваться с критическими ошибками, из-за которых сервис падал и мог простоять от нескольких дней до нескольких недель.
Сначала посмотрел в сторону Elasticsearch, но мои два VPS сервера с конфигурацией 1x2.2 ГГц и 1 Гб RAM каждый заняты БД и несколькими приложениями. Объем логов от сервисов генерируется скромный, и мне хотелось получать ежедневную отчетность, которую я мог бы сам настраивать.
Поэтому решил собрать свой универсальный сервис для мониторинга ошибок, сохранять расширенные логи в БД и получать ежедневный отчет о работе приложений на почту. Возможно решение пригодится кому-то еще.
Архитектура
Функционал сервиса:
Прием логов через API;
Асинхронная обработка и сохранение в PostgreSQL;
Автоматическая отправка отчетов по электронной почте;
Обработка ошибок с уведомлениями.
Общая архитектура системы:

Код диаграммы
--- title: "Система сбора логов и обработки ошибок" --- flowchart TB classDef person fill:#08427b,color:#ffffff,stroke:#052456 classDef system fill:#1168bd,color:#ffffff,stroke:#0b4883 classDef container fill:#438dd5,color:#ffffff,stroke:#2a6ea8 classDef db fill:#438dd5,color:#ffffff,stroke:#2a6ea8 subgraph External["Внешние системы"] App1["Приложение 1<br/>(Генератор событий)"]:::system App2["Приложение 2<br/>(Генератор отчетов)"]:::system App3["Приложение 3<br/>(Генератор событий)"]:::system end subgraph MainService["Сервис сбора логов и обработки ошибок"] APILayer["API<br/>- /error<br/>- /upload_log<br/>- /report/send"]:::container Scheduler["Scheduler<br/>- Обработка файлов<br/>- Пакетная вставка в БД<br/>- Удаление обработанных файлов"]:::container ReportGen["Report Generator<br/>- Агрегация данных по сервисам<br/>- Генерация HTML-отчетов"]:::container MailService["Mail Service<br/>- Отправка отчета<br/>- Отправка уведомлений об ошибках"]:::container end DB[("БД")]:::db User["Email Client<br/>(Пользователь)<br/>- Получение отчетов<br/>- Получение уведомлений об ошибках"]:::person App1 -->|<span style="color:yellow">POST /error</span>| APILayer App2 -->|<span style="color:orange">POST /report/send</span>| APILayer App3 -->|"POST /upload_log"| APILayer APILayer -->|"Передача файлов"| Scheduler APILayer -->|<span style="color:yellow">Передача ошибки</span>| MailService APILayer -->|<span style="color:orange">Формирование отчета</span>| ReportGen ReportGen -->|<span style="color:orange">Запрос данных</span>| DB Scheduler -->|"Сохранение данных"| DB ReportGen -->|<span style="color:orange">Отправка отчета</span>| MailService MailService -->|"Email уведомления"| User
Структура проекта:
processing_service/ │ ├── main.py # Точка входа приложения, FastAPI настройка ├── requirements.txt # Зависимости Python ├── .env # Переменные окружения │ ├── api_contract/ # Контракты API │ └── contract.py # Pydantic модели │ ├── core/ # Основная логика │ ├── handler_file.py # Обработчик файлов логов (шедулер) │ └── utils.py # Формирование отчетов и HTML │ ├── db/ # Работа с базой данных │ ├── db.py # Подключение к PostgreSQL │ └── model.py # ORM модель │ ├── decorators/ # Декораторы │ └── decorators.py # Перехват исключений с отправкой уведомлений │ ├── mail_sender/ # Отправка почты │ └── mail_sender.py # SMTP клиент для отправки email │ ├── config/ # Конфигурационные файлы │ └── name_services.json # Список сервисов для отчетов │ └── logs_to_process/ # Директория для временного хранения лог-файлов ├── app1_2026-08-23.log # Лог-файлы ожидающие обработки
Реализация
1. API (main.py и contract.py)
создал эндпоинты:
POST /error — прием сообщений об ошибках от других сервисов
GET /report/send — формирование и отправка суточного отчета.
POST /upload_log — загрузка файлов логов для обработки.
GET /health — проверка работоспособности сервиса.
В main.py запускаю шедулер который ждет получения файлов с логами. Файлы загружаются методом POST /upload_log.
Добавил фильтрацию по IP-адресами, так как мои сервисы для публикации в VK расположены на двух серверах. Следовательно, только от них я принимаю файлы и запросы на отправку отчетов.
main.py
from fastapi import FastAPI, BackgroundTasks, UploadFile, File, HTTPException, Request from api_contract.contract import error from fastapi.responses import JSONResponse from mail_sender.mail_sender import MailSend from core.utils import HandlerJsonData import asyncio import os import aiofiles from dotenv import load_dotenv from contextlib import asynccontextmanager from core.handler_file import HandlerLogFile # ---------- Настройка обработки лог-файлов ---------- load_dotenv() PROCESSING_DIR = os.getenv("LOG_PROCESSING_DIR") os.makedirs(PROCESSING_DIR, exist_ok=True) MessageSend = MailSend() HandlerData = HandlerJsonData() # ---------- Функция получения реального IP клиента ---------- def GetClientIp(request: Request) -> str: """Пытается определить реальный IP клиента, учитывая прокси-заголовки.""" forwarded = request.headers.get("X-Forwarded-For") if forwarded: return forwarded.split(",")[0].strip()# берём первый IP в цепочке realIp = request.headers.get("X-Real-IP") if realIp: return realIp.strip() client = request.client # fallback на прямой IP соединения return client.host if client else "unknown" @asynccontextmanager async def lifespan(app: FastAPI): # Запускаем синхронный шедулер в фоновом потоке HandlerLogFile.startScheduler() yield app = FastAPI(lifespan=lifespan) # ---------- Глобальная проверка IP (middleware) ---------- ALLOWED_IPS = {"твое IP","твое IP"} @app.middleware("http") async def ip_whitelist_middleware(request: Request, call_next): # Эндпоинт /health доступен всем (для мониторинга) if request.url.path == "/health": return await call_next(request) clientIp = GetClientIp(request) if clientIp not in ALLOWED_IPS: return JSONResponse( status_code=403, content={"detail": "Access denied"} ) return await call_next(request) # ---------- API эндпоинты ---------- @app.post("/error") async def error_sender(input: error, background_tasks: BackgroundTasks): # асинхронная обработка запросов async def task_send(): await asyncio.to_thread(MessageSend.ErrorMessage, input) background_tasks.add_task(task_send) return JSONResponse(content={"message": "success"}, status_code=200) @app.get("/report/send") def report_data(): htmlReport = HandlerData.SendReport() MessageSend.ReportMessage(htmlReport) return JSONResponse(content={"message": "success"}, status_code=200) @app.get("/health") async def health_answer(): return JSONResponse(content={"message": "I am alive"}, status_code=200) @app.post("/upload_log") async def upload_log(file: UploadFile = File(...)): # Валидация расширения if not file.filename.endswith('.log'): raise HTTPException(status_code=400, detail="File not allowed") # Проверка, что файл не пустой content = await file.read(1024) # читаем первые 1KB для валидации if not content: raise HTTPException(status_code=400, detail="Empty file") # Сохраняем файл в директорию для обработки file_path = os.path.join(PROCESSING_DIR, file.filename) await file.seek(0) # возвращаем курсор в начало после чтения async with aiofiles.open(file_path, 'wb') as out_file: content = await file.read() await out_file.write(content) return JSONResponse(content={"message": f"File {file.filename} uploaded successfully"}, status_code=200) if __name__ == "__main__": import uvicorn uvicorn.run(app=app, host="127.0.0.1", port=8000)
contract.py
from pydantic import BaseModel, Field from typing import Dict, Any class error(BaseModel): error_message: str = Field(default=..., min_length=1, max_length=2000, description="Текст ошибки") sender_service: str = Field(default=..., min_length=1, max_length=150, description="Сервис отправитель") class data_report(BaseModel): type_report: str = Field(default=..., min_length=1, max_length=10, description="Тип отчета") data: Dict[str, Any] = Field(default=..., min_length=1, max_length=1000, description="Данные для отчета")
2. Обработка файлов c логами (handler_file.py)
Файлы обрабатываются в порядке создания (FIFO). Реализована пакетная вставка в БД по 30 записей и автоматическое удаление файлов после успешной обработки.
handler_file.py
import os import json import time import threading from typing import List, Dict, Any from db.db import GetSession from db.model import LogMessages from api_contract.contract import error from mail_sender.mail_sender import MailSend from datetime import datetime MessageSend = MailSend() PROCESSING_DIR = os.getenv("LOG_PROCESSING_DIR", "./logs_to_process") class HandlerLogFile: @staticmethod def errorSend(error_text: str, sender_service: str = "Log Handler") -> None: """Отправляет ошибку (синхронно).""" error_msg = error(error_message=error_text, sender_service=sender_service) MessageSend.ErrorMessage(error_msg) # предполагаем, что метод синхронный @staticmethod def insertBatch(batch: List[Dict[str, Any]], maxRetries: int = 2) -> bool: for attempt in range(1, maxRetries + 1): try: with GetSession() as session: objects = [] for item in batch: # 1. Преобразуем поле event -> type_event, если нужно if 'event' in item: item['event'] = item.pop('event') # 2. Преобразуем строку date в объект datetime if 'date' in item and isinstance(item['date'], str): # Заменяем 'Z' на '+00:00' для корректного парсинга UTC date_string = item['date'].replace('Z', '+00:00') item['date'] = datetime.fromisoformat(date_string) # 3. Отбираем только нужные поля allowed_fields = {'service', 'event', 'status', 'status_description', 'message', 'date', 'created_at'} filtered_item = {k: v for k, v in item.items() if k in allowed_fields} objects.append(LogMessages(**filtered_item)) session.add_all(objects) session.commit() return True except Exception as e: if attempt == maxRetries: print(f"Insert failed after {maxRetries} attempts: {e}") return False time.sleep(2 ** attempt) return False @staticmethod def processLogFile(filePath: str) -> None: """Синхронно читает файл, валидирует JSON и вставляет батчами.""" batch = [] batchSize = 30 try: with open(filePath, 'r', encoding='utf-8') as f: for line in f: line = line.strip() if not line: continue try: data = json.loads(line) batch.append(data) except json.JSONDecodeError: continue # пропускаем битые строки if len(batch) >= batchSize: if not HandlerLogFile.insertBatch(batch): print(f"Failed to insert batch from {filePath}, stopping") return # файл не удаляем batch = [] # Оставшиеся записи if batch: if not HandlerLogFile.insertBatch(batch): print(f"Failed to insert final batch from {filePath}") return # Удаляем файл только после полной успешной обработки os.remove(filePath) print(f"Processed and deleted: {filePath}") except Exception as e: HandlerLogFile.errorSend(f"Failed to process file {os.path.basename(filePath)}: {e}") @staticmethod def schedulerLoop() -> None: """Синхронный цикл шедулера (запускается в отдельном потоке).""" while True: try: files = [] if not os.path.exists(PROCESSING_DIR): os.makedirs(PROCESSING_DIR, exist_ok=True) for f in os.listdir(PROCESSING_DIR): filePath = os.path.join(PROCESSING_DIR, f) if os.path.isfile(filePath) and f.endswith('.log'): files.append((filePath, os.path.getctime(filePath))) # Сортируем по времени создания (старые первыми) files.sort(key=lambda x: x[1]) for filePath, _ in files: HandlerLogFile.processLogFile(filePath) except Exception as e: print(f"Scheduler error: {e}") time.sleep(300) # 5 минут @staticmethod def startScheduler() -> None: """Запускает шедулер в фоновом демон-потоке.""" thread = threading.Thread(target=HandlerLogFile.schedulerLoop, daemon=True) thread.start() print("Scheduler thread started")
3. БД (model.py и db.py)
В postgresql создал отдельную базу, схему, таблицу и пользователя с правами на чтение/запись. Добавил индексы для увеличения скорости запросов при сборе ежедневного отчета.
Описание полей таблицы:
id: уникальный идентификатор
service: наименование приложение
event: событие
status: статус события
status_description: описания статуса
message: сообщение
date: дата события
created_at: дата записи в БД
model.py
from sqlalchemy import DateTime, String, Text, Integer, Identity, PrimaryKeyConstraint, Index, text from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column from datetime import datetime from typing import Optional class Base(DeclarativeBase): pass class LogMessages(Base): __tablename__ = 'log_messages' __table_args__ = ( PrimaryKeyConstraint('id', name='log_messages_pkey'), Index('date_index', 'date'), Index('service_index', 'service'), {'schema': 'log'} ) id: Mapped[int] = mapped_column(Integer, Identity(always=True, start=1, increment=1, minvalue=1, maxvalue=2147483647, cycle=False, cache=1), primary_key=True) service: Mapped[Optional[str]] = mapped_column(String(100)) event: Mapped[Optional[str]] = mapped_column(String(100)) status: Mapped[Optional[str]] = mapped_column(String(100)) status_description: Mapped[Optional[str]] = mapped_column(String(100)) message: Mapped[Optional[str]] = mapped_column(Text) date: Mapped[Optional[datetime]] = mapped_column(DateTime(True)) created_at: Mapped[Optional[datetime]] = mapped_column(DateTime(True), server_default=text('now()'))
db.py
from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker from contextlib import contextmanager import os from urllib.parse import quote_plus dbname = os.getenv("DBNAME") user = os.getenv("USER") password = quote_plus(os.getenv("PASSWORD_DB", "")) host = os.getenv("HOST") SYNC_DB_URL = f"postgresql+psycopg2://{user}:{password}@{host}/{dbname}" sync_engine = create_engine(SYNC_DB_URL, echo=False) SyncSessionLocal = sessionmaker(bind=sync_engine) @contextmanager def GetSession(): session = SyncSessionLocal() try: yield session session.commit() except Exception: session.rollback() raise finally: session.close()
4. Отчеты (utils.py и mail_sender.py)
Для получения отчета я использую две функции: одна достает из БД логи по каждому сервису, вторая оборачивает их в HTML. Далее отправляю отчет на почту. Функция вызывается из main.py через метод GET /report/send. На сервере при помощи systemd и timer я создал задачу на вызов метода через curl.
utils.py
from db.db import GetSession from db.model import LogMessages from decorators.decorators import catch_all_exceptions import json import os from sqlalchemy import func, select class HandlerJsonData(): """класс для управления загрузки данных в LocalDB""" def __init__(self): pass @catch_all_exceptions def SendReport(self): with open(os.path.join(os.getcwd(), "config","name_services.json"), 'r', encoding='utf-8') as f: config = json.load(f) services = config.get("NAME_SERVICE") # 2. Используем контекстный менеджер для безопасной работы с БД with GetSession() as session: tables = [] # здесь будем собирать все таблицы для отчёта for service_name in services: # Формируем запрос для конкретного сервиса stmt = select( LogMessages.service, LogMessages.event, LogMessages.status, LogMessages.status_description, func.count().label("count") ).where( LogMessages.service == service_name, # Вчерашний день: данные за 24 часа LogMessages.date >= func.current_date() - 1, LogMessages.date < func.current_date() ).group_by( LogMessages.service, LogMessages.event, LogMessages.status, LogMessages.status_description ).order_by( func.count().desc() ) results = session.execute(stmt).all() # Преобразуем результат в список списков для таблицы data_rows = [] for row in results: data_rows.append([ row.service, row.event, row.status, row.status_description, row.count ]) # Формируем структуру таблицы для этого сервиса table = { "title": service_name, # заголовок таблицы – имя сервиса "headers": ["service", "event", "status", "status_description", "count"], "data": data_rows } tables.append(table) return self.PrepareHtml("Отчеты по работе сервисов за последние 24 часа", tables) def PrepareHtml(self, body: str, tables: list = None): html_body = f"<html><body><p>{body}</p>" # Если есть таблицы if tables: for i, table in enumerate(tables): table_data = table.get('data') table_headers = table.get('headers') table_title = table.get('title', f'Таблица {i+1}') # Добавляем заголовок таблицы html_body += f"<h4 style='margin-bottom: 2px;'>{table_title}</h4>" # Создаем таблицу html_body += "<table border='1' cellpadding='5' style='border-collapse: collapse; margin-bottom: 5px;'>" # Добавляем заголовки столбцов if table_headers: html_body += "<tr>" for header in table_headers: html_body += f"<th style='background-color: #f2f2f2;'>{header}</th>" html_body += "</tr>" # Добавляем данные if table_data: for row in table_data: html_body += "<tr>" for cell in row: html_body += f"<td>{cell}</td>" html_body += "</tr>" html_body += "</table>" html_body += "</body></html>" return html_body
mail_sender.py
import smtplib from email.mime.text import MIMEText from email.header import Header from email.mime.multipart import MIMEMultipart import os from dotenv import load_dotenv from api_contract.contract import error load_dotenv() class MailSend(): def __init__(self): self.server = os.getenv('SERVER') self.port = os.getenv('PORT') self.username = os.getenv('USERNAME_MAIL') self.password = os.getenv('PASSWORD') self.from_address = os.getenv('FROM_ADDRESS') self.to_address = os.getenv('TO_ADDRESS') def ErrorMessage(self, errorMessage: error): # Тема и тело письма subject = f"Ошибка от сервиса {errorMessage.sender_service}" body = f"Сервис {errorMessage.sender_service}.\n\nРасшифровка:\n {errorMessage.error_message}" msg = MIMEText(body.encode('utf-8'), _charset='utf-8') msg['Subject'] = Header(subject, 'utf-8') msg['From'] = self.from_address msg['To'] = self.to_address self.SenderMessage(msg.as_string()) def ReportMessage(self, report: str): # Тема и тело письма subject = f"Отчет за сутки" msg = MIMEMultipart('alternative') msg['Subject'] = Header(subject, 'utf-8') msg['From'] = self.from_address msg['To'] = self.to_address # Создаем HTML-часть с правильным Content-Type html_part = MIMEText(report.encode('utf-8'), 'html', 'utf-8') # Добавляем HTML-часть в письмо msg.attach(html_part) # отправляем по почте self.SenderMessage(msg.as_string()) def SenderMessage(self, msg: str): with smtplib.SMTP_SSL(self.server, self.port) as server: # Авторизация на сервере server.login(self.username, self.password) # Отправка письма server.sendmail(self.from_address, self.to_address, msg) server.quit()
5. Обработка ошибок (decorators.py)
Если во время сбора отчета произойдет ошибка, я решил отправлять себе уведомления. Для этого добавил простой декоратор, чтобы при необходимости использовать его и для другого функционала.
decorators.py
from functools import wraps from mail_sender.mail_sender import MailSend import traceback from api_contract.contract import error # если при обработки произойдет ошибка def catch_all_exceptions(func): @wraps(func) def wrapper(*args, **kwargs): try: return func(*args, **kwargs) except Exception as ex: errorMessage = ( f"Ошибка: {type(ex).__name__}\n" f"Сообщение: {ex}\n" f"Трейсбэк: {''.join(traceback.format_tb(ex.__traceback__))}" ) service_name = "service msgbox" # Создаем объект ошибки согласно модели error_data = error( error_message=errorMessage[:1000], # Обрезаем до максимальной длины sender_service=service_name[:150] # Обрезаем до максимальной длины ) mail = MailSend() mail.ErrorMessage(error_data) raise return wrapper
6. Настройки (name_services.json и .env)
Названия сервисов, которые генерируют логи, храню в файле name_services.json. Чтобы расширить количество приложений достаточно внести изменения в этот файл.
В .env храню настройки для почты, данные для подключения к БД и директорию, куда буду сохранять файлы с логами.
name_services.json
{ "NAME_SERVICE": ["app1", "app2"] }
.env
SERVER = PORT = USERNAME_MAIL = PASSWORD = FROM_ADDRESS = TO_ADDRESS = LOG_PROCESSING_DIR = DBNAME = USER = PASSWORD_DB = HOST = SCHEMA =
Далее создал техническую почту на Mail.ru. Требуется создать пароль для внешних приложений, после чего с техченической почты я отправляю отчеты и ошибки на личную.

7. Развертывание:
Через ssh доставил код на gitlab, с gitlab на сервер
Приложение развернул через docker
Собрал образ docker build -t api_report .
Запустил контейнер docker run --restart=always -d -p 8000:8000 --name=api_report api_report
Dockerfile
FROM python:3.11.4-slim # Создаем рабочую директорию RUN mkdir -p /usr/src/app/ WORKDIR /usr/src/app/ # Копируем только папку src в рабочую директорию COPY app/. /usr/src/app/ # Устанавливаем pip последней версии и необходимые зависимости RUN apt-get update && \ apt-get install -y libpq-dev postgresql-client && \ python -m pip install --upgrade pip && \ pip install --no-cache-dir -r requirements.txt && \ rm -rf /root/.cache/pip # Устанавливаем часовой пояс ENV TZ Europe/Moscow EXPOSE 8000 # Указываем точку входа приложения CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"]
В своих приложениях написал вызов методов
inner_api.py
import requests import os from dotenv import load_dotenv from pathlib import Path import yaml from yaml.loader import SafeLoader load_dotenv() def SendError(nameService: str, descriptionError: str): """отправка ошибки""" param = { "error_message": descriptionError, "sender_service": nameService } EVENT_URL = os.getenv("EVENT_URL") # забираем адрес requests.post(f"{EVENT_URL}/error", json=param) def SendLogs(): """Отправляет все .log файлы из директории LOG_PATH на сервер.""" with open(os.path.join(os.getcwd(),"config","config.yaml"), "r") as stream: configSystem = yaml.load(stream, Loader=SafeLoader) EVENT_URL = os.getenv("EVENT_URL").rstrip('/') logPath = Path(configSystem['LogPath']['LOG_PATH']) # Находим все .log файлы в корне папки logFiles = list(logPath.glob("*.log")) if not logFiles: return for logFile in logFiles: with open(logFile, 'rb') as f: files = {'file': (logFile.name, f, 'application/octet-stream')} response = requests.post(f"{EVENT_URL}/upload_log", files=files) if response.status_code == 200: logFile.unlink()# Удаляем файл только при успешной отправке remaining = list(logPath.glob("*.log"))# После обработки всех файлов проверяем, остались ли ещё .log файлы if not remaining: os.environ["LOG_LATEST_FILE"] = ""
Помощь AI
Код, который выдавала нейросеть, показался мне в каких-то местах излишним, и я старался оптимизировать его без потери основного функционала.
Но самый большой вклад нейросети в скриптах handler_file.py для написания шедулера и подготовки HTML при отправке на почту (функция PrepareHtml utils.py). Редактировать эти скрипты мне было просто лень?.
Значительную экономию времени я получил при тестирование функционала и исправлении багов в коде.
Заключение
Создание собственного сервиса сбора логов и обработки ошибок оказалось интересным экспериментом. На рынке существуют готовые решения, но для малых проектов, как у меня, это напоминает стрельбу из пушки по воробьям.
По итогу получился простой, но рабочий инструмент, который полностью закрывает мои потребности в мониторинге. Самое главное преимущество — возможность дорабатывать его под себя. Думаю, вы так же сможете использовать его и под свои задачи.
Надеюсь, мой опыт вам поможет.