From 0a4012ad4ec22b283f0074e5d5cac0f3ac083156 Mon Sep 17 00:00:00 2001 From: Nixevol Date: Wed, 24 Jun 2026 16:45:28 +0800 Subject: [PATCH] =?UTF-8?q?=EF=BB=BFfeat:=20=E7=BB=9F=E4=B8=80=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E7=9B=AE=E5=BD=95=E6=98=A0=E5=B0=84=E9=85=8D=E7=BD=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Configure.json | 15 +- app/api/routers/config.py | 21 ++- app/config.py | 139 ++++++++-------- app/processor.py | 119 ++++---------- app/services/auto_scheduler.py | 119 +++++++------- app/services/csv_processor.py | 70 ++++---- app/services/pipeline.py | 17 +- docs/project_context.md | 9 + frontend/src/components/SettingsPanel.vue | 191 ++++++++++++++++++---- frontend/src/styles.css | 38 ++++- frontend/src/types.ts | 14 +- 11 files changed, 442 insertions(+), 310 deletions(-) diff --git a/Configure.json b/Configure.json index 63ce54e..485e616 100644 --- a/Configure.json +++ b/Configure.json @@ -29,13 +29,14 @@ "enabled": true, "keep_count": 1 }, - "RJData": { - "enabled": true, - "weekly_directories": [ - "RJ/2.6G/2.6RJGD", - "RJ/2.6G/2.6RJYD", - "RJ/700M/700RJGD", - "RJ/700M/700RJYD" + "DataMappings": { + "directories": [ + {"path": "4G", "table": "4G_UD", "ready_rule": "daily"}, + {"path": "5G", "table": "5G_UD", "ready_rule": "daily"}, + {"path": "RJ/2.6G/2.6RJGD", "table": "2_6GRJGD", "ready_rule": "auto"}, + {"path": "RJ/2.6G/2.6RJYD", "table": "2_6GRJYD", "ready_rule": "auto"}, + {"path": "RJ/700M/700RJGD", "table": "700MRJGD", "ready_rule": "auto"}, + {"path": "RJ/700M/700RJYD", "table": "700MRJYD", "ready_rule": "auto"} ], "table_field_mappings": { "2_6GRJGD": [ diff --git a/app/api/routers/config.py b/app/api/routers/config.py index 96bfad1..88dcecf 100644 --- a/app/api/routers/config.py +++ b/app/api/routers/config.py @@ -8,9 +8,9 @@ from fastapi.responses import Response from app import state from app.config import ( + DataMappingsConfig, HistoryRetentionConfig, MetrixConfig, - RJDataConfig, RemoteDataConfig, SOURCE_TYPES, WAREHOUSE_TYPES, @@ -80,6 +80,18 @@ async def update_metrix_config(config: dict[str, Any] = Body(...)): return {"success": True, "message": "Metrix 连接配置已更新", "update": state.config.update} +@router.post("/api/config/data-mappings") +async def update_data_mappings_config(config: dict[str, Any] = Body(...)): + state.reload_config() + current = state.config.data_mappings.normalized() + payload = dict(config) + if "table_field_mappings" not in payload: + payload["table_field_mappings"] = current.table_field_mappings + state.config.data_mappings = DataMappingsConfig.from_dict(payload) + state.config.save() + return {"success": True, "message": "数据目录映射已更新", "update": state.config.update} + + @router.post("/api/config/history-retention") async def update_history_retention(config: dict[str, Any] = Body(...)): state.reload_config() @@ -149,6 +161,10 @@ def _apply_config_data(data: dict[str, Any]) -> None: if isinstance(metrix_data, dict): state.config.metrix = MetrixConfig.from_dict(metrix_data) + data_mappings = data.get("DataMappings") + if isinstance(data_mappings, dict): + state.config.data_mappings = DataMappingsConfig.from_dict(data_mappings) + mysql_data = data.get("MySQL_DBInfo") if isinstance(mysql_data, dict): for key in ("host", "port", "user", "passwd", "dbname"): @@ -169,6 +185,3 @@ def _apply_config_data(data: dict[str, Any]) -> None: 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/config.py b/app/config.py index dc7cd4a..6902b9e 100644 --- a/app/config.py +++ b/app/config.py @@ -46,7 +46,7 @@ class MetrixConfig: Metrix appears in the UI as two connection types ("存储平台"/"数据库平台") that share the same base_url + token; storage_id is the file source, database_conn_id + target_database - are the warehouse. data_dir_to_table maps data sub-dirs to staging tables (Metrix mode only). + are the warehouse. """ base_url: str = "http://host.docker.internal:8000" token: str = "" @@ -54,18 +54,12 @@ class MetrixConfig: database_conn_id: str = "" target_database: str = "" recent_days: int = 7 - data_dir_to_table: Dict[str, str] = field(default_factory=lambda: {"4G": "4G_UD", "5G": "5G_UD"}) def normalized(self) -> "MetrixConfig": try: recent_days = max(int(self.recent_days), 1) except (TypeError, ValueError): recent_days = 7 - mapping = { - str(k).strip(): str(v).strip() - for k, v in (self.data_dir_to_table or {}).items() - if str(k).strip() and str(v).strip() - } return MetrixConfig( base_url=str(self.base_url or "").strip(), token=str(self.token or "").strip(), @@ -73,7 +67,6 @@ class MetrixConfig: database_conn_id=str(self.database_conn_id or "").strip(), target_database=str(self.target_database or "").strip(), recent_days=recent_days, - data_dir_to_table=mapping or {"4G": "4G_UD", "5G": "5G_UD"}, ) def to_dict(self, include_token: bool = False) -> Dict[str, Any]: @@ -84,7 +77,6 @@ class MetrixConfig: "database_conn_id": n.database_conn_id, "target_database": n.target_database, "recent_days": n.recent_days, - "data_dir_to_table": n.data_dir_to_table, } if include_token: data["token"] = n.token @@ -93,7 +85,6 @@ class MetrixConfig: @classmethod def from_dict(cls, data: Dict[str, Any] | None) -> "MetrixConfig": data = data or {} - mapping = data.get("data_dir_to_table") return cls( base_url=str(data.get("base_url", "http://host.docker.internal:8000")), token=str(data.get("token", "")), @@ -101,7 +92,70 @@ class MetrixConfig: database_conn_id=str(data.get("database_conn_id", "")), target_database=str(data.get("target_database", "")), recent_days=data.get("recent_days", 7), - data_dir_to_table=mapping if isinstance(mapping, dict) else {"4G": "4G_UD", "5G": "5G_UD"}, + ).normalized() + + +@dataclass +class DataMappingsConfig: + """Source directories mapped to staging tables.""" + directories: List[Dict[str, str]] = field( + default_factory=lambda: [ + {"path": "4G", "table": "4G_UD", "ready_rule": "daily"}, + {"path": "5G", "table": "5G_UD", "ready_rule": "daily"}, + ] + ) + table_field_mappings: Dict[str, List[Dict[str, Any]]] = field(default_factory=dict) + + def normalized(self) -> "DataMappingsConfig": + directories = [] + seen = set() + for item in self.directories or []: + if not isinstance(item, dict): + continue + path = str(item.get("path", "")).replace("\\", "/").strip().strip("/") + table = str(item.get("table", "")).strip() + ready_rule = str(item.get("ready_rule", "daily")).strip().lower() + if ready_rule not in {"daily", "auto"}: + ready_rule = "daily" + if not path or not table: + continue + key = (path, table, ready_rule) + if key in seen: + continue + seen.add(key) + directories.append({"path": path, "table": table, "ready_rule": ready_rule}) + + if not directories: + directories = [ + {"path": "4G", "table": "4G_UD", "ready_rule": "daily"}, + {"path": "5G", "table": "5G_UD", "ready_rule": "daily"}, + ] + + 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 DataMappingsConfig(directories=directories, table_field_mappings=mappings) + + def to_dict(self) -> Dict[str, Any]: + normalized = self.normalized() + return { + "directories": normalized.directories, + "table_field_mappings": normalized.table_field_mappings, + } + + @classmethod + def from_dict(cls, data: Dict[str, Any] | None) -> "DataMappingsConfig": + data = data or {} + directories = data.get("directories", []) + return cls( + directories=directories if isinstance(directories, list) else [], + table_field_mappings=data.get("table_field_mappings", {}) if isinstance(data.get("table_field_mappings"), dict) else {}, ).normalized() @@ -238,57 +292,6 @@ 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 @@ -341,9 +344,9 @@ class AppConfig: warehouse_type: str = "mysql" mysql: MySQLConfig = field(default_factory=MySQLConfig) metrix: MetrixConfig = field(default_factory=MetrixConfig) + data_mappings: DataMappingsConfig = field(default_factory=DataMappingsConfig) 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) @@ -366,8 +369,8 @@ class AppConfig: ) remote_config = RemoteDataConfig.from_dict(data.get("RemoteData")) metrix_config = MetrixConfig.from_dict(data.get("Metrix")) + data_mappings = DataMappingsConfig.from_dict(data.get("DataMappings")) history_retention = HistoryRetentionConfig.from_dict(data.get("HistoryRetention")) - rj_data = RJDataConfig.from_dict(data.get("RJData")) return cls( update=data.get("Update", ""), @@ -375,9 +378,9 @@ class AppConfig: warehouse_type=_normalize_warehouse_type(data.get("WarehouseType")), mysql=mysql_config, metrix=metrix_config, + data_mappings=data_mappings, remote_data=remote_config, history_retention=history_retention, - rj_data=rj_data, sheet_filter=data.get("SheetFilter", []), extract_fields=data.get("ExtractField", []) ) @@ -404,9 +407,9 @@ class AppConfig: "dbname": self.mysql.dbname }, "Metrix": self.metrix.normalized().to_dict(include_token=True), + "DataMappings": self.data_mappings.normalized().to_dict(), "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 } @@ -424,9 +427,9 @@ class AppConfig: "dbname": self.mysql.dbname }, "metrix": self.metrix.normalized().to_dict(), + "data_mappings": self.data_mappings.normalized().to_dict(), "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 } @@ -445,9 +448,9 @@ class AppConfig: "dbname": self.mysql.dbname }, "metrix": self.metrix.normalized().to_dict(include_token=True), + "data_mappings": self.data_mappings.normalized().to_dict(), "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 a6a05c1..9a362f3 100644 --- a/app/processor.py +++ b/app/processor.py @@ -123,21 +123,18 @@ class DataProcessor: 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() + # 预编译按表字段映射 + self._table_field_maps: Dict[str, Tuple[Dict[str, str], Dict[str, str]]] = {} + self._build_table_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(): + def _build_table_field_maps(self) -> None: + """预构建按表覆盖的字段映射""" + mappings = self.config.data_mappings.normalized().table_field_mappings + for table_name, fields in mappings.items(): field_map = {} type_map = {} for field_def in fields: @@ -147,7 +144,7 @@ class DataProcessor: if source and target: field_map[source] = target type_map[target] = field_type - self._rj_field_maps[table_name] = (field_map, type_map) + self._table_field_maps[table_name] = (field_map, type_map) def _build_field_map(self) -> Tuple[Dict[str, str], Dict[str, str]]: """ @@ -186,9 +183,8 @@ class DataProcessor: field_map: {源字段名: 目标字段名} type_map: {目标字段名: 字段类型} """ - # 检查是否是RJ表(使用预构建的缓存) - if table_name in self._rj_field_maps: - return self._rj_field_maps[table_name] + if table_name in self._table_field_maps: + return self._table_field_maps[table_name] # 使用全局映射 return self._field_map, self._type_map @@ -520,7 +516,7 @@ class DataProcessor: keep_default_na=False # 不使用默认的 NA 值 ) - # 获取字段映射(RJ表使用专用映射,其他表使用全局映射) + # 获取字段映射(按表配置优先,其他表使用全局映射) field_map, type_map = self._get_field_map_for_table(table_name) # 快速字段匹配(使用预编译的映射表) @@ -829,84 +825,28 @@ 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]: + def _find_data_directories(self) -> Dict[str, List[Path]]: """ - 查找包含数据文件的目录,返回 {表名: 目录路径} - 支持 4G/5G 和 RJ 数据目录 + 查找包含数据文件的目录,返回 {表名: [目录路径]} """ - data_dirs = {} - target_names = {'4G', '5G', '4g', '5g'} - + data_dirs: Dict[str, List[Path]] = {} 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(): - if subdir.is_dir(): - table_name = f"{subdir.name}_UD" - data_dirs[table_name] = subdir - self.logger.info(f"使用直接子目录: {subdir.relative_to(self.work_dir)} -> 表: {table_name}") + for item in self.config.data_mappings.normalized().directories: + subdir = self.work_dir / item["path"] + if subdir.exists() and subdir.is_dir(): + self._add_data_directory(data_dirs, item["table"], subdir) + self.logger.info(f"发现数据目录: {subdir.relative_to(self.work_dir)} -> 表: {item['table']}") return data_dirs - def _find_rj_data_directories(self, data_dirs: Dict[str, Path]) -> None: - """查找RJ数据目录""" - self.logger.info("查找RJ数据目录...") + @staticmethod + def _add_data_directory(data_dirs: Dict[str, List[Path]], table_name: str, directory: Path) -> None: + existing = data_dirs.setdefault(table_name, []) + resolved = directory.resolve() + if all(path.resolve() != resolved for path in existing): + existing.append(directory) - # 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 + 连接复用)""" self.logger.info("正在处理 CSV 文件并上传到数据库...") @@ -919,14 +859,17 @@ class DataProcessor: return # 按目录分组处理,使用连接复用 - for table_name, subdir in data_dirs.items(): - self.logger.info(f"处理目录: {subdir.relative_to(self.work_dir)} -> 表: {table_name}") + for table_name, subdirs in data_dirs.items(): + directory_names = ", ".join(str(subdir.relative_to(self.work_dir)) for subdir in subdirs) + self.logger.info(f"处理目录: {directory_names} -> 表: {table_name}") # 删除旧表 self.db.drop_table(table_name) # 处理该目录下的所有 CSV - csv_files = self._filter_recent_files(list(self._scan_files(subdir, ['.csv'])), "CSV", root=subdir) + csv_files: list[Path] = [] + for subdir in subdirs: + csv_files.extend(self._filter_recent_files(list(self._scan_files(subdir, ['.csv'])), "CSV", root=subdir)) self.logger.info(f"找到 {len(csv_files)} 个 CSV 文件") total_rows = 0 diff --git a/app/services/auto_scheduler.py b/app/services/auto_scheduler.py index a725e1a..a0be2db 100644 --- a/app/services/auto_scheduler.py +++ b/app/services/auto_scheduler.py @@ -24,8 +24,8 @@ FAILURE_RESULTS = {"scan_failed", "trigger_failed", "source_cleanup_failed", "fa @dataclass(frozen=True) -class RJDirectoryReadyStatus: - """RJ 数据目录就绪状态,按最新文件自动识别日粒度或周粒度。""" +class AutoDirectoryReadyStatus: + """自动粒度目录就绪状态,按最新文件自动识别日粒度或周粒度。""" directory: str ready: bool granularity: str | None = None @@ -115,7 +115,12 @@ class AutoScheduler: app_config = state.current_config() config = app_config.remote_data.normalized() scheduler = config.auto_scheduler.normalized() - rj_config = app_config.rj_data.normalized() + data_mappings = app_config.data_mappings.normalized() + auto_directories = [ + item["path"] + for item in data_mappings.directories + if item.get("ready_rule") == "auto" + ] target_days = required_week_days(scheduler.week_offset) ready_flag = self._read_ready_flag() @@ -127,9 +132,7 @@ 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_directories": rj_config.weekly_directories, - "rj_weekly_directories": rj_config.weekly_directories, + "auto_ready_directories": auto_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, @@ -188,15 +191,19 @@ class AutoScheduler: target_days = required_week_days(scheduler.week_offset) downloader = make_source_downloader(app_config) - rj_config = app_config.rj_data.normalized() - rj_directories = set(rj_config.weekly_directories) if rj_config.enabled else set() + data_mappings = app_config.data_mappings.normalized() + auto_directories = { + item["path"] + for item in data_mappings.directories + if item.get("ready_rule") == "auto" + } - # 检查现有7天目录(4G/5G数据),RJ 目录单独按粒度判断,避免被普通7天规则误拦。 + # 自动粒度目录单独判断,避免被普通 7 天规则误拦。 directory_status = self._check_remote_ready( downloader, scheduler, target_days, - excluded_directories=rj_directories, + excluded_directories=auto_directories, ) error_count = sum(1 for item in directory_status.values() if item.error) @@ -210,37 +217,37 @@ class AutoScheduler: active_status = [item for item in directory_status.values() if not item.skipped] skipped_count = len(directory_status) - len(active_status) - daily_ready = bool(active_status) and all(item.ready for item in active_status) - # 检查 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: + # 检查自动粒度目录,按最新文件自动判断日粒度或周粒度。 + auto_status = self._check_auto_ready(downloader, target_days, sorted(auto_directories)) + auto_error_count = sum(1 for item in auto_status.values() if item.error) + self._set_combined_directory_status(directory_status, auto_status) + if auto_error_count: return self._finish_check( "scan_failed", - f"RJ 远程目录扫描失败,{rj_error_count}/{len(rj_status)} 个目录无法访问或扫描失败", + f"自动粒度远程目录扫描失败,{auto_error_count}/{len(auto_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 = "" + active_auto_status = [item for item in auto_status.values() if not item.skipped] + daily_ready = all(item.ready for item in active_status) if active_status else not directory_status and bool(active_auto_status) + auto_ready = all(item.ready for item in active_auto_status) + auto_status_text = "" - 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 active_auto_status: + auto_ready_count = sum(1 for item in active_auto_status if item.ready) + auto_total = len(active_auto_status) + if not auto_ready: + auto_status_text = f",自动粒度数据 {auto_ready_count}/{auto_total} 个有效目录就绪" # 两个条件都满足才触发 - if daily_ready and rj_ready: - self._mark_ready(target_days, directory_status, rj_status) + if daily_ready and auto_ready: + self._mark_ready(target_days, directory_status, auto_status) skipped_text = f",已跳过 {skipped_count} 个停推目录" if skipped_count else "" - rj_text = ",RJ 数据已就绪" if active_rj_status else "" + auto_text = ",自动粒度数据已就绪" if active_auto_status else "" return self._finish_check( "marked_ready", - f"远程数据已满足目标周 7 天{skipped_text}{rj_text},已写入就绪标识,下次检查将自动处理", + f"远程数据已满足目标周 7 天{skipped_text}{auto_text},已写入就绪标识,下次检查将自动处理", manual, ) @@ -254,7 +261,7 @@ class AutoScheduler: if not active_status: return self._finish_check( "waiting", - f"远程普通日数据未发现有效目录,无法触发处理{rj_status_text}", + f"远程普通日数据未发现有效目录,无法触发处理{auto_status_text}", manual, ) @@ -262,7 +269,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}{rj_status_text}", + f"远程数据未就绪,{ready_count}/{len(active_status)} 个有效目录满足目标周 7 天{skipped_text}{auto_status_text}", manual, ) @@ -380,40 +387,40 @@ class AutoScheduler: file_count=len(files), ) - def _check_rj_ready( + def _check_auto_ready( self, downloader: RemoteDataDownloader, target_days: list[date], - ) -> dict[str, RJDirectoryReadyStatus]: - """检查 RJ 数据目录是否就绪,自动识别目录最新文件是日粒度还是周粒度。""" - rj_config = state.current_config().rj_data.normalized() - if not rj_config.enabled: + directories: list[str], + ) -> dict[str, AutoDirectoryReadyStatus]: + """检查自动粒度目录是否就绪,自动识别目录最新文件是日粒度还是周粒度。""" + if not directories: return {} - result: dict[str, RJDirectoryReadyStatus] = {} + result: dict[str, AutoDirectoryReadyStatus] = {} - for directory in rj_config.weekly_directories: - result[directory] = self._check_single_rj_directory(downloader, directory, target_days) + for directory in directories: + result[directory] = self._check_single_auto_directory(downloader, directory, target_days) return result - def _check_single_rj_directory( + def _check_single_auto_directory( self, downloader: RemoteDataDownloader, directory: str, target_days: list[date], - ) -> RJDirectoryReadyStatus: - """检查单个 RJ 目录是否包含目标周数据。""" + ) -> AutoDirectoryReadyStatus: + """检查单个自动粒度目录是否包含目标周数据。""" files = self._safe_list_remote_zip_files(downloader, directory) if files is None: - return RJDirectoryReadyStatus( + return AutoDirectoryReadyStatus( directory=directory, ready=False, error="远程目录不存在或无法访问", ) if not files: - return RJDirectoryReadyStatus( + return AutoDirectoryReadyStatus( directory=directory, ready=True, file_count=0, @@ -427,7 +434,7 @@ class AutoScheduler: if (date_range := parse_file_date_range(remote_file.name)) ] if not parsed_files: - return RJDirectoryReadyStatus( + return AutoDirectoryReadyStatus( directory=directory, ready=False, found_days=[], @@ -452,7 +459,7 @@ class AutoScheduler: missing = [item for item in target_days if item not in found] if granularity == "daily": - return RJDirectoryReadyStatus( + return AutoDirectoryReadyStatus( directory=directory, ready=not missing, granularity=granularity, @@ -465,7 +472,7 @@ class AutoScheduler: 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( + return AutoDirectoryReadyStatus( directory=directory, ready=True, granularity=granularity, @@ -475,7 +482,7 @@ class AutoScheduler: file_count=len(files), ) - return RJDirectoryReadyStatus( + return AutoDirectoryReadyStatus( directory=directory, ready=False, granularity=granularity, @@ -523,7 +530,7 @@ class AutoScheduler: self, target_days: list[date], directory_status: dict[str, DirectoryReadyStatus], - rj_status: dict[str, RJDirectoryReadyStatus] | None = None, + auto_status: dict[str, AutoDirectoryReadyStatus] | None = None, ) -> None: READY_DIR.mkdir(parents=True, exist_ok=True) payload = { @@ -536,14 +543,10 @@ class AutoScheduler: for name, item in directory_status.items() }, } - if rj_status: - payload["rj_directories"] = { + if auto_status: + payload["auto_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_status.items() + for name, item in auto_status.items() } READY_FLAG.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8") @@ -619,7 +622,7 @@ class AutoScheduler: def _set_combined_directory_status( self, directory_status: dict[str, DirectoryReadyStatus], - rj_status: dict[str, RJDirectoryReadyStatus], + auto_status: dict[str, AutoDirectoryReadyStatus], ) -> None: with self._status_lock: combined = { @@ -629,7 +632,7 @@ class AutoScheduler: combined.update( { name: item.to_dict() - for name, item in rj_status.items() + for name, item in auto_status.items() } ) self._directory_status = combined diff --git a/app/services/csv_processor.py b/app/services/csv_processor.py index f354f0c..ee80e42 100644 --- a/app/services/csv_processor.py +++ b/app/services/csv_processor.py @@ -41,13 +41,9 @@ class CsvProcessor: self.log = log self.recent_days = int(config.get("recent_days", 7)) self.sheet_filter = set(config.get("sheet_filter", [])) - self.data_dir_to_table = {k.upper(): v for k, v in (config.get("data_dir_to_table") or {}).items()} + self.directories = _build_directory_mappings(config.get("directories") or []) self.field_map, self.type_map = _build_global_map(config.get("extract_fields", [])) - rj = config.get("rj") or {} - self.rj_enabled = bool(rj.get("enabled")) - self.rj_weekly_dirs = list(rj.get("weekly_directories") or []) - self.rj_dir_to_table = dict(rj.get("dir_to_table") or {}) - self.rj_maps = _build_rj_maps(rj.get("table_field_mappings") or {}) + self.table_maps = _build_table_maps(config.get("table_field_mappings") or {}) self.out_dir = self.work_dir / ".out" def process(self) -> dict[str, Path]: @@ -99,8 +95,10 @@ class CsvProcessor: return {} self.out_dir.mkdir(parents=True, exist_ok=True) result: dict[str, Path] = {} - for table, directory in data_dirs.items(): - csv_files = self._filter_recent(list(self._scan(directory, (".csv",))), "CSV", root=directory) + for table, directories in data_dirs.items(): + csv_files: list[Path] = [] + for directory in directories: + csv_files.extend(self._filter_recent(list(self._scan(directory, (".csv",))), "CSV", root=directory)) if not csv_files: continue field_map, type_map = self._maps_for_table(table) @@ -185,30 +183,17 @@ class CsvProcessor: return out[union] # --- directory detection -------------------------------------------- - def _find_data_dirs(self) -> dict[str, Path]: - data_dirs: dict[str, Path] = {} - target_names = set(self.data_dir_to_table.keys()) - for sub in self.work_dir.rglob("*"): - if sub.is_dir() and sub.name.upper() in target_names: - table = self.data_dir_to_table[sub.name.upper()] - data_dirs.setdefault(table, sub) - if self.rj_enabled: - self._find_rj_dirs(data_dirs) + def _find_data_dirs(self) -> dict[str, list[Path]]: + data_dirs: dict[str, list[Path]] = {} + for item in self.directories: + directory = self.work_dir / item["path"] + if directory.exists() and directory.is_dir(): + _add_data_dir(data_dirs, item["table"], directory) return data_dirs - def _find_rj_dirs(self, data_dirs: dict[str, Path]) -> None: - for weekly in self.rj_weekly_dirs: - path = self.work_dir / weekly - if path.exists() and path.is_dir() and path.name in self.rj_dir_to_table: - data_dirs.setdefault(self.rj_dir_to_table[path.name], path) - if not any(table in data_dirs for table in self.rj_dir_to_table.values()): - for sub in self.work_dir.rglob("*"): - if sub.is_dir() and sub.name in self.rj_dir_to_table: - data_dirs.setdefault(self.rj_dir_to_table[sub.name], sub) - def _maps_for_table(self, table: str): - if table in self.rj_maps: - return self.rj_maps[table] + if table in self.table_maps: + return self.table_maps[table] return self.field_map, self.type_map # --- helpers --------------------------------------------------------- @@ -257,7 +242,32 @@ def _build_global_map(extract_fields: list[dict]) -> tuple[dict[str, str], dict[ return field_map, type_map -def _build_rj_maps(table_field_mappings: dict) -> dict[str, tuple[dict[str, str], dict[str, str]]]: +def _build_directory_mappings(items: list[dict]) -> list[dict[str, str]]: + mappings: list[dict[str, str]] = [] + seen: set[tuple[str, str]] = set() + for item in items: + if not isinstance(item, dict): + continue + path = str(item.get("path", "")).replace("\\", "/").strip().strip("/") + table = str(item.get("table", "")).strip() + if not path or not table: + continue + key = (path, table) + if key in seen: + continue + seen.add(key) + mappings.append({"path": path, "table": table}) + return mappings + + +def _add_data_dir(data_dirs: dict[str, list[Path]], table: str, directory: Path) -> None: + existing = data_dirs.setdefault(table, []) + resolved = directory.resolve() + if all(path.resolve() != resolved for path in existing): + existing.append(directory) + + +def _build_table_maps(table_field_mappings: dict) -> dict[str, tuple[dict[str, str], dict[str, str]]]: maps: dict[str, tuple[dict[str, str], dict[str, str]]] = {} for table, fields in table_field_mappings.items(): field_map: dict[str, str] = {} diff --git a/app/services/pipeline.py b/app/services/pipeline.py index cffb580..1aee6e4 100644 --- a/app/services/pipeline.py +++ b/app/services/pipeline.py @@ -11,12 +11,6 @@ from app.processor import ProcessLogger from app.services.csv_processor import CsvProcessor from app.services.platform import make_client -RJ_DIR_TO_TABLE = { - "2.6RJGD": "2_6GRJGD", - "2.6RJYD": "2_6GRJYD", - "700RJGD": "700MRJGD", - "700RJYD": "700MRJYD", -} RESULT_TABLES = ["4G_结果表", "5G_结果表"] @@ -34,18 +28,13 @@ def validate_metrix(metrix: MetrixConfig) -> None: def build_processor_config(app_config: AppConfig) -> dict: metrix = app_config.metrix.normalized() - rj = app_config.rj_data.normalized() + mappings = app_config.data_mappings.normalized() return { "recent_days": metrix.recent_days, "sheet_filter": list(app_config.sheet_filter), - "data_dir_to_table": dict(metrix.data_dir_to_table), + "directories": list(mappings.directories), "extract_fields": app_config.extract_fields, - "rj": { - "enabled": rj.enabled, - "weekly_directories": rj.weekly_directories, - "dir_to_table": RJ_DIR_TO_TABLE, - "table_field_mappings": rj.table_field_mappings, - }, + "table_field_mappings": mappings.table_field_mappings, } diff --git a/docs/project_context.md b/docs/project_context.md index 4cd48b2..5646bf5 100644 --- a/docs/project_context.md +++ b/docs/project_context.md @@ -694,3 +694,12 @@ - 现状确认:`pipeline.py`/`warehouse.py` 所有 `run_script` 调用本就只传 `content`(报表 SQL 来自 `read_report_sql()` 读取本地 `ReportScript.sql`),从不传 `script_id`,行为已正确。 - 改动:`app/services/platform.py::run_script` 删除一直未被调用的 `script_id` 参数与对应 body 分支,方法只构造 `content`,从代码层面杜绝走平台库内脚本那条路;签名由 `(conn_id, script_id=None, content="", ...)` 改为 `(conn_id, content="", ...)`,现有调用全用 `content=` 关键字、`conn_id` 位置参,未受影响。 - 验证:`python -m compileall app` 通过;全仓 `app` 内除该行注释外无 `script_id` 引用。 + +## 2026-06-24:数据目录映射统一为 DataMappings + +- 背景:UD 与 RJ 在处理阶段本质都是「源目录 -> 暂存表」。此前 UD 使用 `UDData.directories`,RJ 使用 `RJData.weekly_directories` + 代码内置目录名到表名映射,概念重复且前端容易继续堆卡片。 +- 配置:删除 `UDData` / `RJData` 两套配置,改为顶层 `DataMappings`。`DataMappings.directories` 每行结构为 `{path, table, ready_rule}`,`ready_rule=daily` 表示目标周每日 7 天检查,`ready_rule=auto` 表示按目录最新 ZIP 自动识别日粒度或周粒度;RJ 原字段映射并入 `DataMappings.table_field_mappings`,按目标表名覆盖全局字段映射。`MetrixConfig` 仍只保留平台连接信息。 +- 处理链路:直连 MySQL 的 `DataProcessor` 与 Metrix 仓库模式的 `CsvProcessor` 都只读取 `DataMappings.directories`。一个表可对应多个目录:每个目录先按最近日期筛选 CSV,再合并导入同一张暂存表(Metrix 模式写 `.out/{table}.csv`,MySQL 模式逐文件导入同表)。表级字段映射存在时优先使用 `DataMappings.table_field_mappings[table]`,否则使用全局 `ExtractField`。 +- 自动调度:原 RJ 专用检查改为通用自动粒度检查。`ready_rule=auto` 的目录会从普通每日扫描中排除,并单独按最新 ZIP 判断日/周粒度;如果配置里只有自动粒度目录,只要这些目录就绪也可触发调度。 +- 前端:设置页「规则映射」左侧独立滚动配置栏中只保留一个「数据目录映射」卡片,每行可编辑目录、暂存表与就绪规则,并保存到 `/api/config/data-mappings`。后续新增目录类映射继续加同一张表,不再新增配置卡片。 +- 验证:`python -m compileall -q app` 通过;`frontend` `npm run build` 通过(仅既有大 chunk 提示);构建产物与 Python 缓存已清理。 diff --git a/frontend/src/components/SettingsPanel.vue b/frontend/src/components/SettingsPanel.vue index 3078f8e..6338ed6 100644 --- a/frontend/src/components/SettingsPanel.vue +++ b/frontend/src/components/SettingsPanel.vue @@ -337,36 +337,73 @@
- - -

匹配这些关键词的 Sheet 将被跳过处理

-
- - {{ filter }} - - -
- - - 添加 - -
- - -
+ + + + + + +

匹配这些关键词的 Sheet 将被跳过处理

+
+ + {{ filter }} + + +
+ + + 添加 + +
+ + +
+