Clear Spike Librenms

Panduan cara menghapus traffic spike atau data tidak valid pada grafik monitoring LibreNMS menggunakan perintah rrdtool agar visualisasi bandwidth kembali akurat

Clear Spike Librenms
Photo by Luke Chesser / Unsplash

Fungsi utama

Script ini akan membaca file .rrd, lalu untuk DS traffic seperti:

INOCTETS, OUTOCTETS

dia cari sample yang:

nilai sekarang > median sekitar × factor

dan hanya dihapus kalau spike itu pendek, sesuai --max-run.

Contoh:

200M, 210M, 190M, 1G, 205M, 215M

Dengan --factor 4, nilai 1G akan dianggap spike karena jauh lebih besar dari median sekitar.

Tapi kalau begini:

200M, 700M, 850M, 900M, 750M, 250M

Kalau --max-run 1, dia tidak akan hapus, karena itu terlihat seperti traffic tinggi yang sustain beberapa sample.


Cara pasang

Simpan script dari canvas ke:

/opt/librenms/scripts/rrd_remove_short_spikes.py

Lalu:

chmod +x /opt/librenms/scripts/rrd_remove_short_spikes.pychown librenms:librenms /opt/librenms/scripts/rrd_remove_short_spikes.py

Command awal yang paling aman

Dry run dulu:

python3 /opt/librenms/scripts/rrd_remove_short_spikes.py \  /opt/librenms/rrd/103.73.75.25/port-id317.rrd \  --factor 5 \  --window 7 \  --max-run 1 \  --dry-run \  --verbose

Kalau output-nya masuk akal, jalankan real:

python3 /opt/librenms/scripts/rrd_remove_short_spikes.py \  /opt/librenms/rrd/103.73.75.25/port-id317.rrd \  --factor 5 \  --window 7 \  --max-run 1 \  --verbose

Kalau user kamu sudah masuk group librenms dan permission RRD sudah drwxrwxr-x+, bisa tanpa sudo.


Penjelasan fitur

1. --factor

Default:

--factor 5

Artinya sample dianggap spike kalau:

sample > median sekitar × 5

Contoh:

  • traffic sekitar: 200M
  • factor 5
  • batas spike: 1G

Kalau sample tiba-tiba 1G, dia akan kena.

Rekomendasi:

  • --factor 3 = agresif
  • --factor 5 = aman untuk awal
  • --factor 8 atau 10 = konservatif

2. --window

Default:

--window 7

Artinya dia melihat 7 sample sekitar titik tersebut.

Dengan polling 5 menit:

  • window 5 = ±10 menit
  • window 7 = ±15 menit
  • window 9 = ±20 menit

Rekomendasi:

  • --window 5 untuk spike sangat tajam
  • --window 7 untuk default
  • --window 9 kalau traffic agak fluktuatif

Harus angka ganjil minimal 3.


3. --max-run

Default:

--max-run 1

Ini fitur paling penting.

Artinya hanya hapus spike yang terjadi maksimal 1 sample berturut-turut.

Dengan polling 5 menit:

  • --max-run 1 = hapus spike 5 menit
  • --max-run 2 = hapus spike sampai 10 menit
  • --max-run 3 = hapus spike sampai 15 menit

Rekomendasi:

  • pakai 1 untuk paling aman
  • pakai 2 kalau spike palsu sering muncul 2 titik
  • jangan terlalu besar, karena bisa menghapus traffic real

4. --replace

Default:

--replace nan

Pilihan:

--replace nan--replace median

Bedanya:

  • nan = titik spike dibuang dari graph
  • median = titik spike diganti nilai median sekitar

Rekomendasi:

  • pakai nan untuk data palsu
  • pakai median kalau ingin graph terlihat lebih “nyambung”

Contoh:

python3 /opt/librenms/scripts/rrd_remove_short_spikes.py \  /opt/librenms/rrd/103.73.75.25/port-id317.rrd \  --factor 5 \  --window 7 \  --max-run 1 \  --replace median \  --verbose

5. --min-value-bps

Default:

--min-value-bps 0

Ini untuk mencegah sample kecil ikut dianggap spike.

Contoh:
traffic normal rendah, 10Kbps tiba-tiba 100Kbps. Secara rasio 10x, tapi tidak penting.

Kalau kamu cuma mau bersihkan spike besar, set:

--min-value-bps 100M

Artinya hanya nilai di atas 100 Mbps yang boleh dipertimbangkan sebagai spike.

Contoh:

python3 /opt/librenms/scripts/rrd_remove_short_spikes.py \  /opt/librenms/rrd/103.73.75.25/port-id317.rrd \  --factor 4 \  --window 7 \  --max-run 1 \  --min-value-bps 100M \  --dry-run \  --verbose

