264 lines
12 KiB
Python
264 lines
12 KiB
Python
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')
|
|
PROCESS_EVENT_UPDATES = os.getenv('PROCESS_EVENT_UPDATES', 'false').lower() in ('true', '1', 'yes', 'on')
|
|
AUTO_PUB_EVT = os.getenv('AUTO_PUB_EVT', 'false').lower() in ('true', '1', 'yes', 'on')
|
|
|
|
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 create_event_post_links(cursor, event_id, unit_name):
|
|
"""Создание связей между событием и постами по auto_unit"""
|
|
if not unit_name:
|
|
return
|
|
|
|
try:
|
|
# Поиск постов с соответствующим auto_unit
|
|
cursor.execute(
|
|
"SELECT id FROM posts WHERE auto_unit = %s",
|
|
(unit_name,)
|
|
)
|
|
posts = cursor.fetchall()
|
|
|
|
for post in posts:
|
|
post_id = post['id']
|
|
|
|
# Проверяем, существует ли уже связь
|
|
cursor.execute(
|
|
"SELECT id FROM interlinks WHERE post_id = %s AND event_id = %s",
|
|
(post_id, event_id)
|
|
)
|
|
|
|
if not cursor.fetchone():
|
|
# Создаем новую связь с auto_added = 1
|
|
cursor.execute(
|
|
"INSERT INTO interlinks (post_id, event_id, created_at, auto_added) VALUES (%s, %s, %s, 1)",
|
|
(post_id, event_id, datetime.now())
|
|
)
|
|
log_message(f"Создана автоматическая связь: событие {event_id} <-> пост {post_id} (unit: {unit_name})")
|
|
else:
|
|
log_message(f"Связь уже существует: событие {event_id} <-> пост {post_id}")
|
|
|
|
except Error as e:
|
|
log_message(f"Ошибка при создании связей для события {event_id}: {e}", "ERROR")
|
|
|
|
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)
|
|
|
|
# Определение marked_to_publication для новой записи
|
|
marked_to_publication = 1 if (is_visible and AUTO_PUB_EVT) else 0
|
|
|
|
# Проверка существующей записи
|
|
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, marked_to_publication, None, None, None)
|
|
)
|
|
inserted_count += 1
|
|
log_message(f"Добавление: Запись {id_event}, Created: {added}")
|
|
|
|
# Создание связей с постами для новой записи
|
|
cursor.execute("SELECT id FROM events WHERE id_event = %s AND number = %s", (id_event, number))
|
|
event_record = cursor.fetchone()
|
|
if event_record:
|
|
create_event_post_links(cursor, event_record['id'], unit_name)
|
|
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(seconds=2)
|
|
|
|
if time_diff > timedelta(seconds=1):
|
|
# Проверка изменения только по полю name
|
|
name_changed = existing['name'] != name
|
|
# Сброс флага is_posted_tg только если PROCESS_EVENT_UPDATES = true
|
|
is_posted_tg = 0 if (name_changed and PROCESS_EVENT_UPDATES) else existing['is_posted_tg']
|
|
|
|
# Обновление marked_to_publication на основе новых данных
|
|
new_marked_to_publication = 1 if (is_visible and AUTO_PUB_EVT) else 0
|
|
|
|
# log_message(f"PROCESS_EVENT_UPDATES: {PROCESS_EVENT_UPDATES}")
|
|
# log_message(f"name_changed: {name_changed}")
|
|
|
|
# Обновление записи
|
|
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,
|
|
marked_to_publication = %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,
|
|
new_marked_to_publication, id_event, number)
|
|
)
|
|
updated_count += 1
|
|
|
|
# Создание связей с постами для обновленной записи
|
|
cursor.execute("SELECT id FROM events WHERE id_event = %s AND number = %s", (id_event, number))
|
|
event_record = cursor.fetchone()
|
|
if event_record:
|
|
create_event_post_links(cursor, event_record['id'], unit_name)
|
|
|
|
# Формирование сообщения для лога
|
|
update_msg = f"Обновление: Запись {id_event}, Updated: {modified}"
|
|
if name_changed and PROCESS_EVENT_UPDATES:
|
|
update_msg += ". Флаг публикации сброшен."
|
|
elif name_changed and not PROCESS_EVENT_UPDATES:
|
|
update_msg += ". Имя изменено, но флаг публикации сохранен (PROCESS_EVENT_UPDATES=false)."
|
|
if new_marked_to_publication:
|
|
update_msg += f" marked_to_publication установлен в 1 (is_visible={is_visible}, AUTO_PUB_EVT={AUTO_PUB_EVT})."
|
|
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()
|