Upload files to "/"

This commit is contained in:
2025-09-25 20:43:28 +03:00
parent cdb85f61c7
commit ffb2c3fa63
4 changed files with 211 additions and 0 deletions
+199
View File
@@ -0,0 +1,199 @@
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()