Ini cocok untuk kasus:

  • normal 200M
  • spike 1G
  • jangan ganggu fluktuasi kecil

6. --min-baseline-bps

Default:

--min-baseline-bps 0

Ini memastikan baseline sekitar cukup besar dulu sebelum deteksi aktif.

Contoh:
kalau traffic normal 1Mbps, spike ke 20Mbps bisa dianggap 20x, padahal belum tentu penting.

Kalau kamu set:

--min-baseline-bps 100M

maka deteksi hanya aktif kalau median sekitar minimal 100 Mbps.

Untuk case normal 200M spike 1G, ini bagus:

python3 /opt/librenms/scripts/rrd_remove_short_spikes.py \  /opt/librenms/rrd/103.73.75.25/port-id317.rrd \  --factor 4 \  --window 7 \  --max-run 1 \  --min-baseline-bps 100M \  --dry-run \  --verbose

7. --ds

Default:

--ds INOCTETS,OUTOCTETS

Ini hanya proses traffic in/out.

Kalau hanya mau bersihkan inbound:

--ds INOCTETS

Kalau hanya outbound:

--ds OUTOCTETS

8. --cf

Default:

--cf AVERAGE,MAX

Ini penting untuk LibreNMS karena graph sering menggambar:

  • AVERAGE
  • MAX

Kalau spike masih muncul di One Month, biasanya karena MAX masih kotor.

Kalau hanya mau proses average:

--cf AVERAGE

Kalau mau fokus clear spike bulanan, tetap pakai default:

--cf AVERAGE,MAX

Script

#!/usr/bin/env python3
"""
Remove short-lived traffic spikes from LibreNMS/RRD port files.

This script is intentionally pattern-based, not port-speed-based.
It removes short bursts that are far above the local median around the sample.

Typical use:
  python3 rrd_remove_short_spikes.py /opt/librenms/rrd/HOST/port-id123.rrd --dry-run --verbose
  python3 rrd_remove_short_spikes.py /opt/librenms/rrd/HOST/port-id123.rrd --factor 5 --window 7 --max-run 1 --verbose
"""

import argparse
import math
import os
import shutil
import subprocess
import sys
import tempfile
import xml.etree.ElementTree as ET
from dataclasses import dataclass
from pathlib import Path
from statistics import median
from typing import Dict, List, Optional, Sequence, Tuple


DEFAULT_DS = "INOCTETS,OUTOCTETS"


@dataclass
class Sample:
    value: Optional[float]
    element: ET.Element


@dataclass
class KillCandidate:
    idx: int
    ds_name: str
    value: float
    baseline: float
    ratio: float


def run(cmd: List[str]) -> subprocess.CompletedProcess:
    return subprocess.run(cmd, check=True, text=True, capture_output=True)


def parse_args() -> argparse.Namespace:
    parser = argparse.ArgumentParser(
        description="Remove short traffic spikes from RRD files using local-median anomaly detection."
    )
    parser.add_argument("rrdfile", help="Path to target .rrd file")
    parser.add_argument(
        "--ds",
        default=DEFAULT_DS,
        help=f"Comma-separated DS names to process. Default: {DEFAULT_DS}",
    )
    parser.add_argument(
        "--factor",
        type=float,
        default=5.0,
        help="Spike threshold. A sample is suspicious if value > local_median * factor. Default: 5.0",
    )
    parser.add_argument(
        "--window",
        type=int,
        default=7,
        help="Odd-sized local window around each sample. Default: 7",
    )
    parser.add_argument(
        "--max-run",
        type=int,
        default=1,
        help="Maximum consecutive suspicious samples to remove. Longer runs are treated as real traffic. Default: 1",
    )
    parser.add_argument(
        "--min-value-bps",
        default="0",
        help="Minimum traffic value in bps before a sample can be removed. Example: 100M. Default: 0",
    )
    parser.add_argument(
        "--min-baseline-bps",
        default="0",
        help="Minimum local baseline in bps before detection is allowed. Helps avoid false positives on near-zero traffic. Default: 0",
    )
    parser.add_argument(
        "--replace",
        choices=["nan", "median"],
        default="nan",
        help="Replace removed spikes with NaN or local median. Default: nan",
    )
    parser.add_argument(
        "--cf",
        default="AVERAGE,MAX",
        help="Comma-separated RRA consolidation functions to process. Default: AVERAGE,MAX",
    )
    parser.add_argument(
        "--dry-run",
        action="store_true",
        help="Analyze and report only. Do not write changes.",
    )
    parser.add_argument(
        "--no-backup",
        action="store_true",
        help="Do not create .bak backup before restore.",
    )
    parser.add_argument(
        "--verbose",
        action="store_true",
        help="Print detailed changes.",
    )
    return parser.parse_args()


