feat: 统一数据目录映射配置

This commit is contained in:
2026-06-24 16:45:28 +08:00
parent 850573feeb
commit 0a4012ad4e
11 changed files with 442 additions and 310 deletions
+17 -4
View File
@@ -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)
+71 -68
View File
@@ -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
}
+31 -88
View File
@@ -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
+61 -58
View File
@@ -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
+40 -30
View File
@@ -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] = {}
+3 -14
View File
@@ -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,
}