117 lines
3.8 KiB
Python
117 lines
3.8 KiB
Python
"""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)
|