def parse_bps(value: str) -> float:
    raw = str(value).strip().lower().replace(" ", "")
    if raw in {"", "0"}:
        return 0.0
    units = {
        "": 1.0,
        "b": 1.0,
        "bps": 1.0,
        "k": 1_000.0,
        "kb": 1_000.0,
        "kbps": 1_000.0,
        "m": 1_000_000.0,
        "mb": 1_000_000.0,
        "mbps": 1_000_000.0,
        "g": 1_000_000_000.0,
        "gb": 1_000_000_000.0,
        "gbps": 1_000_000_000.0,
        "t": 1_000_000_000_000.0,
        "tb": 1_000_000_000_000.0,
        "tbps": 1_000_000_000_000.0,
    }
    import re

    match = re.fullmatch(r"([0-9]+(?:\.[0-9]+)?)([a-z]*)", raw)
    if not match or match.group(2) not in units:
        raise ValueError(f"Invalid bps value: {value}")
    return float(match.group(1)) * units[match.group(2)]


def dump_rrd(rrd_file: Path, xml_file: Path) -> None:
    cp = run(["rrdtool", "dump", str(rrd_file)])
    xml_file.write_text(cp.stdout, encoding="utf-8")


def restore_rrd(xml_file: Path, rrd_file: Path) -> None:
    tmp_out = rrd_file.with_suffix(rrd_file.suffix + ".new")
    run(["rrdtool", "restore", str(xml_file), str(tmp_out)])
    os.replace(tmp_out, rrd_file)


def get_ds_names(root: ET.Element) -> List[str]:
    names: List[str] = []
    for ds in root.findall("./ds"):
        name = (ds.findtext("name") or "").strip().upper()
        if name:
            names.append(name)
    return names


def parse_float_or_none(text: Optional[str]) -> Optional[float]:
    if text is None:
        return None
    raw = text.strip()
    if raw.lower() == "nan":
        return None
    try:
        val = float(raw)
    except ValueError:
        return None
    if math.isnan(val) or math.isinf(val):
        return None
    return val


def local_baseline(values: Sequence[Optional[float]], idx: int, window: int) -> Optional[float]:
    radius = window // 2
    left = max(0, idx - radius)
    right = min(len(values), idx + radius + 1)
    neighbors = [v for j, v in enumerate(values[left:right], start=left) if j != idx and v is not None and v > 0]
    if len(neighbors) < 2:
        return None
    return float(median(neighbors))


def suspicious_candidates(
    samples: List[Sample],
    ds_name: str,
    factor: float,
    window: int,
    min_value_octets: float,
    min_baseline_octets: float,
) -> List[KillCandidate]:
    values = [s.value for s in samples]
    candidates: List[KillCandidate] = []
    for idx, value in enumerate(values):
        if value is None or value <= 0:
            continue
        if value < min_value_octets:
            continue
        baseline = local_baseline(values, idx, window)
        if baseline is None or baseline <= 0:
            continue
        if baseline < min_baseline_octets:
            continue
        ratio = value / baseline
        if ratio > factor:
            candidates.append(KillCandidate(idx=idx, ds_name=ds_name, value=value, baseline=baseline, ratio=ratio))
    return candidates


def split_runs(candidates: List[KillCandidate]) -> List[List[KillCandidate]]:
    if not candidates:
        return []
    sorted_candidates = sorted(candidates, key=lambda c: c.idx)
    runs: List[List[KillCandidate]] = [[sorted_candidates[0]]]
    for cand in sorted_candidates[1:]:
        if cand.idx == runs[-1][-1].idx + 1:
            runs[-1].append(cand)
        else:
            runs.append([cand])
    return runs


def format_bps_from_octets(value: float) -> str:
    bps = value * 8.0
    units = [("Tbps", 1e12), ("Gbps", 1e9), ("Mbps", 1e6), ("Kbps", 1e3)]
    for suffix, factor in units:
        if abs(bps) >= factor:
            return f"{bps / factor:.2f} {suffix}"
    return f"{bps:.2f} bps"


