diff --git a/ts_anomaly_td/__main__.py b/ts_anomaly_td/__main__.py new file mode 100644 index 0000000..abf5eaa --- /dev/null +++ b/ts_anomaly_td/__main__.py @@ -0,0 +1,5 @@ +"""python -m ts_anomaly_td 入口。""" + +from ts_anomaly_td.cli import main + +main() diff --git a/ts_anomaly_td/cli.py b/ts_anomaly_td/cli.py new file mode 100644 index 0000000..9848bbb --- /dev/null +++ b/ts_anomaly_td/cli.py @@ -0,0 +1,142 @@ +"""CLI argparse 入口:inject / detect-batch / forecast-anomaly / e2e。""" + +import argparse +import logging +import sys + +from ts_anomaly_td.connector import TDConnection +from ts_anomaly_td.schema import setup_schema +from ts_anomaly_td.io_csv import read_csv +from ts_anomaly_td.detection import detect_all_algos, ALL_ALGOS +from ts_anomaly_td.forecast import forecast_anomaly + +logger = logging.getLogger(__name__) + + +def cmd_inject(args): + """inject 子命令:读 CSV → 建表 → 批量插入。""" + logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") + + ts, vals, labels = read_csv(args.csv) + rows = list(zip(ts, vals, labels)) + logger.info("read %d rows from %s", len(rows), args.csv) + + conn = TDConnection(args.url) + try: + db = args.db or "ts_anomaly" + setup_schema(conn, args.stable, db) + conn.batch_insert(args.stable, rows) + logger.info("inject complete: %d rows → %s.s_%s", len(rows), db, args.stable) + finally: + conn.close() + + +def cmd_detect_batch(args): + """detect-batch 子命令:多算法 ANOMALY_WINDOW。""" + logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") + + algos = args.algos.split(",") if args.algos else ALL_ALGOS + + conn = TDConnection(args.url) + try: + db = args.db or "ts_anomaly" + conn.execute_no_result(f"USE {db}") + results = detect_all_algos(conn, args.stable, algos=algos) + for algo, r in results.items(): + if r["error"]: + print(f" {algo}: ERROR - {r['error']}") + else: + print(f" {algo}: {len(r['windows'])} windows") + finally: + conn.close() + + +def cmd_forecast_anomaly(args): + """forecast-anomaly 子命令:FORECAST + 区间外点。""" + logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") + + conn = TDConnection(args.url) + try: + db = args.db or "ts_anomaly" + conn.execute_no_result(f"USE {db}") + points = forecast_anomaly( + conn, args.stable, + algo=args.algo, + rows=args.rows, + conf=args.conf, + ) + anomalies = [p for p in points if p["is_anomaly"]] + print(f"FORECAST: {len(points)} points, {len(anomalies)} anomalies") + for p in anomalies: + print(f" ts={p['ts']} value={p['value']} _flow={p['_flow']} _fhigh={p['_fhigh']}") + finally: + conn.close() + + +def cmd_e2e(args): + """e2e 子命令:编排完整流程。""" + from ts_anomaly_td.e2e import run_e2e + return run_e2e(args) + + +def build_parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser( + prog="ts-anomaly-td", + description="TDengine 时序异常分析 CLI", + ) + sub = parser.add_subparsers(dest="command", required=True) + + # inject + p_inj = sub.add_parser("inject", help="注入 CSV 数据到 TDengine") + p_inj.add_argument("--csv", required=True, help="CSV 文件路径") + p_inj.add_argument("--stable", required=True, help="supertable 名称") + p_inj.add_argument("--db", default="ts_anomaly", help="数据库名 (default: ts_anomaly)") + p_inj.add_argument("--url", default="ws://root:taosdata@localhost:6041", help="TDengine WebSocket URL") + + # detect-batch + p_det = sub.add_parser("detect-batch", help="批量 ANOMALY_WINDOW 6 算法") + p_det.add_argument("--stable", required=True, help="supertable 名称") + p_det.add_argument("--db", default="ts_anomaly", help="数据库名") + p_det.add_argument("--algos", default=None, help="逗号分隔算法列表 (默认全 6 个)") + p_det.add_argument("--url", default="ws://root:taosdata@localhost:6041", help="TDengine WebSocket URL") + + # forecast-anomaly + p_fc = sub.add_parser("forecast-anomaly", help="FORECAST 预测式异常检测") + p_fc.add_argument("--stable", required=True, help="supertable 名称") + p_fc.add_argument("--db", default="ts_anomaly", help="数据库名") + p_fc.add_argument("--algo", default="holtwinters", help="预测算法 (default: holtwinters)") + p_fc.add_argument("--rows", type=int, default=10, help="预测行数 (default: 10)") + p_fc.add_argument("--conf", type=int, default=95, help="置信度 (default: 95)") + p_fc.add_argument("--url", default="ws://root:taosdata@localhost:6041", help="TDengine WebSocket URL") + + # e2e + p_e2e = sub.add_parser("e2e", help="端到端编排") + p_e2e.add_argument("--data-dir", default="data", help="CSV 数据目录 (default: data)") + p_e2e.add_argument("--output-dir", default="render", help="PNG 输出目录 (default: render)") + p_e2e.add_argument("--log-dir", default="logs", help="JSON 日志目录 (default: logs)") + p_e2e.add_argument("--url", default="ws://root:taosdata@localhost:6041", help="TDengine WebSocket URL") + + return parser + + +def main(): + parser = build_parser() + args = parser.parse_args() + + dispatch = { + "inject": cmd_inject, + "detect-batch": cmd_detect_batch, + "forecast-anomaly": cmd_forecast_anomaly, + "e2e": cmd_e2e, + } + + fn = dispatch.get(args.command) + if fn is None: + parser.print_help() + sys.exit(1) + + sys.exit(fn(args) or 0) + + +if __name__ == "__main__": + main() diff --git a/ts_anomaly_td/e2e.py b/ts_anomaly_td/e2e.py new file mode 100644 index 0000000..b534e7e --- /dev/null +++ b/ts_anomaly_td/e2e.py @@ -0,0 +1,3 @@ +"""E2E 编排器(占位 stub)""" +def run_e2e(args): + raise NotImplementedError("e2e 将在 Task 7 实现")