Files
Zilant2025/volk_load.py
T
2025-12-21 12:39:06 +03:00

509 lines
24 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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")
log_message(error_msg)
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")
log_message(error_msg)
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")
log_message(error_msg)
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")
log_message(f"Пропущена запись без id или с пустым id: {item}")
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:
# Явно указываем кодировку UTF-8 при декодировании
raw_data = response.read()
try:
decoded_data = raw_data.decode('utf-8')
except UnicodeDecodeError as e:
error_msg = f"Ошибка декодирования ответа для категории {category_id}: {e}. Попытка декодирования как latin-1."
# log_message(error_msg, "WARNING")
log_message(error_msg)
decoded_data = raw_data.decode('latin-1')
try:
desc_response = json.loads(decoded_data)
except json.JSONDecodeError as e:
error_msg = f"Ошибка парсинга JSON для категории {category_id}: {e}. Первые 500 символов ответа: {decoded_data[:500]}"
# log_message(error_msg, "ERROR")
log_message(error_msg)
print(error_msg)
# Пропускаем эту категорию и продолжаем обработку
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")
log_message(error_msg)
print(error_msg)
# Пропускаем эту категорию и продолжаем обработку
continue
# Проверка наличия блока data
if 'data' not in desc_response:
error_msg = f"В JSON ответе для категории {category_id} отсутствует блок 'data'. Пропуск категории."
# log_message(error_msg, "ERROR")
log_message(error_msg)
print(error_msg)
# Пропускаем эту категорию и продолжаем обработку
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")
log_message(error_msg)
print(error_msg)
# Пропускаем эту категорию и продолжаем обработку
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")
log_message(error_msg)
print(error_msg)
import traceback
traceback_str = traceback.format_exc()
# log_message(f"Трассировка ошибки для категории {category_id}:\n{traceback_str}", "ERROR")
log_message(f"Трассировка ошибки для категории {category_id}:\n{traceback_str}")
# Пропускаем эту категорию и продолжаем обработку остальных
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")
log_message(error_msg)
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")
log_message(error_msg)
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")
log_message(error_msg)
print(error_msg)
return
# Проверка наличия category_id
if not category_id:
error_msg = "ID категории не указан"
# log_message(error_msg, "ERROR")
log_message(error_msg)
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")
log_message(error_msg)
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")
log_message(error_msg)
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 для всей таблицы, чтобы продолжить нумерацию
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")
log_message(f"Пропущено событие без id: {item}")
continue
# Проверка обязательного поля title
if 'title' not in item or not item['title']:
# log_message(f"Пропущено событие без title (id: {item.get('id', 'N/A')}): {item}", "WARNING")
log_message(f"Пропущено событие без title (id: {item.get('id', 'N/A')}): {item}")
continue
# Подготовка значений с маппингом полей
# Преобразуем id в число, если это необходимо
try:
id_event = int(item['id']) # Уникальный идентификатор события из JSON
except (ValueError, TypeError):
# log_message(f"Пропущено событие с некорректным id (не число): {item.get('id', 'N/A')}", "WARNING")
log_message(f"Пропущено событие с некорректным id (не число): {item.get('id', 'N/A')}")
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")
log_message(error_msg)
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")
log_message(error_msg)
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")
log_message("В базе данных нет категорий для загрузки событий")
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")
log_message(error_msg)
print(error_msg)
# Продолжаем обработку других категорий
continue
log_message(f"Завершение загрузки всех событий. Обработано категорий: {processed_count}/{len(categories)}")
except Error as e:
error_msg = f"Ошибка базы данных при загрузке всех событий: {e}"
# log_message(error_msg, "ERROR")
log_message(error_msg)
print(error_msg)
except Exception as e:
error_msg = f"Общая ошибка при загрузке всех событий: {e}"
# log_message(error_msg, "ERROR")
log_message(error_msg)
print(error_msg)
if __name__ == "__main__":
load_categories()