545 lines
28 KiB
Python
545 lines
28 KiB
Python
import os
|
||
import json
|
||
import time
|
||
import traceback
|
||
import socket
|
||
from datetime import datetime
|
||
from urllib.request import urlopen
|
||
from urllib.error import URLError, HTTPError
|
||
from dotenv import load_dotenv
|
||
import mysql.connector
|
||
from mysql.connector import Error
|
||
from ensure_db import ensure_database_structure
|
||
|
||
# Загрузка переменных окружения
|
||
load_dotenv()
|
||
|
||
VOLK_CATEGORY_LIST = os.getenv('VOLK_CATEGORY_LIST')
|
||
VOLK_CATEGORY_DESC_PREFIX = os.getenv('VOLK_CATEGORY_DESC_PREFIX')
|
||
VOLK_EVENT_LIST_PREFIX = os.getenv('VOLK_EVENT_LIST_PREFIX')
|
||
VOLK_API_VERSION = os.getenv('VOLK_API_VERSION')
|
||
MDB_HOST = os.getenv('MDB_HOST')
|
||
MDB_USER = os.getenv('MDB_USER')
|
||
MDB_PW = os.getenv('MDB_PW')
|
||
MDBASE = os.getenv('MDBASE')
|
||
LOG_FILE = os.getenv('LOG_FILE')
|
||
|
||
def log_message(message, category="VOLK_loader"):
|
||
"""Функция для записи в лог-файл с повторными попытками при блокировке"""
|
||
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||
log_entry = f"[{timestamp}] [{category}] {message}\n"
|
||
|
||
# До 5 попыток записи с увеличением задержки
|
||
for attempt in range(5):
|
||
try:
|
||
with open(LOG_FILE, 'a', encoding='utf-8') as f:
|
||
f.write(log_entry)
|
||
return True
|
||
except (IOError, OSError) as e:
|
||
if attempt < 4: # Не последняя попытка
|
||
time.sleep(0.4) # Задержка 400ms между попытками
|
||
else:
|
||
# После 5 неудачных попыток просто возвращаем False
|
||
# Не используем print() чтобы избежать реентерабельных вызовов с gunicorn
|
||
return False
|
||
return False
|
||
|
||
def safe_log_direct(message, category="VOLK_loader"):
|
||
"""Безопасная функция для прямой записи в лог-файл без обработки исключений
|
||
Используется в критических ситуациях, когда обычный log_message может не сработать"""
|
||
try:
|
||
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||
log_entry = f"[{timestamp}] [{category}] {message}\n"
|
||
# Прямая запись без контекстного менеджера для максимальной надежности
|
||
f = open(LOG_FILE, 'a', encoding='utf-8')
|
||
try:
|
||
f.write(log_entry)
|
||
f.flush() # Принудительно сбрасываем буфер
|
||
finally:
|
||
f.close()
|
||
except:
|
||
# Игнорируем все исключения - эта функция должна быть максимально безопасной
|
||
pass
|
||
|
||
def load_categories():
|
||
"""Экспортируемая функция для загрузки категорий из JSON в базу данных"""
|
||
inserted_count = 0
|
||
updated_count = 0
|
||
|
||
try:
|
||
log_message("Начало работы")
|
||
|
||
# Проверка наличия URL
|
||
if not VOLK_CATEGORY_LIST:
|
||
error_msg = "Переменная окружения VOLK_CATEGORY_LIST не установлена"
|
||
log_message(error_msg, "ERROR")
|
||
return
|
||
|
||
# Чтение JSON из URL
|
||
log_message(f"Чтение данных из {VOLK_CATEGORY_LIST}")
|
||
with urlopen(VOLK_CATEGORY_LIST) as response:
|
||
data = json.loads(response.read().decode())
|
||
|
||
# Проверка версии API
|
||
api_version = data.get('api_version')
|
||
if api_version != VOLK_API_VERSION:
|
||
error_msg = f"Несовпадение версии API: получена версия '{api_version}', ожидается '{VOLK_API_VERSION}'. Обработка остановлена."
|
||
log_message(error_msg, "ERROR")
|
||
return
|
||
|
||
log_message(f"Версия API соответствует: {api_version}")
|
||
|
||
# Проверка наличия блока data
|
||
if 'data' not in data:
|
||
error_msg = "В JSON отсутствует блок 'data'"
|
||
log_message(error_msg, "ERROR")
|
||
return
|
||
|
||
categories_data = data['data']
|
||
log_message(f"Получено {len(categories_data)} записей для обработки")
|
||
|
||
# Подключение к базе данных
|
||
connection = mysql.connector.connect(
|
||
host=MDB_HOST,
|
||
user=MDB_USER,
|
||
password=MDB_PW,
|
||
database=MDBASE
|
||
)
|
||
|
||
# Установка временной зоны соединения в UTC
|
||
cursor_temp = connection.cursor()
|
||
cursor_temp.execute("SET time_zone = '+00:00'")
|
||
cursor_temp.close()
|
||
|
||
cursor = connection.cursor(dictionary=True)
|
||
|
||
# Проверка структуры БД
|
||
log_message("Проверка структуры базы данных")
|
||
ensure_database_structure()
|
||
|
||
# Список категорий для последующей загрузки описаний
|
||
loaded_categories = []
|
||
|
||
for item in categories_data:
|
||
# Проверка обязательного поля ID
|
||
if 'id' not in item or not item['id']:
|
||
log_message(f"Пропущена запись без id или с пустым id: {item}", "WARNING")
|
||
continue
|
||
|
||
# Подготовка значений
|
||
category_id = item['id']
|
||
# TITLE может отсутствовать или быть пустым - в этом случае используем пустую строку
|
||
title = item.get('title', '') or ''
|
||
|
||
# Проверка существующей записи
|
||
cursor.execute(
|
||
"SELECT AUTO_ID, ID, TITLE FROM categories WHERE ID = %s",
|
||
(category_id,)
|
||
)
|
||
existing = cursor.fetchone()
|
||
|
||
if not existing:
|
||
# Вставка новой записи
|
||
cursor.execute(
|
||
"INSERT INTO categories (ID, TITLE, description) VALUES (%s, %s, %s)",
|
||
(category_id, title, None)
|
||
)
|
||
inserted_count += 1
|
||
log_message(f"Добавление: Категория {category_id}, Название: {title if title else '(пусто)'}")
|
||
else:
|
||
# Обновление существующей записи (всегда обновляем TITLE, даже если оно пустое)
|
||
if existing['TITLE'] != title:
|
||
cursor.execute(
|
||
"UPDATE categories SET TITLE = %s WHERE ID = %s",
|
||
(title, category_id)
|
||
)
|
||
updated_count += 1
|
||
log_message(f"Обновление: Категория {category_id}, Новое название: {title if title else '(пусто)'}")
|
||
else:
|
||
log_message(f"Категория {category_id} без изменений")
|
||
|
||
connection.commit()
|
||
loaded_categories.append(category_id)
|
||
|
||
log_message(f"Загружено категорий: {len(loaded_categories)}. Добавлено записей: {inserted_count}. Обновлено записей: {updated_count}.")
|
||
|
||
# Загрузка описаний категорий
|
||
if VOLK_CATEGORY_DESC_PREFIX and loaded_categories:
|
||
log_message(f"Начало загрузки описаний для {len(loaded_categories)} категорий")
|
||
desc_updated_count = 0
|
||
|
||
for category_id in loaded_categories:
|
||
try:
|
||
desc_url = VOLK_CATEGORY_DESC_PREFIX + category_id
|
||
log_message(f"Запрос описания для категории {category_id} из {desc_url}")
|
||
log_message("отладочное сообщение для проверки вывода в лог-файл")
|
||
|
||
desc_response = None
|
||
try:
|
||
# Добавляем таймаут для избежания зависаний (30 секунд)
|
||
from urllib.request import Request
|
||
req = Request(desc_url)
|
||
response = urlopen(req, timeout=30)
|
||
log_message("urlopen успешно выполнен")
|
||
except (URLError, HTTPError, OSError, socket.gaierror, socket.timeout, ConnectionError) as e:
|
||
# Используем безопасное логирование для критических ошибок
|
||
safe_log_direct("ПЕРЕХВАЧЕНА СЕТЕВАЯ ОШИБКА urlopen", "ERROR")
|
||
safe_log_direct(f"Категория: {category_id}, URL: {desc_url}", "ERROR")
|
||
safe_log_direct(f"Тип ошибки: {type(e).__name__}, Сообщение: {str(e)}", "ERROR")
|
||
try:
|
||
traceback_str = traceback.format_exc()
|
||
safe_log_direct(f"Трассировка:\n{traceback_str}", "ERROR")
|
||
except:
|
||
pass
|
||
# Также пробуем обычное логирование
|
||
try:
|
||
error_msg = f"Ошибка сетевого подключения при запросе описания категории {category_id} из {desc_url}: {type(e).__name__}: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
except:
|
||
pass
|
||
# Пропускаем эту категорию и продолжаем обработку
|
||
continue
|
||
except Exception as e:
|
||
# Перехватываем любые другие исключения, которые могут возникнуть при urlopen
|
||
safe_log_direct("ПЕРЕХВАЧЕНА НЕОЖИДАННАЯ ОШИБКА urlopen", "ERROR")
|
||
safe_log_direct(f"Категория: {category_id}, URL: {desc_url}", "ERROR")
|
||
safe_log_direct(f"Тип ошибки: {type(e).__name__}, Сообщение: {str(e)}", "ERROR")
|
||
try:
|
||
traceback_str = traceback.format_exc()
|
||
safe_log_direct(f"Трассировка:\n{traceback_str}", "ERROR")
|
||
except:
|
||
pass
|
||
# Также пробуем обычное логирование
|
||
try:
|
||
error_msg = f"Неожиданная ошибка при запросе описания категории {category_id} из {desc_url}: {type(e).__name__}: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
except:
|
||
pass
|
||
# Пропускаем эту категорию и продолжаем обработку
|
||
continue
|
||
|
||
try:
|
||
with response:
|
||
# Явно указываем кодировку UTF-8 при декодировании
|
||
raw_data = response.read()
|
||
log_message("raw data read")
|
||
try:
|
||
decoded_data = raw_data.decode('utf-8')
|
||
log_message("data decoded")
|
||
except UnicodeDecodeError as e:
|
||
error_msg = f"Ошибка декодирования ответа для категории {category_id}: {e}. Попытка декодирования как latin-1."
|
||
log_message(error_msg, "WARNING")
|
||
decoded_data = raw_data.decode('latin-1')
|
||
|
||
try:
|
||
desc_response = json.loads(decoded_data)
|
||
log_message("json loaded")
|
||
except json.JSONDecodeError as e:
|
||
error_msg = f"Ошибка парсинга JSON для категории {category_id}: {e}. Первые 500 символов ответа: {decoded_data[:500]}"
|
||
log_message(error_msg, "ERROR")
|
||
# Пропускаем эту категорию и продолжаем обработку
|
||
continue
|
||
except Exception as e:
|
||
error_msg = f"Ошибка при чтении ответа для категории {category_id}: {type(e).__name__}: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
traceback_str = traceback.format_exc()
|
||
log_message(f"Трассировка ошибки чтения для категории {category_id}:\n{traceback_str}", "ERROR")
|
||
# Пропускаем эту категорию и продолжаем обработку
|
||
continue
|
||
|
||
# Проверяем, что desc_response был успешно получен
|
||
if desc_response is None:
|
||
error_msg = f"Не удалось получить ответ для категории {category_id}"
|
||
log_message(error_msg, "ERROR")
|
||
continue
|
||
|
||
# Проверка версии API
|
||
api_version = desc_response.get('api_version')
|
||
if api_version != VOLK_API_VERSION:
|
||
error_msg = f"Несовпадение версии API при загрузке описания категории {category_id}: получена версия '{api_version}', ожидается '{VOLK_API_VERSION}'. Пропуск категории."
|
||
log_message(error_msg, "ERROR")
|
||
# Пропускаем эту категорию и продолжаем обработку
|
||
continue
|
||
|
||
# Проверка наличия блока data
|
||
if 'data' not in desc_response:
|
||
error_msg = f"В JSON ответе для категории {category_id} отсутствует блок 'data'. Пропуск категории."
|
||
log_message(error_msg, "ERROR")
|
||
# Пропускаем эту категорию и продолжаем обработку
|
||
continue
|
||
|
||
desc_data = desc_response['data']
|
||
|
||
# Проверка типа desc_data - должен быть словарём
|
||
if not isinstance(desc_data, dict):
|
||
error_msg = f"Блок 'data' для категории {category_id} не является словарём (тип: {type(desc_data)}). Пропуск категории."
|
||
log_message(error_msg, "ERROR")
|
||
# Пропускаем эту категорию и продолжаем обработку
|
||
continue
|
||
|
||
# Получаем description из блока data
|
||
description = desc_data.get('description')
|
||
# Если description пустое, None или отсутствует, обнуляем поле в БД
|
||
if not description:
|
||
description = None
|
||
|
||
# Обновляем описание в базе данных
|
||
cursor.execute(
|
||
"UPDATE categories SET description = %s WHERE ID = %s",
|
||
(description, category_id)
|
||
)
|
||
connection.commit()
|
||
desc_updated_count += 1
|
||
|
||
if description:
|
||
log_message(f"Обновлено описание для категории {category_id}")
|
||
else:
|
||
log_message(f"Описание для категории {category_id} обнулено (пустое или отсутствует)")
|
||
|
||
except Exception as e:
|
||
error_msg = f"Ошибка при загрузке описания для категории {category_id}: {type(e).__name__}: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
traceback_str = traceback.format_exc()
|
||
log_message(f"Трассировка ошибки для категории {category_id}:\n{traceback_str}", "ERROR")
|
||
# Пропускаем эту категорию и продолжаем обработку остальных
|
||
continue
|
||
|
||
log_message(f"Завершение загрузки описаний. Обновлено описаний: {desc_updated_count}.")
|
||
|
||
log_message(f"Завершение работы. Добавлено записей: {inserted_count}. Обновлено записей: {updated_count}.")
|
||
|
||
except Error as e:
|
||
error_msg = f"Ошибка базы данных: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
if 'connection' in locals() and connection.is_connected():
|
||
connection.rollback()
|
||
except Exception as e:
|
||
error_msg = f"Общая ошибка: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
finally:
|
||
if 'connection' in locals() and connection.is_connected():
|
||
cursor.close()
|
||
connection.close()
|
||
|
||
def load_events_by_category(category_id):
|
||
"""Экспортируемая функция для загрузки событий категории из JSON в базу данных"""
|
||
inserted_count = 0
|
||
updated_count = 0
|
||
|
||
try:
|
||
log_message(f"Начало загрузки событий для категории {category_id}")
|
||
|
||
# Проверка наличия URL префикса
|
||
if not VOLK_EVENT_LIST_PREFIX:
|
||
error_msg = "Переменная окружения VOLK_EVENT_LIST_PREFIX не установлена"
|
||
log_message(error_msg, "ERROR")
|
||
return
|
||
|
||
# Проверка наличия category_id
|
||
if not category_id:
|
||
error_msg = "ID категории не указан"
|
||
log_message(error_msg, "ERROR")
|
||
return
|
||
|
||
# Формирование URL для запроса
|
||
event_url = VOLK_EVENT_LIST_PREFIX + category_id
|
||
log_message(f"Чтение данных из {event_url}")
|
||
|
||
# Чтение JSON из URL
|
||
with urlopen(event_url) as response:
|
||
data = json.loads(response.read().decode())
|
||
|
||
# Проверка версии API
|
||
api_version = data.get('api_version')
|
||
if api_version != VOLK_API_VERSION:
|
||
error_msg = f"Несовпадение версии API: получена версия '{api_version}', ожидается '{VOLK_API_VERSION}'. Обработка остановлена."
|
||
log_message(error_msg, "ERROR")
|
||
return
|
||
|
||
log_message(f"Версия API соответствует: {api_version}")
|
||
|
||
# Проверка наличия блока data
|
||
if 'data' not in data:
|
||
error_msg = "В JSON отсутствует блок 'data'"
|
||
log_message(error_msg, "ERROR")
|
||
return
|
||
|
||
events_data = data['data']
|
||
log_message(f"Получено {len(events_data)} событий для обработки")
|
||
|
||
# Подключение к базе данных
|
||
connection = mysql.connector.connect(
|
||
host=MDB_HOST,
|
||
user=MDB_USER,
|
||
password=MDB_PW,
|
||
database=MDBASE
|
||
)
|
||
|
||
# Установка временной зоны соединения в UTC
|
||
cursor_temp = connection.cursor()
|
||
cursor_temp.execute("SET time_zone = '+00:00'")
|
||
cursor_temp.close()
|
||
|
||
cursor = connection.cursor(dictionary=True)
|
||
|
||
# Проверка структуры БД
|
||
log_message("Проверка структуры базы данных")
|
||
ensure_database_structure()
|
||
|
||
# Получаем максимальный number для всей таблицы, чтобы продолжить нумерацию
|
||
cursor.execute(
|
||
"SELECT MAX(number) as max_number FROM events"
|
||
)
|
||
result = cursor.fetchone()
|
||
max_number = result['max_number'] or 0
|
||
start_number = max_number + 1
|
||
|
||
# Обработка событий
|
||
current_number = start_number
|
||
for item in events_data:
|
||
# Проверка обязательного поля id
|
||
if 'id' not in item or item['id'] is None:
|
||
log_message(f"Пропущено событие без id: {item}", "WARNING")
|
||
continue
|
||
|
||
# Проверка обязательного поля title
|
||
if 'title' not in item or not item['title']:
|
||
log_message(f"Пропущено событие без title (id: {item.get('id', 'N/A')}): {item}", "WARNING")
|
||
continue
|
||
|
||
# Подготовка значений с маппингом полей
|
||
# Преобразуем id в число, если это необходимо
|
||
try:
|
||
id_event = int(item['id']) # Уникальный идентификатор события из JSON
|
||
except (ValueError, TypeError):
|
||
log_message(f"Пропущено событие с некорректным id (не число): {item.get('id', 'N/A')}", "WARNING")
|
||
continue
|
||
name = item.get('title', '')
|
||
about = item.get('announce', '') or ''
|
||
authors = item.get('authors', '') or ''
|
||
about_social_picture = item.get('image', '') or ''
|
||
unit_name = category_id
|
||
|
||
# Проверка существующей записи по id_event (уникальный идентификатор)
|
||
cursor.execute(
|
||
"SELECT id, id_event, name, authors, about, about_social_picture FROM events WHERE id_event = %s",
|
||
(id_event,)
|
||
)
|
||
existing = cursor.fetchone()
|
||
|
||
if not existing:
|
||
# Вставка новой записи
|
||
# Используем id из JSON для id_event, а для number используем автоинкремент
|
||
cursor.execute(
|
||
"""INSERT INTO events (id_event, number, unit_name, name, about, about_social_picture, authors)
|
||
VALUES (%s, %s, %s, %s, %s, %s, %s)""",
|
||
(id_event, current_number, unit_name, name, about, about_social_picture, authors)
|
||
)
|
||
inserted_count += 1
|
||
log_message(f"Добавление: Событие '{name}', id_event: {id_event}, номер: {current_number}, категория: {category_id}")
|
||
current_number += 1
|
||
else:
|
||
# Обновление существующей записи
|
||
needs_update = False
|
||
update_fields = []
|
||
|
||
if existing['name'] != name:
|
||
needs_update = True
|
||
update_fields.append('name')
|
||
if existing['authors'] != authors:
|
||
needs_update = True
|
||
update_fields.append('authors')
|
||
if existing['about'] != about:
|
||
needs_update = True
|
||
update_fields.append('about')
|
||
if existing['about_social_picture'] != about_social_picture:
|
||
needs_update = True
|
||
update_fields.append('about_social_picture')
|
||
|
||
if needs_update:
|
||
cursor.execute(
|
||
"""UPDATE events SET name = %s, authors = %s, about = %s, about_social_picture = %s
|
||
WHERE id_event = %s""",
|
||
(name, authors, about, about_social_picture, id_event)
|
||
)
|
||
updated_count += 1
|
||
log_message(f"Обновление: Событие id_event={id_event} '{name}', обновлены поля: {', '.join(update_fields)}")
|
||
else:
|
||
log_message(f"Событие id_event={id_event} '{name}' без изменений")
|
||
|
||
connection.commit()
|
||
|
||
log_message(f"Завершение работы. Добавлено событий: {inserted_count}. Обновлено событий: {updated_count}.")
|
||
|
||
except Error as e:
|
||
error_msg = f"Ошибка базы данных: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
if 'connection' in locals() and connection.is_connected():
|
||
connection.rollback()
|
||
except Exception as e:
|
||
error_msg = f"Общая ошибка: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
finally:
|
||
if 'connection' in locals() and connection.is_connected():
|
||
cursor.close()
|
||
connection.close()
|
||
|
||
def load_all_events():
|
||
"""Экспортируемая функция для загрузки всех событий по всем категориям"""
|
||
try:
|
||
log_message("Начало загрузки всех событий по всем категориям")
|
||
|
||
# Подключение к базе данных
|
||
connection = mysql.connector.connect(
|
||
host=MDB_HOST,
|
||
user=MDB_USER,
|
||
password=MDB_PW,
|
||
database=MDBASE
|
||
)
|
||
|
||
cursor = connection.cursor(dictionary=True)
|
||
|
||
# Получаем список всех категорий
|
||
cursor.execute("SELECT ID FROM categories ORDER BY ID")
|
||
categories = cursor.fetchall()
|
||
|
||
if not categories:
|
||
log_message("В базе данных нет категорий для загрузки событий", "WARNING")
|
||
cursor.close()
|
||
connection.close()
|
||
return
|
||
|
||
log_message(f"Найдено {len(categories)} категорий для обработки")
|
||
|
||
# Закрываем соединение перед вызовами load_events_by_category
|
||
# (так как она создает свое соединение)
|
||
cursor.close()
|
||
connection.close()
|
||
|
||
# Загружаем события для каждой категории
|
||
processed_count = 0
|
||
for category in categories:
|
||
category_id = category['ID']
|
||
try:
|
||
log_message(f"Обработка категории {category_id} ({processed_count + 1}/{len(categories)})")
|
||
load_events_by_category(category_id)
|
||
processed_count += 1
|
||
except Exception as e:
|
||
error_msg = f"Ошибка при загрузке событий для категории {category_id}: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
# Продолжаем обработку других категорий
|
||
continue
|
||
|
||
log_message(f"Завершение загрузки всех событий. Обработано категорий: {processed_count}/{len(categories)}")
|
||
|
||
except Error as e:
|
||
error_msg = f"Ошибка базы данных при загрузке всех событий: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
except Exception as e:
|
||
error_msg = f"Общая ошибка при загрузке всех событий: {e}"
|
||
log_message(error_msg, "ERROR")
|
||
|
||
if __name__ == "__main__":
|
||
load_categories()
|
||
|