""" 二期数据引擎 支持两种数据源: 1. 全市数据(city_*.parquet) — 默认,266所学校,16个区 2. 区级数据(era2_*.parquet) — 旧流程兼容,176所学校,13个区 """ import pandas as pd import numpy as np from pathlib import Path from typing import Dict, List, Optional import logging from config_era2 import ( PARQUET_BASIC_INFO, PARQUET_COURSE_IMPL, PARQUET_SUBJECT_IMPL, CITY_PARQUET_BASIC_INFO, CITY_PARQUET_COURSE_IMPL, CITY_PARQUET_SUBJECT_IMPL, SCHOOL_NAME_FULL_TO_SHORT, SUBJECTS, ) logger = logging.getLogger(__name__) class DataEngineEra2: """ 二期数据引擎 默认使用全市数据(city_*.parquet),通过 district_filter 筛选特定区。 同时提供全市基准数据接口供 StatsEngine 使用。 """ def __init__(self, district_filter: Optional[str] = None, use_city_data: bool = True): """ Args: district_filter: 可选,只加载某个区的数据(如 "杨浦区") None 表示加载全部区域 use_city_data: 是否使用全市数据(默认True)。False则用旧的区级parquet """ self.district_filter = district_filter self.use_city_data = use_city_data and CITY_PARQUET_COURSE_IMPL.exists() self._basic_info: Optional[pd.DataFrame] = None self._course_impl: Optional[pd.DataFrame] = None self._subject_impl: Optional[pd.DataFrame] = None # 全市数据(不受 district_filter 限制,用于全市基准) self._city_course_impl: Optional[pd.DataFrame] = None self._city_subject_impl: Optional[pd.DataFrame] = None def load_all(self) -> None: """加载所有数据""" src = "全市" if self.use_city_data else "区级" logger.info(f"开始加载二期parquet数据({src}数据源)...") if self.use_city_data: self._load_city_data() else: self._load_legacy_data() logger.info( f"数据加载完成: 基础信息={len(self._basic_info)}行, " f"课程实施={len(self._course_impl)}行, " f"学科课程={len(self._subject_impl)}行, " f"当前区学校数={len(self.schools)}" ) if self._city_course_impl is not None: logger.info(f"全市基准: {self._city_course_impl['学校名称'].nunique()}校") @staticmethod def _strip_str_columns(df: pd.DataFrame, columns: List[str]) -> pd.DataFrame: """对指定列做 strip(),去除空格和换行符""" for col in columns: if col in df.columns: df[col] = df[col].astype(str).str.strip() return df def _load_city_data(self) -> None: """从全市parquet加载(默认路径)""" # 全市课程实施(全量,用于基准) df_course_all = pd.read_parquet(CITY_PARQUET_COURSE_IMPL) df_course_all = df_course_all.rename(columns={"学校简称": "学校名称"}) df_course_all = self._strip_str_columns(df_course_all, ["学校名称", "字段名称", "所在区"]) self._city_course_impl = df_course_all df_subject_all = pd.read_parquet(CITY_PARQUET_SUBJECT_IMPL) df_subject_all = df_subject_all.rename(columns={"学校简称": "学校名称"}) df_subject_all = self._strip_str_columns(df_subject_all, ["学校名称", "字段名称", "所在区", "学科"]) self._city_subject_impl = df_subject_all # 按区筛选的数据(用于赋分和报告生成) if self.district_filter: self._course_impl = df_course_all[df_course_all["所在区"] == self.district_filter].copy() self._subject_impl = df_subject_all[df_subject_all["所在区"] == self.district_filter].copy() else: self._course_impl = df_course_all.copy() self._subject_impl = df_subject_all.copy() # 基础信息表 df_basic = pd.read_parquet(CITY_PARQUET_BASIC_INFO) df_basic = self._strip_str_columns(df_basic, ["学校名称", "字段名称", "区"]) if self.district_filter: self._basic_info = df_basic[df_basic["区"] == self.district_filter].copy() else: self._basic_info = df_basic def _load_legacy_data(self) -> None: """从旧的区级parquet加载(兼容)""" self._basic_info = pd.read_parquet(PARQUET_BASIC_INFO) self._basic_info = self._strip_str_columns(self._basic_info, ["学校名称", "字段名称", "区"]) if self.district_filter: self._basic_info = self._basic_info[self._basic_info["区"] == self.district_filter].copy() self._basic_info = self._standardize_school_name(self._basic_info, "学校名称") df = pd.read_parquet(PARQUET_COURSE_IMPL) df = self._strip_str_columns(df, ["学校简称", "字段名称", "所在区"]) if self.district_filter: df = df[df["所在区"] == self.district_filter].copy() df = self._standardize_school_name(df, "学校简称") df = df.rename(columns={"学校简称": "学校名称"}) self._course_impl = df df = pd.read_parquet(PARQUET_SUBJECT_IMPL) df = self._strip_str_columns(df, ["学校简称", "字段名称", "所在区", "学科"]) if self.district_filter: df = df[df["所在区"] == self.district_filter].copy() df = self._standardize_school_name(df, "学校简称") df = df.rename(columns={"学校简称": "学校名称"}) self._subject_impl = df def _standardize_school_name(self, df: pd.DataFrame, col: str) -> pd.DataFrame: """将学校全称映射为简称(旧数据兼容用)""" df[col] = df[col].map(lambda x: SCHOOL_NAME_FULL_TO_SHORT.get(x, x)) return df @property def city_schools(self) -> List[str]: """获取全市学校列表(用于全市基准计算)""" if self._city_course_impl is not None: return sorted(self._city_course_impl["学校名称"].unique().tolist()) return self.schools @property def city_course_data(self) -> Optional[pd.DataFrame]: """全市课程实施数据(不受district_filter限制)""" return self._city_course_impl @property def city_subject_data(self) -> Optional[pd.DataFrame]: """全市学科课程数据(不受district_filter限制)""" return self._city_subject_impl # ===== 与一期 DataEngine 完全一致的接口 ===== @property def schools(self) -> List[str]: """获取所有学校名称列表""" if self._course_impl is None: self.load_all() return sorted(self._course_impl["学校名称"].unique().tolist()) def get_school_basic_info(self, school: str) -> Dict: """获取学校基本信息""" if self._basic_info is None: self.load_all() df = self._basic_info[self._basic_info["学校名称"] == school] def _get_field(field_name): rows = df[df["字段名称"] == field_name] if len(rows) > 0: return rows.iloc[0]["字段值"] return None return { "school_name": school, "建校年份": _get_field("建校年份"), "占地面积": _get_field("占地面积"), "建筑面积": _get_field("建筑面积"), } def get_school_course_data(self, school: str) -> pd.DataFrame: """获取某学校的课程实施情况数据""" if self._course_impl is None: self.load_all() return self._course_impl[self._course_impl["学校名称"] == school].copy() def get_school_subject_data(self, school: str, subject: Optional[str] = None) -> pd.DataFrame: """获取某学校的学科课程实施数据""" if self._subject_impl is None: self.load_all() df = self._subject_impl[self._subject_impl["学校名称"] == school].copy() if subject: df = df[df["学科"] == subject] return df # ===== 课程领导力相关数据提取(与一期一致) ===== def get_weekly_hours(self, school: str) -> pd.DataFrame: """获取学校各学科各年级各学期的周课时数据""" df = self.get_school_course_data(school) hours_fields = ["学科必修课周课时", "学科选择性必修课周课时", "学科类选修课周课时"] result = df[df["字段名称"].isin(hours_fields)].copy() result["字段取值"] = pd.to_numeric(result["字段取值"], errors="coerce") return result def get_course_norms(self, school: str) -> List[Dict]: """获取学校课程规范落实相关数据""" df = self.get_school_course_data(school) norm_keywords = ["建设规范文本", "档案", "已经建成并使用", "尚未建成"] mask = df["字段名称"].apply( lambda x: any(k in str(x) for k in norm_keywords) if pd.notna(x) else False ) return df[mask][["字段名称", "字段取值"]].to_dict("records") # ===== 教学变革力相关(与一期一致) ===== def get_teaching_reform_data(self, school: str) -> pd.DataFrame: df = self.get_school_subject_data(school) reform_keywords = [ "认识程度", "落实程度", "认识", "落实", "理解式学习", "自主性学习", "实践性学习", "跨学科学习", "信息技术与教学融合", "信息融入教学", ] mask = df["字段名称"].apply( lambda x: any(k in str(x) for k in reform_keywords) if pd.notna(x) else False ) return df[mask] def get_homework_data(self, school: str) -> pd.DataFrame: df = self.get_school_subject_data(school) hw_keywords = [ "作业", "实践类", "表现类", "跨学科", "团队合作", "批改", "评价", "属性标注", "时长控制", ] mask = df["字段名称"].apply( lambda x: any(k in str(x) for k in hw_keywords) if pd.notna(x) else False ) return df[mask] # ===== 通用方法 ===== def get_field_value(self, school: str, source: str, field_name: str, subject: Optional[str] = None) -> Optional[str]: if source == "basic": df = self._basic_info[self._basic_info["学校名称"] == school] col = "字段名称" val_col = "字段值" elif source == "course": df = self.get_school_course_data(school) col = "字段名称" val_col = "字段取值" elif source == "subject": df = self.get_school_subject_data(school, subject) col = "字段名称" val_col = "字段取值" else: return None rows = df[df[col] == field_name] if len(rows) > 0: return rows.iloc[0][val_col] return None def summary(self) -> Dict: if self._basic_info is None: self.load_all() return { "era": 2, "district_filter": self.district_filter, "schools": self.schools, "school_count": len(self.schools), "basic_info_rows": len(self._basic_info), "course_impl_rows": len(self._course_impl), "subject_impl_rows": len(self._subject_impl), "subjects": SUBJECTS, "districts": sorted(self._course_impl["所在区"].unique().tolist()) if "所在区" in self._course_impl.columns else [], }