feat: 新增 RJ 周数据处理功能

- 新增 RJData 配置块,支持周数据目录和字段映射
- 自动调度增加 RJ 周文件检查(与现有 7 天检查为 AND 关系)
- 处理器支持 RJ 表专用字段映射
- 目标表: 2_6GRJGD, 2_6GRJYD
- 字段: 开始时间, 结束时间, gNBId, cellId, gNBplmn, 上下行总流量_GB

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-05-29 13:26:13 +08:00
co-authored by Claude Opus 4.8
parent 272ff38070
commit 44f4aaa79c
6 changed files with 354 additions and 56 deletions
+60 -3
View File
@@ -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
}
+92 -29
View File
@@ -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 + 连接复用)"""
+139 -5
View File
@@ -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: