import os import json import time from datetime import datetime, timedelta, timezone 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() EVT_PREFIX = os.getenv('EVT_PREFIX') 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="JSON_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 parse_datetime(dt_str): """Парсинг datetime с обработкой временных зон""" if not dt_str: return None # Преобразуем строку в datetime с временной зоной dt = datetime.fromisoformat(dt_str.replace('Z', '+00:00')) # Если datetime имеет временную зону, преобразуем в UTC и удаляем информацию о зоне if dt.tzinfo is not None: dt = dt.astimezone(timezone.utc).replace(tzinfo=None) return dt def process_json_data(): inserted_count = 0 updated_count = 0 try: log_message("Начало работы") # Чтение JSON из URL log_message(f"Чтение данных из {EVT_PREFIX}") with urlopen(EVT_PREFIX) as response: data = json.loads(response.read().decode()) log_message(f"Получено {len(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() for item in data: # Проверка обязательных полей if 'id' not in item or 'number' not in item: log_message(f"Пропущена запись без id или number: {item}", "WARNING") continue # Подготовка значений id_event = item['id'] number = item['number'] unit_name = item.get('unit_name', '') # Убрана обработка пробелов и добавление # name = item.get('name', '') about = item.get('about', '') about_social_picture = item.get('about_social_picture', '') # Обработка tags tags = item.get('tags', '') if tags: tags = ' '.join(f"#{tag}" for tag in tags.split() if len(tag) > 1) # Обработка datetime modified = parse_datetime(item.get('modified')) added = parse_datetime(item.get('added')) # Булевы значения is_canceled = bool(item.get('is_canceled', False)) accepted = bool(item.get('accepted', False)) denied = bool(item.get('denied', False)) # Вычисление is_visible is_visible = bool(name and about and not is_canceled and not denied) # Проверка существующей записи cursor.execute( "SELECT id_event, number, modified, name, about, is_posted_tg FROM events WHERE id_event = %s AND number = %s", (id_event, number) ) existing = cursor.fetchone() if not existing: # Вставка новой записи cursor.execute( """INSERT INTO events ( id_event, number, unit_name, name, about, about_social_picture, modified, added, tags, is_canceled, accepted, denied, is_visible, is_posted_tg, marked_to_publication, tg_message_id, tg_posted_date, announcement_link ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)""", (id_event, number, unit_name, name, about, about_social_picture, modified, added, tags, is_canceled, accepted, denied, is_visible, 0, 0, None, None, None) ) inserted_count += 1 log_message(f"Добавление: Запись {id_event}, Created: {added}") else: # Преобразуем дату из БД к наивному datetime для сравнения existing_modified = existing['modified'] if existing_modified and isinstance(existing_modified, datetime): if existing_modified.tzinfo is not None: existing_modified = existing_modified.replace(tzinfo=None) # Проверка расхождения во времени time_diff = abs(modified - existing_modified) if modified and existing_modified else timedelta(minutes=2) if time_diff > timedelta(minutes=1): # Проверка изменения только по полю name name_changed = existing['name'] != name is_posted_tg = 0 if name_changed else existing['is_posted_tg'] # Обновление записи cursor.execute( """UPDATE events SET unit_name = %s, name = %s, about = %s, about_social_picture = %s, modified = %s, added = %s, tags = %s, is_canceled = %s, accepted = %s, denied = %s, is_visible = %s, is_posted_tg = %s WHERE id_event = %s AND number = %s""", (unit_name, name, about, about_social_picture, modified, added, tags, is_canceled, accepted, denied, is_visible, is_posted_tg, id_event, number) ) updated_count += 1 # Формирование сообщения для лога update_msg = f"Обновление: Запись {id_event}, Updated: {modified}" if name_changed: update_msg += ". Флаг публикации сброшен." log_message(update_msg) 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_json_all(): """Экспортируемая функция для загрузки JSON данных в базу""" process_json_data() if __name__ == "__main__": load_json_all()