main.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426
  1. import argparse
  2. import csv
  3. import json
  4. import logging
  5. import sys
  6. import time
  7. from datetime import datetime, timedelta
  8. from logging.handlers import TimedRotatingFileHandler
  9. from pathlib import Path
  10. import requests
  11. import yaml
  12. LOGIN_RETRIES = 3
  13. LOGIN_RETRY_INTERVAL_SECONDS = 60
  14. QUERY_TIMEOUT_RETRIES = 3
  15. QUERY_TIMEOUT_RETRY_INTERVAL_SECONDS = 10
  16. DEFAULT_LOG_FILE = "logs/run.log"
  17. DEFAULT_LOG_RETENTION_DAYS = 3
  18. DEFAULT_ADDPOINTDATUM_TIMEOUT_SECONDS = 5
  19. DEFAULT_CALCAGG_TIMEOUT_SECONDS = 60 * 50 # 50minute
  20. DEFAULT_TOKEN_REFRESH_INTERVAL_HOURS = 6
  21. LOGGER = logging.getLogger("jdf_energy_collector")
  22. def load_config(config_path: str) -> dict:
  23. with open(config_path, "r", encoding="utf-8") as file:
  24. config = yaml.safe_load(file) or {}
  25. required_sections = ["energy_chaowang", "points", "services", "schedule"]
  26. missing_sections = [section for section in required_sections if section not in config]
  27. if missing_sections:
  28. raise ValueError(f"配置文件缺少配置段: {', '.join(missing_sections)}")
  29. return config
  30. def setup_logging(config: dict) -> None:
  31. log_config = config.get("logging") or {}
  32. log_file = Path(log_config.get("file", DEFAULT_LOG_FILE))
  33. retention_days = int(log_config.get("retention_days", DEFAULT_LOG_RETENTION_DAYS))
  34. log_file.parent.mkdir(parents=True, exist_ok=True)
  35. formatter = logging.Formatter("%(asctime)s %(levelname)s %(message)s")
  36. file_handler = TimedRotatingFileHandler(
  37. log_file,
  38. when="midnight",
  39. interval=1,
  40. backupCount=retention_days,
  41. encoding="utf-8",
  42. )
  43. file_handler.setFormatter(formatter)
  44. console_handler = logging.StreamHandler(sys.stdout)
  45. console_handler.setFormatter(formatter)
  46. LOGGER.handlers.clear()
  47. LOGGER.setLevel(logging.INFO)
  48. LOGGER.addHandler(file_handler)
  49. LOGGER.addHandler(console_handler)
  50. LOGGER.propagate = False
  51. def mask_sensitive(value):
  52. if isinstance(value, dict):
  53. return {
  54. key: "***" if str(key).lower() in {"authorization", "password", "token"} else mask_sensitive(item)
  55. for key, item in value.items()
  56. }
  57. if isinstance(value, list):
  58. return [mask_sensitive(item) for item in value]
  59. return value
  60. def to_log_text(value) -> str:
  61. return json.dumps(mask_sensitive(value), ensure_ascii=False, default=str)
  62. def post_json(name: str, url: str, payload: dict, headers: dict | None = None, timeout: int | float | None = None):
  63. LOGGER.info(
  64. "HTTP请求开始: name=%s url=%s timeout=%s headers=%s payload=%s",
  65. name,
  66. url,
  67. timeout,
  68. to_log_text(headers or {}),
  69. to_log_text(payload),
  70. )
  71. started_at = time.monotonic()
  72. try:
  73. response = requests.post(url, headers=headers, json=payload, timeout=timeout)
  74. elapsed_ms = int((time.monotonic() - started_at) * 1000)
  75. try:
  76. response_body = response.json()
  77. except ValueError:
  78. response_body = response.text[:2000]
  79. LOGGER.info(
  80. "HTTP请求结束: name=%s status_code=%s elapsed_ms=%s response=%s",
  81. name,
  82. response.status_code,
  83. elapsed_ms,
  84. to_log_text(response_body),
  85. )
  86. return response
  87. except Exception as exc:
  88. elapsed_ms = int((time.monotonic() - started_at) * 1000)
  89. LOGGER.exception("HTTP请求异常: name=%s elapsed_ms=%s error=%s", name, elapsed_ms, exc)
  90. raise
  91. def read_points(points_file: str) -> tuple[list[str], dict[str, str]]:
  92. with open(points_file, encoding="utf-8-sig", newline="") as file:
  93. reader = csv.DictReader(file, delimiter=",")
  94. if reader.fieldnames is None:
  95. raise ValueError(f"点位文件为空: {points_file}")
  96. reader.fieldnames = [name.strip() if name else name for name in reader.fieldnames]
  97. if "id" not in reader.fieldnames or "point_id" not in reader.fieldnames:
  98. raise ValueError(f"{points_file} 必须包含 id 和 point_id 两列")
  99. tag_ids = []
  100. id_to_point_id = {}
  101. seen_tag_ids = set()
  102. seen_point_ids = set()
  103. duplicated_tag_ids = set()
  104. duplicated_point_ids = set()
  105. for row in reader:
  106. tag_id = (row.get("id") or "").strip()
  107. point_id = (row.get("point_id") or "").strip()
  108. if not tag_id or not point_id:
  109. continue
  110. if tag_id in seen_tag_ids:
  111. duplicated_tag_ids.add(tag_id)
  112. if point_id in seen_point_ids:
  113. duplicated_point_ids.add(point_id)
  114. seen_tag_ids.add(tag_id)
  115. seen_point_ids.add(point_id)
  116. tag_ids.append(tag_id)
  117. id_to_point_id[tag_id] = point_id
  118. if duplicated_tag_ids:
  119. raise ValueError(f"{points_file} 的 id 列存在重复值: {sorted(duplicated_tag_ids)}")
  120. if duplicated_point_ids:
  121. raise ValueError(f"{points_file} 的 point_id 列存在重复值: {sorted(duplicated_point_ids)}")
  122. if not tag_ids:
  123. raise ValueError(f"点位文件没有有效点位: {points_file}")
  124. LOGGER.info("点位文件读取完成: file=%s count=%s", points_file, len(tag_ids))
  125. return tag_ids, id_to_point_id
  126. def load_points_from_config(config: dict) -> tuple[list[str], dict[str, str]]:
  127. return read_points(str(Path(config["points"]["file"])))
  128. def calculate_hour_range(now: datetime | None = None) -> tuple[datetime, datetime, datetime]:
  129. current = now or datetime.now()
  130. current_hour = current.replace(minute=0, second=0, microsecond=0)
  131. start_time = current_hour - timedelta(hours=1)
  132. end_time = current_hour - timedelta(seconds=1)
  133. return start_time, end_time, current_hour
  134. def format_api_time(value: datetime) -> str:
  135. return value.strftime("%Y-%m-%d %H:%M:%S")
  136. def local_timestamp(value: datetime) -> int:
  137. return int(value.timestamp())
  138. def parse_full_time(full_time: str) -> datetime:
  139. return datetime.strptime(full_time, "%Y-%m-%d %H:%M:%S")
  140. def login(energy_config: dict) -> str:
  141. login_url = f"{energy_config['base_url'].rstrip('/')}/api/thingshome-admin/sys/service/getToken"
  142. payload = {
  143. "username": energy_config["username"],
  144. "password": energy_config["password"],
  145. }
  146. timeout = energy_config.get("login_timeout_seconds", 15)
  147. for attempt in range(1, LOGIN_RETRIES + 1):
  148. try:
  149. response = post_json("energy_chaowang_login", login_url, payload, timeout=timeout)
  150. response.raise_for_status()
  151. result = response.json()
  152. if result.get("code") == 0 and result.get("token"):
  153. return result["token"]
  154. LOGGER.warning("登录失败,第 %s 次: %s", attempt, to_log_text(result))
  155. except Exception as exc:
  156. LOGGER.warning("登录异常,第 %s 次: %s", attempt, exc)
  157. if attempt < LOGIN_RETRIES:
  158. time.sleep(LOGIN_RETRY_INTERVAL_SECONDS)
  159. raise RuntimeError("登录失败,已达到最大重试次数")
  160. class TokenManager:
  161. def __init__(self, energy_config: dict):
  162. self.energy_config = energy_config
  163. self.token = None
  164. self.last_login_at = None
  165. self.refresh_interval = timedelta(
  166. hours=float(energy_config.get("token_refresh_interval_hours", DEFAULT_TOKEN_REFRESH_INTERVAL_HOURS))
  167. )
  168. def get_token(self) -> str:
  169. now = datetime.now()
  170. if self.token and self.last_login_at and now < self.last_login_at + self.refresh_interval:
  171. LOGGER.info("复用登录 Token,下次登录时间: %s", format_api_time(self.last_login_at + self.refresh_interval))
  172. return self.token
  173. return self.refresh_token()
  174. def refresh_token(self) -> str:
  175. LOGGER.info("开始获取登录 Token")
  176. self.token = login(self.energy_config)
  177. self.last_login_at = datetime.now()
  178. LOGGER.info("登录 Token 获取成功,下次登录时间: %s", format_api_time(self.last_login_at + self.refresh_interval))
  179. return self.token
  180. def query_hour_data(energy_config: dict, token: str, tag_ids: list[str], start_time: datetime, end_time: datetime) -> dict:
  181. query_url = f"{energy_config['base_url'].rstrip('/')}/api/thingshome-ems/realTime/selectByCustom"
  182. headers = {
  183. "Content-Type": "application/json",
  184. "token": token,
  185. }
  186. payload = {
  187. "type": "1",
  188. "tagNameList": tag_ids,
  189. "energyTypeList": ["electric power"],
  190. "tb": "0",
  191. "startTime": format_api_time(start_time),
  192. "endTime": format_api_time(end_time),
  193. "dateCode": "h",
  194. }
  195. timeout = energy_config.get("query_timeout_seconds", 30)
  196. for attempt in range(1, QUERY_TIMEOUT_RETRIES + 1):
  197. try:
  198. response = post_json("energy_chaowang_query_hour", query_url, payload, headers=headers, timeout=timeout)
  199. response.raise_for_status()
  200. result = response.json()
  201. if result.get("code") != 0:
  202. raise RuntimeError(f"查询接口返回失败: {result}")
  203. return result
  204. except requests.Timeout as exc:
  205. LOGGER.warning("查询接口 timeout,第 %s 次: %s", attempt, exc)
  206. if attempt < QUERY_TIMEOUT_RETRIES:
  207. time.sleep(QUERY_TIMEOUT_RETRY_INTERVAL_SECONDS)
  208. raise RuntimeError("查询接口 timeout,已达到最大重试次数")
  209. def query_hour_data_with_token_refresh(
  210. energy_config: dict,
  211. token_manager: TokenManager,
  212. tag_ids: list[str],
  213. start_time: datetime,
  214. end_time: datetime,
  215. ) -> dict:
  216. token = token_manager.get_token()
  217. try:
  218. return query_hour_data(energy_config, token, tag_ids, start_time, end_time)
  219. except requests.HTTPError as exc:
  220. if exc.response is None or exc.response.status_code != 401:
  221. raise
  222. LOGGER.warning("查询接口返回 401,重新登录获取 Token 后重试")
  223. token = token_manager.refresh_token()
  224. return query_hour_data(energy_config, token, tag_ids, start_time, end_time)
  225. def extract_records(result: dict) -> list[dict]:
  226. records = []
  227. data_root = result.get("data") or {}
  228. for point_data in data_root.get("data") or []:
  229. records.extend(point_data.get("dataList") or [])
  230. if records:
  231. return records
  232. return data_root.get("trend") or []
  233. def call_addpointdatum(basedataportal_base_url: str, records: list[dict], id_to_point_id: dict[str, str]) -> list[str]:
  234. point_ids = []
  235. point_id_set = set()
  236. data_by_point_id = {}
  237. for record in records:
  238. tag_id = str(record.get("id", "")).strip()
  239. point_id = id_to_point_id.get(tag_id)
  240. if not point_id:
  241. LOGGER.warning("跳过未在 CSV 中配置的点位: %s", tag_id)
  242. continue
  243. name = record.get("name")
  244. full_time = record.get("fullTime")
  245. this_value = record.get("thisValue")
  246. if full_time is None or this_value is None:
  247. LOGGER.warning("跳过字段不完整的点位: %s", to_log_text(record))
  248. continue
  249. #if this_value == 0:
  250. #LOGGER.info("跳过 thisValue 为 0 的点位: %s", to_log_text(record))
  251. #continue
  252. full_time_ts = local_timestamp(parse_full_time(full_time))
  253. LOGGER.info("点位数据: name=%s fullTime=%s thisValue=%s", name, full_time, this_value)
  254. data_by_point_id.setdefault(point_id, []).append({"ts": full_time_ts, "value": str(this_value)})
  255. if point_id not in point_id_set:
  256. point_ids.append(point_id)
  257. point_id_set.add(point_id)
  258. if not point_ids:
  259. return []
  260. url = f"{basedataportal_base_url.rstrip('/')}/ai/addpointdatum"
  261. payload = [{"point_id": point_id, "data": data_by_point_id[point_id]} for point_id in point_ids]
  262. response = post_json(
  263. "basedataportal_addpointdatum",
  264. url,
  265. payload,
  266. timeout=DEFAULT_ADDPOINTDATUM_TIMEOUT_SECONDS,
  267. )
  268. response.raise_for_status()
  269. result = response.json()
  270. if result.get("state") != 0:
  271. raise RuntimeError(f"addpointdatum 接口返回异常: {to_log_text(result)}")
  272. LOGGER.info("addpointdatum 写入完成: count=%s", len(point_ids))
  273. return point_ids
  274. def call_calcagg(calcagg_base_url: str, point_ids: list[str], begin: int, end: int) -> None:
  275. url = f"{calcagg_base_url.rstrip('/')}/api/calcagg/calc_agg_points_range"
  276. payload = {
  277. "end": end,
  278. "begin": begin,
  279. "sync_run": True,
  280. "point_ids": point_ids,
  281. "operator_name": "jdf-energy-collector",
  282. }
  283. response = post_json("calcagg_calc_agg_points_range", url, payload, timeout=DEFAULT_CALCAGG_TIMEOUT_SECONDS)
  284. response.raise_for_status()
  285. result = response.json()
  286. if result.get("state") != 0:
  287. LOGGER.warning("计算聚合接口返回异常: %s", to_log_text(result))
  288. def run_once(config: dict, token_manager: TokenManager, tag_ids: list[str], id_to_point_id: dict[str, str]) -> None:
  289. start_time, end_time, current_hour = calculate_hour_range()
  290. LOGGER.info(
  291. "开始执行小时数据采集: %s - %s, ts:%s-%s",
  292. format_api_time(start_time),
  293. format_api_time(end_time),
  294. start_time,
  295. end_time,
  296. )
  297. result = query_hour_data_with_token_refresh(
  298. config["energy_chaowang"], token_manager, tag_ids, start_time, end_time
  299. )
  300. records = extract_records(result)
  301. written_point_ids = call_addpointdatum(config["services"]["basedataportal"], records, id_to_point_id)
  302. if not written_point_ids:
  303. LOGGER.info("没有成功写入 addpointdatum 的点位,跳过计算聚合接口调用")
  304. return
  305. begin = local_timestamp(start_time)
  306. end = local_timestamp(end_time)
  307. call_calcagg(config["services"]["calcagg"], written_point_ids, begin, end)
  308. LOGGER.info("本轮采集完成,写入点位数量: %s", len(written_point_ids))
  309. def wait_until_next_run(schedule_minute: int) -> None:
  310. now = datetime.now()
  311. next_run = now.replace(minute=schedule_minute, second=0, microsecond=0)
  312. if now > next_run:
  313. next_run += timedelta(hours=1)
  314. sleep_seconds = max(0, (next_run - now).total_seconds())
  315. LOGGER.info("下次执行时间: %s", format_api_time(next_run))
  316. time.sleep(sleep_seconds)
  317. def main() -> None:
  318. parser = argparse.ArgumentParser(description="小时能源数据采集入库程序")
  319. parser.add_argument("--config", default="config.yaml", help="YAML 配置文件路径")
  320. args = parser.parse_args()
  321. config = load_config(args.config)
  322. setup_logging(config)
  323. tag_ids, id_to_point_id = load_points_from_config(config)
  324. token_manager = TokenManager(config["energy_chaowang"])
  325. schedule_minute = int(config["schedule"].get("minute", 30))
  326. if schedule_minute < 0 or schedule_minute > 59:
  327. raise ValueError("schedule.minute 必须在 0 到 59 之间")
  328. while True:
  329. wait_until_next_run(schedule_minute)
  330. try:
  331. run_once(config, token_manager, tag_ids, id_to_point_id)
  332. except Exception as exc:
  333. LOGGER.exception("本轮采集失败: %s", exc)
  334. if __name__ == "__main__":
  335. main()