From 9711aa6270bc630017b9b91ab79db36d1069a964 Mon Sep 17 00:00:00 2001 From: Nixevol Date: Fri, 29 May 2026 13:45:09 +0800 Subject: [PATCH] =?UTF-8?q?refactor:=20=E4=BC=98=E5=8C=96RJ=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E5=A4=84=E7=90=86=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 优化 get_status() 中重复调用 state.current_config() - 优化 _check_and_run_locked() 中重复创建 RemoteDataDownloader - 预构建RJ字段映射缓存,避免每次调用 _get_field_map_for_table 时重新解析 - 优化 _find_rj_data_directories(),优先使用配置路径而非全量递归 Co-Authored-By: Claude Opus 4.8 --- app/processor.py | 80 ++++++++++++++++++++++------------ app/services/auto_scheduler.py | 25 ++++++----- 2 files changed, 66 insertions(+), 39 deletions(-) diff --git a/app/processor.py b/app/processor.py index c47f270..5e84ddc 100644 --- a/app/processor.py +++ b/app/processor.py @@ -118,14 +118,36 @@ class DataProcessor: self._explicit_type_fields: set[str] = set() self._generated_csv_files: set[Path] = set() self._generated_csv_lock = Lock() - + # 预编译字段映射,避免重复查找 self._field_map, self._type_map = self._build_field_map() self._apply_sql_script_type_hints() - + + # 预编译RJ字段映射 + self._rj_field_maps: Dict[str, Tuple[Dict[str, str], Dict[str, str]]] = {} + self._build_rj_field_maps() + # LOAD DATA INFILE 支持状态(在首次使用时检测) self._load_data_supported: Optional[bool] = None self._load_data_checked = False + + def _build_rj_field_maps(self) -> None: + """预构建RJ表的字段映射""" + rj_config = self.config.rj_data.normalized() + if not rj_config.enabled: + return + + for table_name, fields in rj_config.table_field_mappings.items(): + field_map = {} + type_map = {} + for field_def in fields: + 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 + self._rj_field_maps[table_name] = (field_map, type_map) def _build_field_map(self) -> Tuple[Dict[str, str], Dict[str, str]]: """ @@ -164,20 +186,9 @@ class DataProcessor: 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 + # 检查是否是RJ表(使用预构建的缓存) + if table_name in self._rj_field_maps: + return self._rj_field_maps[table_name] # 使用全局映射 return self._field_map, self._type_map @@ -864,17 +875,32 @@ class DataProcessor: """查找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}") + # RJ目录结构: RJ/2.6G/2.6RJGD, RJ/2.6G/2.6RJYD 等 + # 只搜索工作目录下的特定路径,避免全量递归 + rj_config = self.config.rj_data.normalized() + for weekly_dir in rj_config.weekly_directories: + # 构建本地路径: work_dir / RJ/2.6G/2.6RJGD + rj_path = self.work_dir / weekly_dir + if rj_path.exists() and rj_path.is_dir(): + dir_name = rj_path.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_path + self.logger.info(f"发现RJ数据目录: {rj_path.relative_to(self.work_dir)} -> 表: {table_name}") + + # 兜底: 如果配置的路径不存在,尝试从工作目录中查找 + if not any(k in data_dirs for k in self.RJ_DIR_TO_TABLE.values()): + self.logger.info("配置的RJ路径不存在,尝试从工作目录中查找...") + for rj_dir in self.work_dir.rglob('*'): + if not rj_dir.is_dir(): + continue + 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 ae39f7f..db94699 100644 --- a/app/services/auto_scheduler.py +++ b/app/services/auto_scheduler.py @@ -105,9 +105,10 @@ class AutoScheduler: self._thread.join(timeout=5) def get_status(self) -> dict[str, Any]: - config = state.current_config().remote_data.normalized() + app_config = state.current_config() + config = app_config.remote_data.normalized() scheduler = config.auto_scheduler.normalized() - rj_config = state.current_config().rj_data.normalized() + rj_config = app_config.rj_data.normalized() target_days = required_week_days(scheduler.week_offset) ready_flag = self._read_ready_flag() @@ -152,7 +153,8 @@ class AutoScheduler: def _check_and_run_locked(self, manual: bool) -> dict[str, Any]: now = datetime.now() - remote_config = state.current_config().remote_data.normalized() + app_config = state.current_config() + remote_config = app_config.remote_data.normalized() scheduler = remote_config.auto_scheduler.normalized() self._last_check_at = now @@ -178,7 +180,8 @@ class AutoScheduler: # 检查现有7天目录(4G/5G数据) target_days = required_week_days(scheduler.week_offset) - directory_status = self._check_remote_ready(remote_config, scheduler, target_days) + downloader = RemoteDataDownloader(remote_config) + directory_status = self._check_remote_ready(downloader, scheduler, target_days) self._set_directory_status(directory_status) error_count = sum(1 for item in directory_status.values() if item.error) @@ -193,8 +196,8 @@ 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周数据目录 - rj_weekly_status = self._check_rj_weekly_ready(remote_config, scheduler.week_offset) + # 检查RJ周数据目录(复用同一个 downloader) + rj_weekly_status = self._check_rj_weekly_ready(downloader, scheduler.week_offset) rj_weekly_ready = True # 默认就绪(如果没有配置RJ目录) rj_status_text = "" @@ -204,9 +207,9 @@ class AutoScheduler: 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: - for name, status in rj_weekly_status.items(): - self._directory_status[f"rj_weekly:{name}"] = status.to_dict() + self._directory_status.update(rj_status_dict) if not rj_weekly_ready: rj_status_text = f",RJ周数据 {rj_ready_count}/{rj_total} 个目录就绪" @@ -239,11 +242,10 @@ class AutoScheduler: def _check_remote_ready( self, - remote_config: RemoteDataConfig, + downloader: RemoteDataDownloader, scheduler: AutoSchedulerConfig, target_days: list[date], ) -> dict[str, DirectoryReadyStatus]: - downloader = RemoteDataDownloader(remote_config) expected_directories = scheduler.expected_directories if expected_directories: return { @@ -343,7 +345,7 @@ class AutoScheduler: def _check_rj_weekly_ready( self, - remote_config: RemoteDataConfig, + downloader: RemoteDataDownloader, week_offset: int, ) -> dict[str, RJWeeklyDirectoryStatus]: """检查RJ周数据目录是否就绪""" @@ -351,7 +353,6 @@ class AutoScheduler: if not rj_config.enabled: return {} - downloader = RemoteDataDownloader(remote_config) target_week_end = self._calculate_week_end_date(week_offset) result: dict[str, RJWeeklyDirectoryStatus] = {}