Files
Zilant2025/volk_load.py
T

445 lines
20 KiB
Python

import os
import json
import time
from datetime import datetime
from urllib.request import urlopen
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 неудачных попыток выводим в stdout
print(f"Не удалось записать в лог после 5 попыток: {log_entry.strip()}")
return False
return False
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")
print(error_msg)
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")
print(error_msg)
return
log_message(f"Версия API соответствует: {api_version}")
# Проверка наличия блока data
if 'data' not in data:
error_msg = "В JSON отсутствует блок 'data'"
log_message(error_msg, "ERROR")
print(error_msg)
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}")
with urlopen(desc_url) as response:
desc_response = json.loads(response.read().decode())
# Проверка версии 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")
print(error_msg)
# Останавливаем обработку
break
# Проверка наличия блока data
if 'data' not in desc_response:
error_msg = f"В JSON ответе для категории {category_id} отсутствует блок 'data'. Обработка остановлена."
log_message(error_msg, "ERROR")
print(error_msg)
# Останавливаем обработку
break
desc_data = desc_response['data']
# Получаем 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}: {e}"
log_message(error_msg, "ERROR")
print(error_msg)
# Останавливаем обработку при ошибке
break
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")
print(error_msg)
if 'connection' in locals() and connection.is_connected():
connection.rollback()
except Exception as e:
error_msg = f"Общая ошибка: {e}"
log_message(error_msg, "ERROR")
print(error_msg)
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")
print(error_msg)
return
# Проверка наличия category_id
if not category_id:
error_msg = "ID категории не указан"
log_message(error_msg, "ERROR")
print(error_msg)
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")
print(error_msg)
return
log_message(f"Версия API соответствует: {api_version}")
# Проверка наличия блока data
if 'data' not in data:
error_msg = "В JSON отсутствует блок 'data'"
log_message(error_msg, "ERROR")
print(error_msg)
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 и id_event для данной категории, чтобы продолжить нумерацию
cursor.execute(
"SELECT MAX(number) as max_number, MAX(id_event) as max_id_event FROM events WHERE unit_name = %s",
(category_id,)
)
result = cursor.fetchone()
max_number = result['max_number'] or 0
max_id_event = result['max_id_event'] or 0
# Используем максимальное значение из обоих полей для начала нумерации
start_number = max(max_number, max_id_event) + 1
# Обработка событий
current_number = start_number
for item in events_data:
# Проверка обязательного поля title
if 'title' not in item or not item['title']:
log_message(f"Пропущено событие без title: {item}", "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
# Проверка существующей записи по name (title)
cursor.execute(
"SELECT id, name, authors, about, about_social_picture FROM events WHERE name = %s",
(name,)
)
existing = cursor.fetchone()
if not existing:
# Вставка новой записи
# Используем одинаковое значение для 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)""",
(current_number, current_number, unit_name, name, about, about_social_picture, authors)
)
inserted_count += 1
log_message(f"Добавление: Событие '{name}', номер: {current_number}, категория: {category_id}")
current_number += 1
else:
# Обновление существующей записи
needs_update = False
update_fields = []
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 authors = %s, about = %s, about_social_picture = %s
WHERE name = %s""",
(authors, about, about_social_picture, name)
)
updated_count += 1
log_message(f"Обновление: Событие '{name}', обновлены поля: {', '.join(update_fields)}")
else:
log_message(f"Событие '{name}' без изменений")
connection.commit()
log_message(f"Завершение работы. Добавлено событий: {inserted_count}. Обновлено событий: {updated_count}.")
except Error as e:
error_msg = f"Ошибка базы данных: {e}"
log_message(error_msg, "ERROR")
print(error_msg)
if 'connection' in locals() and connection.is_connected():
connection.rollback()
except Exception as e:
error_msg = f"Общая ошибка: {e}"
log_message(error_msg, "ERROR")
print(error_msg)
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")
print(error_msg)
# Продолжаем обработку других категорий
continue
log_message(f"Завершение загрузки всех событий. Обработано категорий: {processed_count}/{len(categories)}")
except Error as e:
error_msg = f"Ошибка базы данных при загрузке всех событий: {e}"
log_message(error_msg, "ERROR")
print(error_msg)
except Exception as e:
error_msg = f"Общая ошибка при загрузке всех событий: {e}"
log_message(error_msg, "ERROR")
print(error_msg)
if __name__ == "__main__":
load_categories()