99 lines
4.0 KiB
Python
99 lines
4.0 KiB
Python
"""
|
|
===============================================================================
|
|
FILE: services/scud_etl/pipeline.py
|
|
ROLE: Загрузка наилучших срезов СКУД и штата/отсутствий 1С с умным fallback-ом.
|
|
===============================================================================
|
|
"""
|
|
|
|
import os
|
|
import logging
|
|
from typing import Optional, Dict, Any, Tuple
|
|
import pandas as pd
|
|
|
|
from core.connection import get_connection
|
|
from core.database import load_scud_from_db_by_snapshot
|
|
from config import DATA_DIR
|
|
from services.data_loader import load_1c_data_smart, load_staff_data, load_absent_data
|
|
|
|
logger = logging.getLogger("SCUD_PIPELINE")
|
|
|
|
|
|
def load_best_snapshot_for_date(date_str: str, prefer_final_y: bool = False) -> Optional[pd.DataFrame]:
|
|
"""
|
|
Загружает наилучший срез СКУД за дату.
|
|
Если prefer_final_y=True — отдает предпочтение финишному Y (23:59:59).
|
|
"""
|
|
with get_connection() as conn:
|
|
cursor = conn.cursor()
|
|
target_snap_id = None
|
|
|
|
if prefer_final_y:
|
|
cursor.execute("""
|
|
SELECT snapshot_id
|
|
FROM scud_logs
|
|
WHERE log_date = ?
|
|
AND (snapshot_id LIKE 'Y%' OR snapshot_time LIKE '%23:59:59' OR snapshot_time LIKE '%22:00:00')
|
|
ORDER BY id DESC LIMIT 1
|
|
""", (date_str,))
|
|
row = cursor.fetchone()
|
|
if row:
|
|
target_snap_id = row[0]
|
|
|
|
if not target_snap_id:
|
|
cursor.execute("""
|
|
SELECT snapshot_id
|
|
FROM scud_logs
|
|
WHERE log_date = ?
|
|
ORDER BY id DESC LIMIT 1
|
|
""", (date_str,))
|
|
row = cursor.fetchone()
|
|
if row:
|
|
target_snap_id = row[0]
|
|
|
|
if not target_snap_id:
|
|
return None
|
|
|
|
return load_scud_from_db_by_snapshot(date_str, snapshot_param=target_snap_id)
|
|
|
|
|
|
def load_1c_files_for_date(date_str: str) -> Tuple[Optional[pd.DataFrame], Optional[pd.DataFrame]]:
|
|
"""
|
|
Загружает штат и отсутствия с каскадным fallback:
|
|
1. Синхронизирует свежие файлы с сетевой шары.
|
|
2. Загружает штат и отсутствия через load_1c_data_smart(..., use_db=True).
|
|
3. Если штат пуст — берет последний доступный срез штата из zup_staff в SQLite.
|
|
"""
|
|
clean_date = date_str.replace('_', '.')
|
|
|
|
try:
|
|
from services.share_copier import copy_1c_files_from_share
|
|
copy_1c_files_from_share()
|
|
except Exception as e:
|
|
logger.warning(f"[Pipeline] Ошибка копирования с шары: {e}")
|
|
|
|
df_staff, df_abs = load_1c_data_smart(clean_date, use_db=True)
|
|
|
|
# Fallback за штат: ако данашњи штат још увек није доступан, користи се претходни из базе
|
|
if df_staff is None or df_staff.empty:
|
|
with get_connection(row_factory=True) as conn:
|
|
cursor = conn.cursor()
|
|
cursor.execute("""
|
|
SELECT snapshot_date
|
|
FROM zup_staff
|
|
ORDER BY
|
|
SUBSTR(snapshot_date, 7, 4) DESC,
|
|
SUBSTR(snapshot_date, 4, 2) DESC,
|
|
SUBSTR(snapshot_date, 1, 2) DESC,
|
|
id DESC
|
|
LIMIT 1
|
|
""")
|
|
row = cursor.fetchone()
|
|
if row and row[0]:
|
|
fallback_date = row[0]
|
|
df_staff = pd.read_sql_query(
|
|
"SELECT fio as 'ФИО', fio_clean, department as 'Подразделение', position as 'Должность' FROM zup_staff WHERE snapshot_date = ?",
|
|
conn, params=(fallback_date,)
|
|
)
|
|
logger.info(f"[Pipeline] Для даты {clean_date} применен штат за {fallback_date} ({len(df_staff)} чел.)")
|
|
|
|
return df_staff, df_abs |