diff --git a/.gitignore b/.gitignore index ffae178..d91f4ec 100644 --- a/.gitignore +++ b/.gitignore @@ -69,11 +69,13 @@ api_tokens.json *.csv # IDE and OS +.claude/ .vscode/ .idea/ *.swp *.swo *~ +nul .DS_Store Thumbs.db desktop.ini diff --git a/Configure.json b/Configure.json index 3cbb25a..63ce54e 100644 --- a/Configure.json +++ b/Configure.json @@ -33,7 +33,9 @@ "enabled": true, "weekly_directories": [ "RJ/2.6G/2.6RJGD", - "RJ/2.6G/2.6RJYD" + "RJ/2.6G/2.6RJYD", + "RJ/700M/700RJGD", + "RJ/700M/700RJYD" ], "table_field_mappings": { "2_6GRJGD": [ @@ -51,6 +53,23 @@ {"Source": "cellId", "Target": "cellId", "Type": "string"}, {"Source": "gNBplmn", "Target": "gNBplmn", "Type": "string"}, {"Source": "5G上下行总流量(上行PDCP PDU数据量+下行PDCP成功发送数据量)-YD(GB)", "Target": "上下行总流量_GB", "Type": "float"} + ], + "700MRJGD": [ + {"Source": "开始时间", "Target": "开始时间", "Type": "datetime"}, + {"Source": "结束时间", "Target": "结束时间", "Type": "datetime"}, + {"Source": "gNBId", "Target": "gNBId", "Type": "string"}, + {"Source": "cellId", "Target": "cellId", "Type": "string"}, + {"Source": "gNBplmn", "Target": "gNBplmn", "Type": "string"}, + {"Source": "5G上下行总流量(上行PDCP PDU数据量+下行PDCP成功发送数据量)-GD(GB)", "Target": "上下行总流量_GB", "Type": "float"} + ], + "700MRJYD": [ + {"Source": "开始时间", "Target": "开始时间", "Type": "datetime"}, + {"Source": "结束时间", "Target": "结束时间", "Type": "datetime"}, + {"Source": "gNBId", "Target": "gNBId", "Type": "string"}, + {"Source": "cellId", "Target": "cellId", "Type": "string"}, + {"Source": "gNBplmn", "Target": "gNBplmn", "Type": "string"}, + {"Source": "5G上下行总流量(上行PDCP PDU数据量+下行PDCP成功发送数据量)-YD(GB)", "Target": "上下行总流量_GB", "Type": "float"}, + {"Source": "5G上下行总流量(上行PDCP PDU数据量+下行PDCP成功发送数据量)(GB)-YD(GB)", "Target": "上下行总流量_GB", "Type": "float"} ] } }, @@ -314,4 +333,4 @@ "Type": "float" } ] -} \ No newline at end of file +} diff --git a/app/api/routers/config.py b/app/api/routers/config.py index fb29995..a61a099 100644 --- a/app/api/routers/config.py +++ b/app/api/routers/config.py @@ -7,7 +7,7 @@ from fastapi import APIRouter, Body, HTTPException, UploadFile, File from fastapi.responses import Response from app import state -from app.config import HistoryRetentionConfig, RemoteDataConfig +from app.config import HistoryRetentionConfig, RJDataConfig, RemoteDataConfig from app.services.api_tokens import export_tokens, import_tokens @@ -132,3 +132,7 @@ def _apply_config_data(data: dict[str, Any]) -> None: history_retention = data.get("HistoryRetention") if isinstance(history_retention, dict): state.config.history_retention = HistoryRetentionConfig.from_dict(history_retention) + + rj_data = data.get("RJData") + if isinstance(rj_data, dict): + state.config.rj_data = RJDataConfig.from_dict(rj_data) diff --git a/app/processor.py b/app/processor.py index 5e84ddc..a6a05c1 100644 --- a/app/processor.py +++ b/app/processor.py @@ -525,9 +525,14 @@ class DataProcessor: # 快速字段匹配(使用预编译的映射表) col_mapping = {} + mapped_targets = set() for col in df.columns: if col in field_map: - col_mapping[col] = field_map[col] + target_col = field_map[col] + if target_col in mapped_targets: + continue + col_mapping[col] = target_col + mapped_targets.add(target_col) if len(col_mapping) < 1: if 'kpis' in str(csv_file).lower(): diff --git a/app/services/auto_scheduler.py b/app/services/auto_scheduler.py index db94699..4e2c2e8 100644 --- a/app/services/auto_scheduler.py +++ b/app/services/auto_scheduler.py @@ -1,7 +1,6 @@ from __future__ import annotations import json -import re import threading from collections import defaultdict from dataclasses import dataclass @@ -11,9 +10,9 @@ from typing import Any from fastapi import HTTPException from app import state -from app.config import CACHE_DIR, AutoSchedulerConfig, RemoteDataConfig +from app.config import CACHE_DIR, AutoSchedulerConfig from app.services.remote_download import RemoteDataDownloader, RemoteFileInfo -from app.utils.file_dates import extract_file_date, required_week_days +from app.utils.file_dates import parse_file_date_range, required_week_days READY_DIR = CACHE_DIR / "auto_scheduler" @@ -22,29 +21,36 @@ DISABLED_CHECK_SECONDS = 60 STARTUP_CHECK_SECONDS = 5 FAILURE_RESULTS = {"scan_failed", "trigger_failed", "source_cleanup_failed", "failed", "invalid_flag"} -# RJ周文件名模式: CapacityReportData2.6RJ_GD_YYYYMMDDHHMM_YYYYMMDDHHMM.zip -RJ_WEEKLY_FILE_PATTERN = re.compile( - r"CapacityReportData2\.6RJ_(?:GD|YD)_(\d{12})_(\d{12})\.zip$" -) - @dataclass(frozen=True) -class RJWeeklyDirectoryStatus: - """RJ周数据目录就绪状态""" +class RJDirectoryReadyStatus: + """RJ 数据目录就绪状态,按最新文件自动识别日粒度或周粒度。""" directory: str ready: bool - found: bool = False + granularity: str | None = None + found_days: list[date] | None = None + missing_days: list[date] | None = None file_name: str | None = None file_count: int = 0 error: str | None = None + skipped: bool = False + skip_reason: str | None = None def to_dict(self) -> dict[str, Any]: + found_days = self.found_days or [] + missing_days = self.missing_days or [] return { "ready": self.ready, - "found": self.found, + "granularity": self.granularity, + "found_days": [item.isoformat() for item in found_days], + "missing_days": [item.isoformat() for item in missing_days], + "found_count": len(found_days), + "required_count": len(found_days) + len(missing_days), "file_name": self.file_name, "file_count": self.file_count, "error": self.error, + "skipped": self.skipped, + "skip_reason": self.skip_reason, } @@ -121,6 +127,7 @@ class AutoScheduler: "week_offset": scheduler.week_offset, "auto_delete_source": config.auto_delete_source, "rj_data_enabled": rj_config.enabled, + "rj_directories": rj_config.weekly_directories, "rj_weekly_directories": rj_config.weekly_directories, "next_check_at": self._format_dt(self._next_check_at), "last_check_at": self._format_dt(self._last_check_at), @@ -178,14 +185,22 @@ class AutoScheduler: target_dates = self._target_dates_from_flag(ready_flag, scheduler) return self._trigger_processing(target_dates, ready_flag, manual) - # 检查现有7天目录(4G/5G数据) target_days = required_week_days(scheduler.week_offset) downloader = RemoteDataDownloader(remote_config) - directory_status = self._check_remote_ready(downloader, scheduler, target_days) - self._set_directory_status(directory_status) + rj_config = app_config.rj_data.normalized() + rj_directories = set(rj_config.weekly_directories) if rj_config.enabled else set() + + # 检查现有7天目录(4G/5G数据),RJ 目录单独按粒度判断,避免被普通7天规则误拦。 + directory_status = self._check_remote_ready( + downloader, + scheduler, + target_days, + excluded_directories=rj_directories, + ) error_count = sum(1 for item in directory_status.values() if item.error) if error_count: + self._set_combined_directory_status(directory_status, {}) return self._finish_check( "scan_failed", f"远程目录扫描失败,{error_count}/{len(directory_status)} 个目录无法访问或扫描失败", @@ -196,29 +211,32 @@ class AutoScheduler: skipped_count = len(directory_status) - len(active_status) daily_ready = bool(active_status) and all(item.ready for item in active_status) - # 检查RJ周数据目录(复用同一个 downloader) - rj_weekly_status = self._check_rj_weekly_ready(downloader, scheduler.week_offset) - rj_weekly_ready = True # 默认就绪(如果没有配置RJ目录) + # 检查 RJ 数据目录,按最新文件自动判断日粒度或周粒度。 + rj_status = self._check_rj_ready(downloader, target_days) + rj_error_count = sum(1 for item in rj_status.values() if item.error) + self._set_combined_directory_status(directory_status, rj_status) + if rj_error_count: + return self._finish_check( + "scan_failed", + f"RJ 远程目录扫描失败,{rj_error_count}/{len(rj_status)} 个目录无法访问或扫描失败", + manual, + ) + + active_rj_status = [item for item in rj_status.values() if not item.skipped] + rj_ready = all(item.ready for item in active_rj_status) rj_status_text = "" - if rj_weekly_status: - rj_ready_count = sum(1 for s in rj_weekly_status.values() if s.ready) - rj_total = len(rj_weekly_status) - rj_weekly_ready = rj_ready_count == rj_total - - # 保存RJ状态到目录状态中 - rj_status_dict = {f"rj_weekly:{name}": status.to_dict() for name, status in rj_weekly_status.items()} - with self._status_lock: - self._directory_status.update(rj_status_dict) - - if not rj_weekly_ready: - rj_status_text = f",RJ周数据 {rj_ready_count}/{rj_total} 个目录就绪" + if active_rj_status: + rj_ready_count = sum(1 for item in active_rj_status if item.ready) + rj_total = len(active_rj_status) + if not rj_ready: + rj_status_text = f",RJ 数据 {rj_ready_count}/{rj_total} 个有效目录就绪" # 两个条件都满足才触发 - if daily_ready and rj_weekly_ready: - self._mark_ready(target_days, directory_status, rj_weekly_status) + if daily_ready and rj_ready: + self._mark_ready(target_days, directory_status, rj_status) skipped_text = f",已跳过 {skipped_count} 个停推目录" if skipped_count else "" - rj_text = ",RJ周数据已就绪" if rj_weekly_status else "" + rj_text = ",RJ 数据已就绪" if active_rj_status else "" return self._finish_check( "marked_ready", f"远程数据已满足目标周 7 天{skipped_text}{rj_text},已写入就绪标识,下次检查将自动处理", @@ -232,6 +250,13 @@ class AutoScheduler: manual, ) + if not active_status: + return self._finish_check( + "waiting", + f"远程普通日数据未发现有效目录,无法触发处理{rj_status_text}", + manual, + ) + ready_count = sum(1 for item in active_status if item.ready) skipped_text = f",跳过 {skipped_count} 个停推目录" if skipped_count else "" return self._finish_check( @@ -245,7 +270,9 @@ class AutoScheduler: downloader: RemoteDataDownloader, scheduler: AutoSchedulerConfig, target_days: list[date], + excluded_directories: set[str] | None = None, ) -> dict[str, DirectoryReadyStatus]: + excluded = {self._normalize_directory_name(item) for item in (excluded_directories or set())} expected_directories = scheduler.expected_directories if expected_directories: return { @@ -255,6 +282,7 @@ class AutoScheduler: target_days, ) for directory in expected_directories + if not self._is_excluded_directory(directory, excluded) } files = self._safe_list_remote_zip_files(downloader, None) @@ -270,11 +298,7 @@ class AutoScheduler: ) } - grouped: dict[str, list[RemoteFileInfo]] = defaultdict(list) - for remote_file in files: - grouped[remote_file.parent or "."].append(remote_file) - - if not grouped: + if not files: return { ".": DirectoryReadyStatus( directory=".", @@ -286,6 +310,16 @@ class AutoScheduler: ) } + grouped: dict[str, list[RemoteFileInfo]] = defaultdict(list) + for remote_file in files: + parent = remote_file.parent or "." + if self._is_excluded_directory(parent, excluded): + continue + grouped[parent].append(remote_file) + + if not grouped: + return {} + return { directory: self._directory_ready_status(directory, directory_files, target_days) for directory, directory_files in sorted(grouped.items(), key=lambda item: item[0]) @@ -330,9 +364,11 @@ class AutoScheduler: required = set(target_days) found = { - file_date + target_day for remote_file in files - if (file_date := extract_file_date(remote_file.name)) in required + if (date_range := parse_file_date_range(remote_file.name)) + for target_day in date_range.covered_days() + if target_day in required } missing = [item for item in target_days if item not in found] return DirectoryReadyStatus( @@ -343,79 +379,108 @@ class AutoScheduler: file_count=len(files), ) - def _check_rj_weekly_ready( + def _check_rj_ready( self, downloader: RemoteDataDownloader, - week_offset: int, - ) -> dict[str, RJWeeklyDirectoryStatus]: - """检查RJ周数据目录是否就绪""" + target_days: list[date], + ) -> dict[str, RJDirectoryReadyStatus]: + """检查 RJ 数据目录是否就绪,自动识别目录最新文件是日粒度还是周粒度。""" rj_config = state.current_config().rj_data.normalized() if not rj_config.enabled: return {} - target_week_end = self._calculate_week_end_date(week_offset) - result: dict[str, RJWeeklyDirectoryStatus] = {} + result: dict[str, RJDirectoryReadyStatus] = {} for directory in rj_config.weekly_directories: - result[directory] = self._check_single_rj_directory( - downloader, directory, target_week_end - ) + result[directory] = self._check_single_rj_directory(downloader, directory, target_days) return result - def _calculate_week_end_date(self, week_offset: int) -> date: - """计算目标周的结束日期(周日)""" - today = date.today() - # 找到本周的周日 - days_since_sunday = today.weekday() + 1 # weekday(): 0=周一, 6=周日 - if days_since_sunday == 7: - days_since_sunday = 0 - this_sunday = today - timedelta(days=days_since_sunday) - # 根据偏移计算目标周的周日 - target_sunday = this_sunday + timedelta(weeks=week_offset) - return target_sunday - def _check_single_rj_directory( self, downloader: RemoteDataDownloader, directory: str, - target_week_end: date, - ) -> RJWeeklyDirectoryStatus: - """检查单个RJ目录是否包含目标周的文件""" + target_days: list[date], + ) -> RJDirectoryReadyStatus: + """检查单个 RJ 目录是否包含目标周数据。""" files = self._safe_list_remote_zip_files(downloader, directory) if files is None: - return RJWeeklyDirectoryStatus( + return RJDirectoryReadyStatus( directory=directory, ready=False, error="远程目录不存在或无法访问", ) if not files: - return RJWeeklyDirectoryStatus( + return RJDirectoryReadyStatus( directory=directory, - ready=False, + ready=True, file_count=0, + skipped=True, + skip_reason="目录为空,视为已停推并跳过", ) - # 查找匹配目标周的文件 - target_end_str = target_week_end.strftime("%Y%m%d") + "0000" - for remote_file in files: - match = RJ_WEEKLY_FILE_PATTERN.search(remote_file.name) - if match: - file_end_date = match.group(2) - if file_end_date == target_end_str: - return RJWeeklyDirectoryStatus( - directory=directory, - ready=True, - found=True, - file_name=remote_file.name, - file_count=len(files), - ) + parsed_files = [ + (remote_file, date_range) + for remote_file in files + if (date_range := parse_file_date_range(remote_file.name)) + ] + if not parsed_files: + return RJDirectoryReadyStatus( + directory=directory, + ready=False, + found_days=[], + missing_days=target_days, + file_count=len(files), + error="目录中未找到可识别日期的 ZIP 文件", + ) - return RJWeeklyDirectoryStatus( + latest_file, latest_range = max( + parsed_files, + key=lambda item: (item[1].start, item[1].end_exclusive, item[0].name), + ) + granularity = "daily" if latest_range.span_days <= 1 else "weekly" + required = set(target_days) + + found = { + target_day + for _, date_range in parsed_files + for target_day in date_range.covered_days() + if target_day in required + } + missing = [item for item in target_days if item not in found] + + if granularity == "daily": + return RJDirectoryReadyStatus( + directory=directory, + ready=not missing, + granularity=granularity, + found_days=sorted(found), + missing_days=missing, + file_name=latest_file.name, + file_count=len(files), + ) + + for remote_file in files: + date_range = parse_file_date_range(remote_file.name) + if date_range and date_range.covers_all(required): + return RJDirectoryReadyStatus( + directory=directory, + ready=True, + granularity=granularity, + found_days=target_days, + missing_days=[], + file_name=remote_file.name, + file_count=len(files), + ) + + return RJDirectoryReadyStatus( directory=directory, ready=False, - found=False, + granularity=granularity, + found_days=sorted(found), + missing_days=missing, + file_name=latest_file.name, file_count=len(files), ) @@ -457,7 +522,7 @@ class AutoScheduler: self, target_days: list[date], directory_status: dict[str, DirectoryReadyStatus], - rj_weekly_status: dict[str, RJWeeklyDirectoryStatus] | None = None, + rj_status: dict[str, RJDirectoryReadyStatus] | None = None, ) -> None: READY_DIR.mkdir(parents=True, exist_ok=True) payload = { @@ -470,10 +535,14 @@ class AutoScheduler: for name, item in directory_status.items() }, } - if rj_weekly_status: + if rj_status: + payload["rj_directories"] = { + name: item.to_dict() + for name, item in rj_status.items() + } payload["rj_weekly_directories"] = { name: item.to_dict() - for name, item in rj_weekly_status.items() + for name, item in rj_status.items() } READY_FLAG.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8") @@ -546,12 +615,39 @@ class AutoScheduler: "status": self.get_status(), } - def _set_directory_status(self, directory_status: dict[str, DirectoryReadyStatus]) -> None: + def _set_combined_directory_status( + self, + directory_status: dict[str, DirectoryReadyStatus], + rj_status: dict[str, RJDirectoryReadyStatus], + ) -> None: with self._status_lock: - self._directory_status = { + combined = { name: item.to_dict() for name, item in directory_status.items() } + combined.update( + { + f"rj:{name}": item.to_dict() + for name, item in rj_status.items() + } + ) + self._directory_status = combined + + @staticmethod + def _normalize_directory_name(directory: str) -> str: + normalized = str(directory or "").replace("\\", "/").strip().strip("/") + return "." if normalized in {"", "."} else normalized + + @classmethod + def _is_excluded_directory(cls, directory: str, excluded: set[str]) -> bool: + if not excluded: + return False + normalized = cls._normalize_directory_name(directory) + return any( + normalized == item or normalized.startswith(f"{item}/") + for item in excluded + if item != "." + ) def _set_next_check(self, value: datetime) -> None: with self._status_lock: diff --git a/app/services/remote_download.py b/app/services/remote_download.py index f10102a..cf8fa03 100644 --- a/app/services/remote_download.py +++ b/app/services/remote_download.py @@ -9,7 +9,7 @@ from pathlib import Path from typing import Callable, Iterable from app.config import RemoteDataConfig -from app.utils.file_dates import extract_file_date, select_recent_items_by_directory +from app.utils.file_dates import parse_file_date_range, select_recent_items_by_directory LogFn = Callable[[str], None] @@ -93,7 +93,13 @@ class RemoteDataDownloader: selected = [ remote_file for remote_file in files - if extract_file_date(remote_file.name) in target_dates + if ( + date_range := parse_file_date_range(remote_file.name) + ) and ( + date_range.covers_all(target_dates) + if date_range.span_days > 1 + else date_range.covers_any(target_dates) + ) ] selected_files.extend(selected) skipped_count = len(files) - len(selected) diff --git a/app/utils/file_dates.py b/app/utils/file_dates.py index ed1534c..7bfd132 100644 --- a/app/utils/file_dates.py +++ b/app/utils/file_dates.py @@ -23,18 +23,84 @@ class DirectoryDateSelection: start_date: date | None = None +@dataclass(frozen=True) +class FileDateRange: + start: date + end_exclusive: date + timestamp_count: int + + @property + def span_days(self) -> int: + return max((self.end_exclusive - self.start).days, 1) + + def covered_days(self) -> list[date]: + return [self.start + timedelta(days=offset) for offset in range(self.span_days)] + + def covers_any(self, target_days: set[date]) -> bool: + return any(day in target_days for day in self.covered_days()) + + def covers_all(self, target_days: set[date]) -> bool: + return target_days.issubset(set(self.covered_days())) + + def extract_file_date(filename: str | Path) -> date | None: """Extract the first timestamp date from a filename.""" - match = FILENAME_DATE_RE.search(Path(str(filename)).name) - if not match: + date_range = parse_file_date_range(filename) + return date_range.start if date_range else None + + +def parse_file_date_range(filename: str | Path) -> FileDateRange | None: + """Extract the covered natural date range from a filename. + + Files with one timestamp are treated as one-day data files. Files with two + timestamps use the second timestamp as an exclusive end when it is midnight, + matching names like xxx_202605110000_202605120000.zip. + """ + values = FILENAME_DATE_RE.findall(Path(str(filename)).name) + if not values: return None - return parse_file_timestamp(match.group(1)) + + start_dt = parse_file_timestamp_datetime(values[0]) + if not start_dt: + return None + + if len(values) == 1: + return FileDateRange( + start=start_dt.date(), + end_exclusive=start_dt.date() + timedelta(days=1), + timestamp_count=1, + ) + + end_dt = parse_file_timestamp_datetime(values[1]) + if not end_dt or end_dt <= start_dt: + return FileDateRange( + start=start_dt.date(), + end_exclusive=start_dt.date() + timedelta(days=1), + timestamp_count=len(values), + ) + + end_exclusive = end_dt.date() + if end_dt.time() != datetime.min.time(): + end_exclusive += timedelta(days=1) + if end_exclusive <= start_dt.date(): + end_exclusive = start_dt.date() + timedelta(days=1) + + return FileDateRange( + start=start_dt.date(), + end_exclusive=end_exclusive, + timestamp_count=len(values), + ) def parse_file_timestamp(value: str) -> date | None: + parsed = parse_file_timestamp_datetime(value) + return parsed.date() if parsed else None + + +def parse_file_timestamp_datetime(value: str) -> datetime | None: fmt = "%Y%m%d%H%M%S" if len(value) == 14 else "%Y%m%d%H%M" try: - return datetime.strptime(value, fmt).date() + return datetime.strptime(value, fmt) except ValueError: return None diff --git a/docs/auto_scheduler_plan.md b/docs/auto_scheduler_plan.md index c5209a8..7d6d75d 100644 --- a/docs/auto_scheduler_plan.md +++ b/docs/auto_scheduler_plan.md @@ -21,25 +21,22 @@ | 格式 | 示例 | 数据日期 | |------|------|----------| -| `XXX_YYYYMMDDHHMM_YYYYMMDDHHMM` | `CapacityReportData2.6_202605110000_202605120000.zip` | 第一个时间戳 `202605110000` → 2026-05-11 | -| `XXX_YYYYMMDDHHMM` | `CapacityReportData2.6_202605110000.zip` | 时间戳 `202605110000` → 2026-05-11 | +| `XXX_YYYYMMDDHHMM_YYYYMMDDHHMM` | `CapacityReportData2.6_202605110000_202605120000.zip` | 覆盖起止时间范围,结束零点按右开区间处理 → 2026-05-11 | +| `XXX_YYYYMMDDHHMM` | `CapacityReportData2.6_202605110000.zip` | 视为单日文件 → 2026-05-11 | ### 2.2 解析逻辑 -复用项目已有的正则 `_ZIP_DATE_RE = re.compile(r"(? date | None: - """从文件名提取数据日期(取第一个匹配的时间戳)""" - match = _ZIP_DATE_RE.search(filename) - if not match: - return None - timestamp = match.group(1) # 12位: YYYYMMDDHHMM 或 14位: YYYYMMDDHHMMSS - fmt = "%Y%m%d%H%M%S" if len(timestamp) == 14 else "%Y%m%d%H%M" - try: - return datetime.strptime(timestamp, fmt).date() - except ValueError: +def parse_file_date_range(filename: str) -> FileDateRange | None: + """从文件名提取数据覆盖范围。""" + values = _ZIP_DATE_RE.findall(filename) + if not values: return None + start = parse_timestamp(values[0]) + end = parse_timestamp(values[1]) if len(values) > 1 else start + timedelta(days=1) + return FileDateRange(start=start.date(), end_exclusive=normalize_end(end)) ``` --- @@ -130,10 +127,12 @@ def get_target_week_range(week_offset: int = 0) -> tuple[date, date]: **关键**: 1. 不依赖系统日期,而是从文件名中提取日期 -2. 文件名只取**第一个**时间戳作为数据日期 -3. 例如文件 `CapacityReportData2.6_202605120000_202605130000.zip` → 数据日期为 2026-05-12 +2. 文件名带两个时间戳时按覆盖范围判断;只有一个时间戳时视为单日文件 +3. 例如文件 `CapacityReportData2.6_202605120000_202605130000.zip` → 覆盖 2026-05-12 4. 如果某个日期没有对应的 ZIP 文件,该目录判定为未就绪 +RJ 目录会额外按最新文件自动判断粒度:最新文件为单日时要求 7 个目标日都齐全;最新文件为多日/周文件时要求存在一个 ZIP 覆盖完整目标自然周。 + ### 3.4 标识文件 - **位置**:`cache/auto_scheduler/ready.flag` diff --git a/docs/project_context.md b/docs/project_context.md index 094bddf..bb52dda 100644 --- a/docs/project_context.md +++ b/docs/project_context.md @@ -1,5 +1,15 @@ # 项目上下文记录 +## 2026-06-01:完善 RJ 自动调度日/周粒度识别 + +- `app/utils/file_dates.py` 新增文件日期范围解析:`XXX_YYYYMMDDHHMM` 视为单日文件,`XXX_YYYYMMDDHHMM_YYYYMMDDHHMM` 按起止时间展开自然日,结束时间为零点时按右开区间处理。 +- `app/services/auto_scheduler.py` 的 RJ 检查改为按目录最新 ZIP 自动识别 `daily` 或 `weekly`:日粒度目录要求目标自然周 7 天都存在,周粒度目录要求有一个 ZIP 覆盖目标自然周;空 RJ 目录继续视为停推并跳过。 +- 自动调度普通 4G/5G 目录扫描会排除已配置的 RJ 目录,避免 `expected_directories=[]` 时 RJ 周目录被普通 7 天规则误判阻塞。 +- `app/services/remote_download.py` 的调度下载筛选改为使用日期覆盖范围:单日文件只要覆盖目标日即下载,多日/周文件必须覆盖完整目标周才下载,避免 ready 后漏下或误下 RJ 周文件。 +- `Configure.json` 将 `RJ/700M/700RJGD`、`RJ/700M/700RJYD` 加入 RJ 数据目录,并补充 `700MRJGD`、`700MRJYD` 字段映射;`processor.py` 避免多个源字段别名映射到同一目标字段时生成重复列。 +- 配置上传接口现在会导入/保存 `RJData`,前端类型补充 `rj_data` 和调度状态中的 `granularity` 字段。 +- 已验证:后端 AST 语法检查、`Configure.json` JSON 解析、`npm run build`、真实 SFTP 清单识别、MySQL 临时导入 700RJYD 样本并清理测试表均通过;前端构建仅保留既有大 chunk 警告。 + ## 2026-05-29:新增 RJ 周数据处理功能 - 新增 `RJData` 配置块到 `app/config.py`,支持 `enabled`、`weekly_directories` 和 `table_field_mappings` 配置项,用于管理 RJ 周数据目录和字段映射。 diff --git a/frontend/src/types.ts b/frontend/src/types.ts index 2f15e05..125ab3c 100644 --- a/frontend/src/types.ts +++ b/frontend/src/types.ts @@ -79,10 +79,17 @@ export interface AppConfig { }; remote_data: RemoteDataConfig; history_retention: HistoryRetentionConfig; + rj_data?: RJDataConfig; sheet_filter: string[]; extract_fields: Array>; } +export interface RJDataConfig { + enabled: boolean; + weekly_directories: string[]; + table_field_mappings: Record>>; +} + export interface RemoteDataConfig { enabled: boolean; protocol: 'ftp' | 'sftp'; @@ -134,6 +141,7 @@ export interface RemoteSchedulerStatus { string, { ready: boolean; + granularity?: 'daily' | 'weekly' | string | null; found_days: string[]; missing_days: string[]; found_count: number;