384 lines
17 KiB
Python
384 lines
17 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 для данной категории, чтобы продолжить нумерацию
|
|
cursor.execute(
|
|
"SELECT MAX(number) as max_number FROM events WHERE unit_name = %s",
|
|
(category_id,)
|
|
)
|
|
result = cursor.fetchone()
|
|
start_number = (result['max_number'] or 0) + 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:
|
|
# Вставка новой записи
|
|
# Используем 0 для id_event, так как его нет в JSON
|
|
cursor.execute(
|
|
"""INSERT INTO events (id_event, number, unit_name, name, about, about_social_picture, authors)
|
|
VALUES (%s, %s, %s, %s, %s, %s, %s)""",
|
|
(0, 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()
|
|
|
|
if __name__ == "__main__":
|
|
load_categories()
|
|
|