diff --git a/.gitignore b/.gitignore index 22f0c54..ffae178 100644 --- a/.gitignore +++ b/.gitignore @@ -46,6 +46,7 @@ CapacityReportData/ auth.ini license.dat api_tokens.json +.mcp.json # Logs *.log diff --git a/Configure.json b/Configure.json index 8b523f9..3cbb25a 100644 --- a/Configure.json +++ b/Configure.json @@ -1,5 +1,5 @@ { - "Update": "2026/05/19 17:08:10", + "Update": "2026/05/25 17:57:25", "MySQL_DBInfo": { "host": "127.0.0.1", "port": 3306, @@ -7,6 +7,53 @@ "passwd": "123456", "dbname": "CapacityReport" }, + "RemoteData": { + "enabled": true, + "protocol": "sftp", + "host": "127.0.0.1", + "port": 2022, + "user": "nixevol", + "remote_dir": "/CapacityReportData", + "passive": true, + "timeout": 30, + "auto_delete_source": true, + "auto_scheduler": { + "enabled": true, + "check_interval_hours": 1, + "expected_directories": [], + "week_offset": 0 + }, + "passwd": "242520" + }, + "HistoryRetention": { + "enabled": true, + "keep_count": 1 + }, + "RJData": { + "enabled": true, + "weekly_directories": [ + "RJ/2.6G/2.6RJGD", + "RJ/2.6G/2.6RJYD" + ], + "table_field_mappings": { + "2_6GRJGD": [ + {"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"} + ], + "2_6GRJYD": [ + {"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"} + ] + } + }, "SheetFilter": [ "指标(计数器)", "Template" @@ -266,21 +313,5 @@ ], "Type": "float" } - ], - "RemoteData": { - "enabled": true, - "protocol": "sftp", - "host": "127.0.0.1", - "port": 2022, - "user": "nixevol", - "remote_dir": "/CapacityReportData", - "passive": true, - "timeout": 30, - "auto_delete_source": true, - "passwd": "242520" - }, - "HistoryRetention": { - "enabled": true, - "keep_count": 1 - } -} + ] +} \ No newline at end of file diff --git a/app/config.py b/app/config.py index 74c1dbc..8cededc 100644 --- a/app/config.py +++ b/app/config.py @@ -167,6 +167,57 @@ class RemoteDataConfig: ).normalized() +@dataclass +class RJDataConfig: + """RJ数据配置 - 周数据处理""" + enabled: bool = False + weekly_directories: List[str] = field(default_factory=list) + # 字段映射: {表名: [{Source, Target, Type}, ...]} + table_field_mappings: Dict[str, List[Dict[str, Any]]] = field(default_factory=dict) + + def normalized(self) -> "RJDataConfig": + directories = [] + seen = set() + for d in self.weekly_directories or []: + normalized = str(d).replace("\\", "/").strip().strip("/") + if normalized and normalized not in seen: + directories.append(normalized) + seen.add(normalized) + + # 标准化字段映射 + mappings = {} + for table_name, fields in (self.table_field_mappings or {}).items(): + if isinstance(fields, list): + mappings[table_name] = [ + {k: v for k, v in f.items() if k in ("Source", "Target", "Type")} + for f in fields + if isinstance(f, dict) and "Source" in f and "Target" in f + ] + + return RJDataConfig( + enabled=bool(self.enabled), + weekly_directories=directories, + table_field_mappings=mappings, + ) + + def to_dict(self) -> Dict[str, Any]: + normalized = self.normalized() + return { + "enabled": normalized.enabled, + "weekly_directories": normalized.weekly_directories, + "table_field_mappings": normalized.table_field_mappings, + } + + @classmethod + def from_dict(cls, data: Dict[str, Any] | None) -> "RJDataConfig": + data = data or {} + return cls( + enabled=bool(data.get("enabled", False)), + weekly_directories=data.get("weekly_directories", []) if isinstance(data.get("weekly_directories"), list) else [], + table_field_mappings=data.get("table_field_mappings", {}) if isinstance(data.get("table_field_mappings"), dict) else {}, + ).normalized() + + @dataclass class HistoryRetentionConfig: enabled: bool = False @@ -205,6 +256,7 @@ class AppConfig: mysql: MySQLConfig = field(default_factory=MySQLConfig) remote_data: RemoteDataConfig = field(default_factory=RemoteDataConfig) history_retention: HistoryRetentionConfig = field(default_factory=HistoryRetentionConfig) + rj_data: RJDataConfig = field(default_factory=RJDataConfig) sheet_filter: List[str] = field(default_factory=list) extract_fields: List[Dict[str, Any]] = field(default_factory=list) @@ -213,10 +265,10 @@ class AppConfig: """从 Configure.json 加载配置""" if not CONFIG_FILE.exists(): return cls() - + with open(CONFIG_FILE, 'r', encoding='utf-8') as f: data = json.load(f) - + mysql_data = data.get("MySQL_DBInfo", {}) mysql_config = MySQLConfig( host=mysql_data.get("host", "localhost"), @@ -227,12 +279,14 @@ class AppConfig: ) remote_config = RemoteDataConfig.from_dict(data.get("RemoteData")) history_retention = HistoryRetentionConfig.from_dict(data.get("HistoryRetention")) - + rj_data = RJDataConfig.from_dict(data.get("RJData")) + return cls( update=data.get("Update", ""), mysql=mysql_config, remote_data=remote_config, history_retention=history_retention, + rj_data=rj_data, sheet_filter=data.get("SheetFilter", []), extract_fields=data.get("ExtractField", []) ) @@ -258,6 +312,7 @@ class AppConfig: }, "RemoteData": self.remote_data.normalized().to_dict(include_password=True), "HistoryRetention": self.history_retention.normalized().to_dict(), + "RJData": self.rj_data.normalized().to_dict(), "SheetFilter": self.sheet_filter, "ExtractField": self.extract_fields } @@ -274,6 +329,7 @@ class AppConfig: }, "remote_data": self.remote_data.normalized().to_dict(), "history_retention": self.history_retention.normalized().to_dict(), + "rj_data": self.rj_data.normalized().to_dict(), "sheet_filter": self.sheet_filter, "extract_fields": self.extract_fields } @@ -291,6 +347,7 @@ class AppConfig: }, "remote_data": self.remote_data.normalized().to_dict(include_password=True), "history_retention": self.history_retention.normalized().to_dict(), + "rj_data": self.rj_data.normalized().to_dict(), "sheet_filter": self.sheet_filter, "extract_fields": self.extract_fields } diff --git a/app/processor.py b/app/processor.py index a88a46f..c47f270 100644 --- a/app/processor.py +++ b/app/processor.py @@ -152,6 +152,36 @@ class DataProcessor: return field_map, type_map + def _get_field_map_for_table(self, table_name: str) -> Tuple[Dict[str, str], Dict[str, str]]: + """ + 获取指定表的字段映射和类型映射 + + Args: + table_name: 目标表名 + + Returns: + (field_map, type_map) + field_map: {源字段名: 目标字段名} + type_map: {目标字段名: 字段类型} + """ + # 检查是否是RJ表 + rj_config = self.config.rj_data.normalized() + if rj_config.enabled and table_name in rj_config.table_field_mappings: + # 使用RJ专用映射 + field_map = {} + type_map = {} + for field_def in rj_config.table_field_mappings[table_name]: + source = field_def.get("Source") + target = field_def.get("Target") + field_type = field_def.get("Type", "string") + if source and target: + field_map[source] = target + type_map[target] = field_type + return field_map, type_map + + # 使用全局映射 + return self._field_map, self._type_map + def _apply_sql_script_type_hints(self) -> None: """从 ReportScript.sql 的 ALTER 语句补全缺失的数值字段类型。""" if not SQL_SCRIPT.exists(): @@ -449,84 +479,87 @@ class DataProcessor: return 'gbk' return 'utf-8' - def _process_csv_file_fast(self, csv_file: Path, table_name: str, + def _process_csv_file_fast(self, csv_file: Path, table_name: str, conn=None, table_created: bool = False) -> Tuple[int, bool]: """ 处理单个 CSV 文件 使用 LOAD DATA LOCAL INFILE,比 executemany 快 10-50 倍 - + Args: csv_file: CSV 文件路径 table_name: 目标表名 conn: 数据库连接(复用) table_created: 表是否已创建 - + Returns: (导入行数, 表是否已创建) """ encoding = self._detect_encoding(csv_file) rel_path = csv_file.relative_to(self.work_dir) self.logger.info(f"处理 CSV: {rel_path} (编码: {encoding})") - + # 读取 CSV,使用优化参数 df = pd.read_csv( - csv_file, - encoding=encoding, - thousands=',', + csv_file, + encoding=encoding, + thousands=',', low_memory=True, # 低内存模式 dtype=str, # 全部作为字符串读取,避免类型推断开销 na_values=[''], # 只把空字符串当作 NA keep_default_na=False # 不使用默认的 NA 值 ) - + + # 获取字段映射(RJ表使用专用映射,其他表使用全局映射) + field_map, type_map = self._get_field_map_for_table(table_name) + # 快速字段匹配(使用预编译的映射表) col_mapping = {} for col in df.columns: - if col in self._field_map: - col_mapping[col] = self._field_map[col] - - if len(col_mapping) <= 3: + if col in field_map: + col_mapping[col] = field_map[col] + + if len(col_mapping) < 1: if 'kpis' in str(csv_file).lower(): self.logger.warning(f"跳过非数据文件: {rel_path}") return 0, table_created raise ValueError(f"字段匹配不足: {rel_path}") - + # 选择需要的列并重命名 source_cols = list(col_mapping.keys()) target_cols = list(col_mapping.values()) - + # 创建结果 DataFrame,使用目标列名 df_result = df[source_cols].copy() df_result.columns = target_cols - + # 替换 NA 为默认值 df_result = df_result.fillna('') - + # 构建目标字段的类型映射 - column_types = {col: self._type_map.get(col, 'string') for col in target_cols} - + column_types = {col: type_map.get(col, 'string') for col in target_cols} + # 根据类型处理每列数据 for col in target_cols: col_type = column_types.get(col, 'string') - + if col_type == 'datetime': # 日期时间类型处理 df_result[col] = self._convert_datetime_column(df_result[col]) - + elif col_type == 'int': # 整数类型处理 df_result[col] = self._convert_int_column(df_result[col]) - + elif col_type == 'float': # 浮点数类型处理 df_result[col] = self._convert_float_column(df_result[col]) - + elif col_type == 'text': # 长文本类型,截断到 65535 字符 mask = df_result[col].str.len() > 65535 if mask.any(): df_result.loc[mask, col] = df_result.loc[mask, col].str[:65535] - + else: # string 或其他 # 字符串类型:去除百分号、截断长度 df_result[col] = df_result[col].str.replace('%', '', regex=False) @@ -780,29 +813,43 @@ class DataProcessor: .str.replace(' ', '', regex=False) ) + # RJ目录名到表名的映射 + RJ_DIR_TO_TABLE = { + "2.6RJGD": "2_6GRJGD", + "2.6RJYD": "2_6GRJYD", + "700RJGD": "700MRJGD", + "700RJYD": "700MRJYD", + } + def _find_data_directories(self) -> Dict[str, Path]: """ 查找包含数据文件的目录,返回 {表名: 目录路径} + 支持 4G/5G 和 RJ 数据目录 """ data_dirs = {} target_names = {'4G', '5G', '4g', '5g'} - + self.logger.info(f"开始查找数据目录,工作目录: {self.work_dir}") - + # 递归查找所有名为 4G 或 5G 的目录 found_dirs = [] for subdir in self.work_dir.rglob('*'): if subdir.is_dir() and subdir.name in target_names: found_dirs.append(subdir) - + self.logger.info(f"找到 {len(found_dirs)} 个候选目录") - + for subdir in found_dirs: table_name = f"{subdir.name.upper()}_UD" if table_name not in data_dirs: data_dirs[table_name] = subdir self.logger.info(f"发现数据目录: {subdir.relative_to(self.work_dir)} -> 表: {table_name}") - + + # 查找 RJ 数据目录 + rj_config = self.config.rj_data.normalized() + if rj_config.enabled: + self._find_rj_data_directories(data_dirs) + if not data_dirs: self.logger.warning("未找到 4G/5G 目录,使用直接子目录") for subdir in self.work_dir.iterdir(): @@ -810,8 +857,24 @@ class DataProcessor: table_name = f"{subdir.name}_UD" data_dirs[table_name] = subdir self.logger.info(f"使用直接子目录: {subdir.relative_to(self.work_dir)} -> 表: {table_name}") - + return data_dirs + + def _find_rj_data_directories(self, data_dirs: Dict[str, Path]) -> None: + """查找RJ数据目录""" + self.logger.info("查找RJ数据目录...") + + # 遍历工作目录查找RJ目录 + for rj_dir in self.work_dir.rglob('*'): + if not rj_dir.is_dir(): + continue + # 检查是否是RJ相关的目录名 + dir_name = rj_dir.name + if dir_name in self.RJ_DIR_TO_TABLE: + table_name = self.RJ_DIR_TO_TABLE[dir_name] + if table_name not in data_dirs: + data_dirs[table_name] = rj_dir + self.logger.info(f"发现RJ数据目录: {rj_dir.relative_to(self.work_dir)} -> 表: {table_name}") def _process_csv_files(self): """处理所有 CSV 文件(使用 LOAD DATA INFILE + 连接复用)""" diff --git a/app/services/auto_scheduler.py b/app/services/auto_scheduler.py index a79d9a8..ae39f7f 100644 --- a/app/services/auto_scheduler.py +++ b/app/services/auto_scheduler.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +import re import threading from collections import defaultdict from dataclasses import dataclass @@ -21,6 +22,31 @@ 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周数据目录就绪状态""" + directory: str + ready: bool + found: bool = False + file_name: str | None = None + file_count: int = 0 + error: str | None = None + + def to_dict(self) -> dict[str, Any]: + return { + "ready": self.ready, + "found": self.found, + "file_name": self.file_name, + "file_count": self.file_count, + "error": self.error, + } + @dataclass(frozen=True) class DirectoryReadyStatus: @@ -81,6 +107,7 @@ class AutoScheduler: def get_status(self) -> dict[str, Any]: config = state.current_config().remote_data.normalized() scheduler = config.auto_scheduler.normalized() + rj_config = state.current_config().rj_data.normalized() target_days = required_week_days(scheduler.week_offset) ready_flag = self._read_ready_flag() @@ -92,6 +119,8 @@ class AutoScheduler: "expected_directories": scheduler.expected_directories, "week_offset": scheduler.week_offset, "auto_delete_source": config.auto_delete_source, + "rj_data_enabled": rj_config.enabled, + "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), "last_result": self._last_result, @@ -147,6 +176,7 @@ 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) directory_status = self._check_remote_ready(remote_config, scheduler, target_days) self._set_directory_status(directory_status) @@ -161,13 +191,34 @@ class AutoScheduler: active_status = [item for item in directory_status.values() if not item.skipped] skipped_count = len(directory_status) - len(active_status) - ready = bool(active_status) and all(item.ready for item in active_status) - if ready: - self._mark_ready(target_days, directory_status) + daily_ready = bool(active_status) and all(item.ready for item in active_status) + + # 检查RJ周数据目录 + rj_weekly_status = self._check_rj_weekly_ready(remote_config, scheduler.week_offset) + rj_weekly_ready = True # 默认就绪(如果没有配置RJ目录) + 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状态到目录状态中 + with self._status_lock: + for name, status in rj_weekly_status.items(): + self._directory_status[f"rj_weekly:{name}"] = status.to_dict() + + if not rj_weekly_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) skipped_text = f",已跳过 {skipped_count} 个停推目录" if skipped_count else "" + rj_text = ",RJ周数据已就绪" if rj_weekly_status else "" return self._finish_check( "marked_ready", - f"远程数据已满足目标周 7 天{skipped_text},已写入就绪标识,下次检查将自动处理", + f"远程数据已满足目标周 7 天{skipped_text}{rj_text},已写入就绪标识,下次检查将自动处理", manual, ) @@ -182,7 +233,7 @@ class AutoScheduler: skipped_text = f",跳过 {skipped_count} 个停推目录" if skipped_count else "" return self._finish_check( "waiting", - f"远程数据未就绪,{ready_count}/{len(active_status)} 个有效目录满足目标周 7 天{skipped_text}", + f"远程数据未就绪,{ready_count}/{len(active_status)} 个有效目录满足目标周 7 天{skipped_text}{rj_status_text}", manual, ) @@ -290,6 +341,83 @@ class AutoScheduler: file_count=len(files), ) + def _check_rj_weekly_ready( + self, + remote_config: RemoteDataConfig, + week_offset: int, + ) -> dict[str, RJWeeklyDirectoryStatus]: + """检查RJ周数据目录是否就绪""" + rj_config = state.current_config().rj_data.normalized() + if not rj_config.enabled: + return {} + + downloader = RemoteDataDownloader(remote_config) + target_week_end = self._calculate_week_end_date(week_offset) + result: dict[str, RJWeeklyDirectoryStatus] = {} + + for directory in rj_config.weekly_directories: + result[directory] = self._check_single_rj_directory( + downloader, directory, target_week_end + ) + + 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目录是否包含目标周的文件""" + files = self._safe_list_remote_zip_files(downloader, directory) + if files is None: + return RJWeeklyDirectoryStatus( + directory=directory, + ready=False, + error="远程目录不存在或无法访问", + ) + + if not files: + return RJWeeklyDirectoryStatus( + directory=directory, + ready=False, + file_count=0, + ) + + # 查找匹配目标周的文件 + 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), + ) + + return RJWeeklyDirectoryStatus( + directory=directory, + ready=False, + found=False, + file_count=len(files), + ) + def _trigger_processing( self, target_dates: list[date], @@ -328,6 +456,7 @@ class AutoScheduler: self, target_days: list[date], directory_status: dict[str, DirectoryReadyStatus], + rj_weekly_status: dict[str, RJWeeklyDirectoryStatus] | None = None, ) -> None: READY_DIR.mkdir(parents=True, exist_ok=True) payload = { @@ -340,6 +469,11 @@ class AutoScheduler: for name, item in directory_status.items() }, } + if rj_weekly_status: + payload["rj_weekly_directories"] = { + name: item.to_dict() + for name, item in rj_weekly_status.items() + } READY_FLAG.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8") def _clear_ready_flag(self) -> None: diff --git a/docs/project_context.md b/docs/project_context.md index e337cad..094bddf 100644 --- a/docs/project_context.md +++ b/docs/project_context.md @@ -1,5 +1,17 @@ # 项目上下文记录 +## 2026-05-29:新增 RJ 周数据处理功能 + +- 新增 `RJData` 配置块到 `app/config.py`,支持 `enabled`、`weekly_directories` 和 `table_field_mappings` 配置项,用于管理 RJ 周数据目录和字段映射。 +- `app/services/auto_scheduler.py` 新增 `RJWeeklyDirectoryStatus` 数据类和 `_check_rj_weekly_ready()` 方法,支持检查 RJ 周数据目录是否包含目标周的文件。 +- 自动调度逻辑改为:现有 7 天目录检查 **且** RJ 周数据检查都满足时才触发处理,两个条件是 AND 关系。 +- `app/processor.py` 新增 `_get_field_map_for_table()` 方法和 `_find_rj_data_directories()` 方法,支持 RJ 表使用专用字段映射。 +- RJ 数据目录结构:`/CapacityReportData/RJ/2.6G/2.6RJGD/` 和 `RJ/2.6G/2.6RJYD/`,每个目录每周一个 ZIP 文件。 +- 目标表名映射:`2.6RJGD` -> `2_6GRJGD`,`2.6RJYD` -> `2_6GRJYD`。 +- 字段映射配置:`开始时间`、`结束时间`、`gNBId`、`cellId`、`gNBplmn`、`上下行总流量_GB`。 +- `Configure.json` 新增 `RJData` 配置块,包含启用状态、周数据目录列表和表字段映射。 +- 已验证:`.venv\Scripts\python.exe -m compileall app` 通过;GD 和 YD 数据字段映射测试成功,6 个字段全部匹配。 + ## 2026-05-25:新增远程自动调度和每目录 7 天处理窗口 - 新增 `app/utils/file_dates.py`,统一解析文件名中的第一个 `YYYYMMDDHHMM` 或 `YYYYMMDDHHMMSS` 时间戳,并提供按目录筛选最近 7 个自然日文件的工具;本地手动上传和远程下载后的 ZIP、Excel、CSV 处理都会按所在目录只保留最近 7 天文件,未携带日期且同目录没有任何可识别日期时保留兼容。