"""TMQ 准实时订阅模块:订阅 supertable topic + 周期 ANOMALY_WINDOW。""" import logging import signal import time logger = logging.getLogger(__name__) class WatchRunner: """TMQ watch 运行器(带轮询降级)。 Args: conn: TDConnection 实例 stable: supertable 名称 window_sec: ANOMALY_WINDOW 触发间隔(秒),默认 30 window_size: 滑动窗口大小(条),默认 1000 """ def __init__(self, conn, stable: str, window_sec: int = 30, window_size: int = 1000): self.conn = conn self.stable = stable self.window_sec = window_sec self.window_size = window_size self._stop = False self._mode = "unknown" def stop(self): """优雅停止。""" self._stop = True def run(self): """主循环。先尝试 TMQ,失败则降级为轮询。""" signal.signal(signal.SIGINT, self._handle_signal) signal.signal(signal.SIGTERM, self._handle_signal) try: self._run_tmq() except Exception as e: logger.warning("TMQ unavailable (%s), falling back to polling", e) print("watch: using polling mode (TMQ unavailable)") self._run_polling() def _handle_signal(self, signum, frame): logger.info("received signal %d, shutting down", signum) self.stop() def _run_tmq(self): """TMQ 模式:订阅 supertable topic → 滑动窗口 → 周期 ANOMALY_WINDOW。""" self._mode = "tmq" import taosws consumer = taosws.Consumer() consumer.subscribe([f"ds_{self.stable}"]) buffer = [] last_check = time.monotonic() while not self._stop: try: records = consumer.consume(timeout=2.0) for record in records: buffer.append((record.ts, record.value)) if len(buffer) > self.window_size: buffer = buffer[-self.window_size:] now = time.monotonic() if now - last_check >= self.window_sec and buffer: self._check_window(buffer) last_check = now except Exception as e: logger.error("TMQ consume error: %s", e) time.sleep(1) consumer.unsubscribe() consumer.close() logger.info("watch stopped (TMQ mode)") def _run_polling(self): """轮询模式:周期 SELECT + ANOMALY_WINDOW。""" self._mode = "polling" while not self._stop: try: sql = ( f"SELECT ts, value FROM s_{self.stable} " f"ORDER BY ts DESC LIMIT {self.window_size}" ) rows = self.conn.execute(sql) buffer = [(int(r[0]), float(r[1])) for r in rows] if buffer: self._check_window(buffer) except Exception as e: logger.error("polling error: %s", e) time.sleep(self.window_sec) logger.info("watch stopped (polling mode)") def _check_window(self, buffer: list): """对缓冲数据跑 ANOMALY_WINDOW + ANSI 输出。""" try: from ts_anomaly_td.detection import detect_all_algos results = detect_all_algos(self.conn, self.stable) for algo, r in results.items(): if r["error"]: continue for wstart, wend in r["windows"]: print( f"\033[33m[ANOMALY]\033[0m " f"\033[36m{algo}\033[0m " f"{wstart}-{wend}" ) except Exception as e: logger.error("check_window failed: %s", e)