def process_xml(
    xml_file: Path,
    selected_ds: set,
    selected_cf: set,
    factor: float,
    window: int,
    max_run: int,
    min_value_bps: float,
    min_baseline_bps: float,
    replacement: str,
    verbose: bool,
) -> Tuple[int, Dict[str, int]]:
    tree = ET.parse(xml_file)
    root = tree.getroot()
    ds_names = get_ds_names(root)
    if not ds_names:
        raise RuntimeError("No DS names found in RRD XML")

    min_value_octets = min_value_bps / 8.0
    min_baseline_octets = min_baseline_bps / 8.0
    total_changed = 0
    changed_by_key: Dict[str, int] = {}

    for rra_idx, rra in enumerate(root.findall("./rra")):
        cf = (rra.findtext("cf") or "").strip().upper()
        if cf not in selected_cf:
            continue
        database = rra.find("database")
        if database is None:
            continue
        rows = database.findall("row")
        if not rows:
            continue

        for ds_idx, ds_name in enumerate(ds_names):
            if ds_name not in selected_ds:
                continue
            samples: List[Sample] = []
            for row in rows:
                values = row.findall("v")
                if ds_idx >= len(values):
                    continue
                elem = values[ds_idx]
                samples.append(Sample(value=parse_float_or_none(elem.text), element=elem))

            candidates = suspicious_candidates(
                samples=samples,
                ds_name=ds_name,
                factor=factor,
                window=window,
                min_value_octets=min_value_octets,
                min_baseline_octets=min_baseline_octets,
            )
            runs = split_runs(candidates)
            for run in runs:
                if len(run) > max_run:
                    if verbose:
                        print(
                            f"Skip run: RRA#{rra_idx} CF={cf} DS={ds_name} len={len(run)} "
                            f"idx={run[0].idx}-{run[-1].idx}; treated as sustained traffic"
                        )
                    continue
                for cand in run:
                    sample = samples[cand.idx]
                    if replacement == "nan":
                        sample.element.text = " NaN "
                    else:
                        sample.element.text = f" {cand.baseline:.10e} "
                    total_changed += 1
                    key = f"{cf}:{ds_name}"
                    changed_by_key[key] = changed_by_key.get(key, 0) + 1
                    if verbose:
                        print(
                            f"Remove spike: RRA#{rra_idx} CF={cf} DS={ds_name} idx={cand.idx} "
                            f"value={format_bps_from_octets(cand.value)} baseline={format_bps_from_octets(cand.baseline)} "
                            f"ratio={cand.ratio:.2f}x"
                        )

    tree.write(xml_file, encoding="utf-8", xml_declaration=True)
    return total_changed, changed_by_key


def main() -> int:
    args = parse_args()
    rrd_file = Path(args.rrdfile)
    if args.window < 3 or args.window % 2 == 0:
        print("--window must be an odd integer >= 3", file=sys.stderr)
        return 2
    if args.factor <= 1:
        print("--factor must be > 1", file=sys.stderr)
        return 2
    if args.max_run < 1:
        print("--max-run must be >= 1", file=sys.stderr)
        return 2
    if not rrd_file.is_file():
        print(f"RRD file not found: {rrd_file}", file=sys.stderr)
        return 2

    selected_ds = {x.strip().upper() for x in args.ds.split(",") if x.strip()}
    selected_cf = {x.strip().upper() for x in args.cf.split(",") if x.strip()}
    min_value_bps = parse_bps(args.min_value_bps)
    min_baseline_bps = parse_bps(args.min_baseline_bps)

    with tempfile.TemporaryDirectory(prefix="rrdspike_") as tmpdir:
        xml_file = Path(tmpdir) / (rrd_file.stem + ".xml")
        dump_rrd(rrd_file, xml_file)

        if args.verbose:
            print(f"RRD file: {rrd_file}")
            print(f"DS selected: {', '.join(sorted(selected_ds))}")
            print(f"CF selected: {', '.join(sorted(selected_cf))}")
            print(f"Factor: {args.factor}x")
            print(f"Window: {args.window} samples")
            print(f"Max run: {args.max_run} samples")
            print(f"Minimum value: {args.min_value_bps}")
            print(f"Minimum baseline: {args.min_baseline_bps}")
            print(f"Replacement: {args.replace}")

        changed, changed_by_key = process_xml(
            xml_file=xml_file,
            selected_ds=selected_ds,
            selected_cf=selected_cf,
            factor=args.factor,
            window=args.window,
            max_run=args.max_run,
            min_value_bps=min_value_bps,
            min_baseline_bps=min_baseline_bps,
            replacement=args.replace,
            verbose=args.verbose,
        )

        print(f"Values changed: {changed}")
        for key, count in sorted(changed_by_key.items()):
            print(f"  {key}: {count}")

        if changed == 0:
            print("No short spikes matched the configured rules.")
            return 0

        if args.dry_run:
            print("Dry run only. No changes written.")
            return 0

        if not args.no_backup:
            backup = rrd_file.with_suffix(rrd_file.suffix + ".bak")
            shutil.copy2(rrd_file, backup)
            print(f"Backup created: {backup}")

        restore_rrd(xml_file, rrd_file)
        print("RRD restored successfully.")

    return 0


if __name__ == "__main__":
    sys.exit(main())