第313篇:Telemetry 订阅配置
关键词
Telemetry、订阅配置、流式数据、Push 模式、gRPC Dial-out、采样路径、On-Change、数据采集、华为 Telemetry
一、Telemetry 概述
1.1 什么是 Telemetry
Telemetry 是一种 设备主动推送数据 的监控机制,替代传统的 Polling(轮询)模式:
Polling vs Telemetry:
Polling(SNMP 轮询): | 管理站 每 5 分钟轮询 | ──get→ ←resp── | 网络设备 | | --- | --- | --- | 问题: ┌─ 5 分钟粒度,无法捕捉瞬时故障 ├─ 大量轮询浪费带宽(指标未变也查询) ├─ 设备负载高(响应轮询请求) └─ NMS 成为瓶颈(1000+ 设备难以扩展)
Telemetry(流式推送): | 采集系统 被动接收 | ←push── | 网络设备 主动上报 | | --- | --- | --- | 优势: ┌─ 秒/毫秒级粒度 ├─ 数据有变化才推送(On-Change) ├─ 设备主动推送,降低查询负载 └─ 水平扩展,支持大规模网络
1.2 Telemetry 的推送模式
Telemetry 两种推送模式:
1. Dial-out(设备主动连接采集器)
┌──────────────────────────────────────────┐
│ 设备 ──TCP/gRPC──→ 采集器 │
│ ├─ 设备作为客户端,主动连接采集器 │
│ ├─ 采集器作为服务端,监听端口 │
│ ├─ 适合防火墙后的网络 │
│ ├─ 无需 NAT 映射 │
│ └─ 华为/思科主流支持方式 │
│ │
│ 数据流向:设备 → 采集器 │
│ 连接发起:设备 → 采集器 │
└──────────────────────────────────────────┘
2. Dial-in(采集器主动拉取)
┌──────────────────────────────────────────┐
│ 采集器 ──gRPC──→ 设备 │
│ ├─ 采集器作为客户端,连接设备 │
│ ├─ 设备作为 gRPC 服务端 │
│ ├─ 使用 gNMI Subscribe RPC │
│ ├─ 适合公有云/集中管理场景 │
│ └─ 需设备可达 │
│ │
│ 数据流向:设备 → 采集器 │
│ 连接发起:采集器 → 设备 │
└──────────────────────────────────────────┘
1.3 订阅模式
Telemetry 订阅模式:
1. SAMPLE(定时采样)
┌─ 设备按固定间隔采样并推送
├─ 适用:计数器、利用率等连续指标
├─ 间隔:1s - 10min 可配
└─ 数据量可控
2. ON-CHANGE(变化触发)
┌─ 数据值变化时才推送
├─ 适用:状态变化(up/down)、配置变更
├─ 即时通知
└─ 数据量小
3. TARGET-DEFINED(设备默认)
┌─ 由设备决定推送策略
├─ 适用:无需精细控制的场景
└─ 较少使用
二、华为设备 Telemetry 配置
2.1 Dial-out 订阅配置
# 华为设备 Telemetry 配置(Dial-out 模式)
# ===== 1. 配置采集器 =====
system-view
telemetry
# 定义采集器组
destination-group COLLECTOR_GRP
ip-address type ipv4 10.0.0.100 # 采集器 IP
port 10031 # 采集器端口
protocol grpc # gRPC 协议
quit
# ===== 2. 定义传感器组(采集路径) =====
sensor-group INTERFACE_STATS
sensor-path huawei-ifm:ifm/interfaces/interface/statistics
quit
sensor-group SYSTEM_BASIC
sensor-path huawei-system:system/deviceInfo
sensor-path huawei-system:system/currentTime
quit
sensor-group BGP_STATE
sensor-path huawei-bgp:bgp/peers/peer/state
quit
# ===== 3. 订阅配置 =====
subscription INTERFACE_SUB
sensor-group INTERFACE_STATS
sample-interval 10000 # 采样间隔(毫秒)
destination-group COLLECTOR_GRP
quit
subscription SYSTEM_SUB
sensor-group SYSTEM_BASIC
sample-interval 60000 # 60 秒
destination-group COLLECTOR_GRP
quit
subscription BGP_SUB
sensor-group BGP_STATE
sample-interval 30000 # 30 秒
destination-group COLLECTOR_GRP
quit
# ===== 4. 提交 =====
commit
quit
# ===== 5. 验证 =====
display telemetry subscription
display telemetry sensor-group
display telemetry destination-group
2.2 采样路径说明
# 华为设备常用采样路径
# 接口统计
huawei-ifm:ifm/interfaces/interface/statistics
└─ 包含:in/out octets, packets, errors, discards
# CPU/内存
huawei-device:device/system/cpu-infos
huawei-device:device/system/memory-infos
# 环境信息
huawei-device:device/system/power-infos
huawei-device:device/system/fan-infos
huawei-device:device/system/temperature-infos
# BGP 状态
huawei-bgp:bgp/peers/peer/state
# OSPF 状态
huawei-ospf:ospf/ospfv2-instances/ospfv2-instance
# 接口状态
huawei-ifm:ifm/interfaces/interface/state
# ARP 表
huawei-arp:arp/arpEntries/arpEntry
# MAC 表
huawei-bridge:bridge/vlan/bridge-vlan
# 路由表
huawei-route:route/routeTable/routeEntry
# 查看设备支持的所有路径
display telemetry sensor-path
2.3 高级配置:On-Change 订阅
# On-Change 订阅(接口状态变化)
system-view
telemetry
sensor-group INTF_STATUS
sensor-path huaemon-ifm:ifm/interfaces/interface/state
quit
subscription INTF_STATUS_SUB
sensor-group INTF_STATUS
sample-interval 0 # 0 = On-Change 模式
destination-group COLLECTOR_GRP
quit
commit
# 注意:
# ┌─ On-Change 的 sample-interval 设为 0
# ├─ 接口 up/down 时立即推送
# └─ 减少无效数据量
三、Telemetry 采集系统搭建
3.1 采集架构
完整 Telemetry 采集架构:
| 设备 1 gRPC Dial-out | ────┐ |
|---|---|
| ┌────────────┐ ├──→ 设备 2 gRPC Dial-out | Telegraf |
| --- | --- |
| ┌────────────┐ 设备 N gRPC Dial-out └────────────┘ | ▼ ───┘ ┌──────────┐ 可视化 |
| --- | --- |
| └──────────┘ |
组件说明: ┌─ Telegraf:接收 gRPC Telemetry 数据 ├─ Kafka:缓冲削峰,保证数据不丢失 ├─ InfluxDB:时序数据存储 └─ Grafana:可视化仪表盘
3.2 Telegraf 采集配置
# telegraf.conf — Telemetry 采集器配置
# 全局配置
[agent]
interval = "10s"
flush_interval = "5s"
metric_batch_size = 5000
# ===== Telemetry 输入:gRPC Dial-out 服务端 =====
[[inputs.socket_listener]]
service_address = "tcp://0.0.0.0:10031"
data_format = "gpb"
# 或者使用 JSON
# data_format = "json"
# ===== Telemetry 输入:gNMI Dial-in 采集 =====
[[inputs.gnmi]]
addresses = [
"192.168.1.1:9339",
"192.168.1.2:9339",
]
username = "admin"
password = "admin123"
insecure = true
# 订阅路径
[[inputs.gnmi.subscription]]
path = "/interfaces/interface[name=*]/state/counters"
sample_interval = "10s"
[[inputs.gnmi.subscription]]
path = "/interfaces/interface[name=*]/state/oper-status"
sample_interval = "0s" # On-Change
[[inputs.gnmi.subscription]]
path = "/components/component/cpu"
sample_interval = "30s"
# ===== 输出到 InfluxDB =====
[[outputs.influxdb]]
urls = ["http://localhost:8086"]
database = "telemetry"
username = "admin"
password = "admin123"
retention_policy = "autogen"
# 标签
tagpass = {
host = ["*"]
}
# ===== 输出到 Kafka =====
[[outputs.kafka]]
brokers = ["localhost:9092"]
topic = "network-telemetry"
data_format = "json"
compression = "snappy"
# ===== 处理:聚合计算 =====
[[aggregators.minmax]]
period = "60s"
drop_original = false
[[aggregators.histogram]]
period = "300s"
buckets = [0.0, 10.0, 50.0, 100.0, 500.0]
3.3 Python 实现 Telemetry 采集器
#!/usr/bin/env python3
# telemetry_collector.py — 简单 Telemetry 采集器
import json
import time
import threading
from http.server import HTTPServer, BaseHTTPRequestHandler
from datetime import datetime
# ===== 内存存储 =====
telemetry_data = {}
data_lock = threading.Lock()
class TelemetryReceiver(BaseHTTPRequestHandler):
"""gRPC HTTP2 接收器(模拟)"""
def do_POST(self):
"""接收 Telemetry 推送数据"""
content_length = int(self.headers.get("Content-Length", 0))
body = self.rfile.read(content_length).decode("utf-8")
try:
data = json.loads(body)
self.process_telemetry(data)
self.send_response(200)
self.send_header("Content-Type", "application/json")
self.end_headers()
self.wfile.write(json.dumps({"status": "ok"}).encode())
except Exception as e:
self.send_response(400)
self.send_header("Content-Type", "application/json")
self.end_headers()
self.wfile.write(json.dumps({"error": str(e)}).encode())
def process_telemetry(self, data):
"""处理 Telemetry 数据"""
with data_lock:
device = data.get("device", "unknown")
timestamp = datetime.now().isoformat()
if device not in telemetry_data:
telemetry_data[device] = []
# 添加时间戳
data["received_at"] = timestamp
telemetry_data[device].append(data)
# 只保留最近 100 条
if len(telemetry_data[device]) > 100:
telemetry_data[device] = telemetry_data[device][-100:]
# 打印摘要
paths = data.get("paths", [])
print(f"[{timestamp}] {device}: {len(paths)} paths received")
def log_message(self, format, *args):
"""静默日志"""
pass
def start_collector(host="0.0.0.0", port=10031):
"""启动采集器"""
server = HTTPServer((host, port), TelemetryReceiver)
print(f"Telemetry 采集器启动: {host}:{port}")
try:
server.serve_forever()
except KeyboardInterrupt:
print("\n采集器停止")
server.server_close()
def get_latest_metrics(device=None, minutes=5):
"""获取最近指标"""
with data_lock:
if device:
data = telemetry_data.get(device, [])
else:
data = []
for d in telemetry_data.values():
data.extend(d)
# 按时间过滤
cutoff = datetime.now().timestamp() - minutes * 60
return [
item for item in data
if datetime.fromisoformat(item["received_at"]).timestamp() > cutoff
]
def analyze_telemetry():
"""分析 Telemetry 数据"""
while True:
time.sleep(5)
with data_lock:
for device, records in telemetry_data.items():
if not records:
continue
last = records[-1]
paths = last.get("paths", [])
# 分析接口状态
intf_status = [
p for p in paths
if "oper-status" in json.dumps(p)
]
if intf_status:
for status in intf_status:
print(f" [分析] {device}: {json.dumps(status)[:100]}")
if __name__ == "__main__":
# 启动分析线程
analyzer = threading.Thread(target=analyze_telemetry, daemon=True)
analyzer.start()
# 启动采集器
start_collector()
3.4 模拟设备推送
#!/usr/bin/env python3
# telemetry_simulator.py — 模拟 Telemetry 设备推送
import json
import time
import random
import requests
from datetime import datetime
# 采集器地址
COLLECTOR_URL = "http://10.0.0.100:10031"
# 模拟设备
DEVICES = ["CORE-RT01", "CORE-SW01", "ACC-SW01"]
def generate_interface_stats(device_name, if_name):
"""生成接口统计"""
return {
"interface": if_name,
"in-octets": random.randint(1000000, 9999999),
"out-octets": random.randint(500000, 5000000),
"in-packets": random.randint(10000, 99999),
"out-packets": random.randint(5000, 50000),
"in-errors": random.choice([0, 0, 0, 0, 1, 0]),
"out-errors": 0,
"in-discards": 0,
"out-discards": 0,
"utilization-in": round(random.uniform(10, 80), 1),
"utilization-out": round(random.uniform(5, 60), 1),
}
def generate_device_telemetry(device_name):
"""生成设备 Telemetry 数据"""
return {
"device": device_name,
"timestamp": datetime.now().isoformat(),
"paths": [
{
"path": "huawei-ifm:ifm/interfaces/interface/statistics",
"values": [
generate_interface_stats(device_name, "GE0/0/1"),
generate_interface_stats(device_name, "GE0/0/2"),
],
},
{
"path": "huawei-device:device/system/cpu-infos",
"values": {
"cpu-usage": round(random.uniform(10, 60), 1),
"cpu-usage-5sec": round(random.uniform(10, 60), 1),
"cpu-usage-1min": round(random.uniform(10, 50), 1),
},
},
{
"path": "huawei-device:device/system/memory-infos",
"values": {
"memory-usage": round(random.uniform(40, 80), 1),
"total-memory": 8192,
"used-memory": int(8192 * random.uniform(0.4, 0.8)),
},
},
],
}
def push_telemetry(device_name):
"""推送 Telemetry 数据"""
data = generate_device_telemetry(device_name)
try:
resp = requests.post(
COLLECTOR_URL,
json=data,
timeout=5,
)
if resp.status_code == 200:
print(f"✓ {device_name}: 数据推送成功")
else:
print(f"✗ {device_name}: 推送失败 {resp.status_code}")
except requests.exceptions.ConnectionError:
print(f"✗ {device_name}: 无法连接采集器 {COLLECTOR_URL}")
except Exception as e:
print(f"✗ {device_name}: 错误 {e}")
def run_simulation():
"""运行模拟"""
print("=== Telemetry 模拟器开始 ===")
print(f"目标采集器: {COLLECTOR_URL}")
cycle = 0
while True:
cycle += 1
print(f"\n--- 第 {cycle} 轮推送 ---")
for device in DEVICES:
push_telemetry(device)
time.sleep(0.5) # 避免同时推送
# 10 秒后再次推送
print("等待 10 秒...")
time.sleep(10)
if __name__ == "__main__":
run_simulation()
四、Telemetry 数据解析
4.1 GPB 编码解析
#!/usr/bin/env python3
# telemetry_parse.py — 解析 Telemetry 数据
import json
import struct
# 模拟 GPB(Google Protocol Buffers)编码数据
# 实际使用 protobuf 库解析
class TelemetryDecoder:
"""Telemetry 数据解码器"""
@staticmethod
def decode_gpb(data):
"""解码 GPB 格式(简化演示)"""
# 实际使用 protobuf 的 ParseFromString
# 这里演示结构
parsed = {
"header": {
"device": data.get("device", "unknown"),
"timestamp": data.get("timestamp", 0),
"encoding": "gpb",
},
"data": [],
}
return parsed
@staticmethod
def decode_json(data_str):
"""解码 JSON 格式"""
return json.loads(data_str)
@staticmethod
def extract_counters(telemetry_data):
"""提取计数器数据"""
counters = []
if "paths" not in telemetry_data:
return counters
for path_entry in telemetry_data["paths"]:
path = path_entry.get("path", "")
values = path_entry.get("values", [])
if "statistics" in path:
for val in values if isinstance(values, list) else [values]:
counters.append({
"path": path,
"interface": val.get("interface", "unknown"),
"in_octets": val.get("in-octets", 0),
"out_octets": val.get("out-octets", 0),
"in_errors": val.get("in-errors", 0),
"in_packets": val.get("in-packets", 0),
"out_packets": val.get("out-packets", 0),
"utilization_in": val.get("utilization-in", 0),
"utilization_out": val.get("utilization-out", 0),
})
return counters
@staticmethod
def delta_calculate(current, previous):
"""计算增量(用于计数器)"""
if not previous:
return current
delta = {}
for key in current:
if key in previous and isinstance(current[key], (int, float)):
if current[key] >= previous[key]:
delta[key] = current[key] - previous[key]
else:
# 计数器翻转
delta[key] = current[key]
else:
delta[key] = current[key]
return delta
# 使用示例
sample_data = {
"device": "CORE-RT01",
"timestamp": "2025-01-15T10:30:00Z",
"paths": [
{
"path": "huawei-ifm:ifm/interfaces/interface/statistics",
"values": [
{"interface": "GE0/0/1", "in-octets": 5000000, "out-octets": 2000000},
],
}
],
}
decoder = TelemetryDecoder()
counters = decoder.extract_counters(sample_data)
print("=== 提取的计数器 ===")
for c in counters:
print(f" {c['interface']}: in={c['in_octets']}, out={c['out_octets']}")
# 计算增量
prev = {"in_octets": 4900000, "out_octets": 1950000}
curr = {"in_octets": 5000000, "out_octets": 2000000}
delta = decoder.delta_calculate(curr, prev)
print(f"\n=== 增量计算 ===")
print(f" 入流量增量: {delta['in_octets']} bytes")
print(f" 出流量增量: {delta['out_octets']} bytes")
print(f" 速率: {delta['in_octets'] / 10} bytes/s(假设 10 秒间隔)")
4.2 数据入库与存储
#!/usr/bin/env python3
# telemetry_store.py — Telemetry 数据存储
import json
import sqlite3
from datetime import datetime
class TelemetryDB:
"""Telemetry 数据持久化"""
def __init__(self, db_path="telemetry.db"):
self.conn = sqlite3.connect(db_path)
self._create_tables()
def _create_tables(self):
"""建表"""
cursor = self.conn.cursor()
# 接口计数器表
cursor.execute("""
CREATE TABLE IF NOT EXISTS interface_counters (
id INTEGER PRIMARY KEY AUTOINCREMENT,
device TEXT,
interface TEXT,
timestamp TEXT,
in_octets INTEGER,
out_octets INTEGER,
in_errors INTEGER,
in_packets INTEGER,
out_packets INTEGER,
utilization_in REAL,
utilization_out REAL
)
""")
# CPU 表
cursor.execute("""
CREATE TABLE IF NOT EXISTS cpu_usage (
id INTEGER PRIMARY KEY AUTOINCREMENT,
device TEXT,
timestamp TEXT,
cpu_usage REAL,
cpu_5sec REAL,
cpu_1min REAL
)
""")
# 内存表
cursor.execute("""
CREATE TABLE IF NOT EXISTS memory_usage (
id INTEGER PRIMARY KEY AUTOINCREMENT,
device TEXT,
timestamp TEXT,
memory_usage REAL,
total_memory INTEGER,
used_memory INTEGER
)
""")
# 索引
cursor.execute("""
CREATE INDEX IF NOT EXISTS idx_counters_device
ON interface_counters(device, timestamp)
""")
self.conn.commit()
def save_telemetry(self, data):
"""保存 Telemetry 数据"""
cursor = self.conn.cursor()
device = data.get("device", "unknown")
timestamp = data.get("received_at", datetime.now().isoformat())
for path_entry in data.get("paths", []):
path = path_entry.get("path", "")
values = path_entry.get("values", [])
if not isinstance(values, list):
values = [values]
for val in values:
if "statistics" in path:
cursor.execute("""
INSERT INTO interface_counters
(device, interface, timestamp, in_octets, out_octets,
in_errors, in_packets, out_packets,
utilization_in, utilization_out)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""", (
device, val.get("interface"),
timestamp,
val.get("in-octets", 0),
val.get("out-octets", 0),
val.get("in-errors", 0),
val.get("in-packets", 0),
val.get("out-packets", 0),
val.get("utilization-in", 0),
val.get("utilization-out", 0),
))
elif "cpu" in path:
cursor.execute("""
INSERT INTO cpu_usage
(device, timestamp, cpu_usage, cpu_5sec, cpu_1min)
VALUES (?, ?, ?, ?, ?)
""", (
device, timestamp,
val.get("cpu-usage", 0),
val.get("cpu-usage-5sec", 0),
val.get("cpu-usage-1min", 0),
))
self.conn.commit()
def query_interface_metrics(self, device, start_time, end_time):
"""查询接口指标"""
cursor = self.conn.cursor()
cursor.execute("""
SELECT timestamp, interface, in_octets, out_octets,
utilization_in, utilization_out
FROM interface_counters
WHERE device = ? AND timestamp BETWEEN ? AND ?
ORDER BY timestamp
""", (device, start_time, end_time))
return cursor.fetchall()
def get_latest_cpu(self, device):
"""获取最新 CPU"""
cursor = self.conn.cursor()
cursor.execute("""
SELECT timestamp, cpu_usage, cpu_5sec, cpu_1min
FROM cpu_usage
WHERE device = ?
ORDER BY timestamp DESC
LIMIT 1
""", (device,))
return cursor.fetchone()
def close(self):
self.conn.close()
五、Telemetry 数据可视化
5.1 Grafana 仪表盘配置
{
"title": "Network Telemetry Dashboard",
"panels": [
{
"title": "接口利用率 TOP 10",
"type": "table",
"datasource": "InfluxDB",
"targets": [
{
"query": "SELECT last(\"utilization_in\") as \"in\", last(\"utilization_out\") as \"out\" FROM \"interface_counters\" WHERE $timeFilter GROUP BY \"device\", \"interface\""
}
]
},
{
"title": "设备 CPU 趋势",
"type": "timeseries",
"datasource": "InfluxDB",
"targets": [
{
"query": "SELECT mean(\"cpu_usage\") FROM \"cpu_usage\" WHERE $timeFilter GROUP BY \"device\""
}
]
},
{
"title": "接口错误率",
"type": "timeseries",
"targets": [
{
"query": "SELECT non_negative_derivative(mean(\"in_errors\")) FROM \"interface_counters\" WHERE $timeFilter GROUP BY \"device\", \"interface\""
}
]
}
]
}
六、Telemetry 最佳实践
6.1 配置规范
Telemetry 配置最佳实践:
1. 采样间隔
┌─ 接口计数器:10-30 秒(太频繁浪费带宽)
├─ CPU/内存:30-60 秒
├─ 环境状态:60-300 秒
├─ 接口状态(up/down):On-Change
└─ BGP 邻居状态:On-Change
2. 路径选择
┌─ 按需订阅,不要全量
├─ 精确路径减少传输量
├─ 使用 get 验证路径正确性
└─ 定期审查订阅路径
3. 采集器
┌─ 多个采集器做负载均衡
├─ Kafka 缓冲削峰
├─ 数据至少备份两份
└─ 采集器高可用
4. 数据管理
┌─ 时序数据库保留策略:7天原始、30天聚合
├─ 数据压缩(Snappy)
├─ 异常检测告警
└─ 定期清理过期数据
6.2 常见问题
Telemetry 常见问题排查:
设备未推送数据:
┌─ 检查 telemetry commit 是否提交
├─ display telemetry subscription 验证
├─ 检查采集器 IP:Port 可达性
└─ 防火墙是否放行端口
数据延迟:
┌─ 采样间隔是否过大
├─ 采集器是否负载过高
├─ Kafka 是否积压
└─ 网络带宽是否充足
部分路径无数据:
┌─ sensor-path 路径是否正确
├─ 设备是否支持该 YANG 模型
├─ 使用 display telemetry sensor-path 验证
└─ 尝试用 NETCONF get 该路径
采集器崩溃:
┌─ 增加内存/CPU
├─ 添加 Kafka 缓冲
├─ 数据限流
└─ 水平扩展采集器
七、总结
Telemetry 的核心价值:
从「你问我答」到「主动上报」
┌─ 实时性:秒/毫秒级数据
├─ 效率:变化才推送
├─ 扩展:水平扩展支持大规模
└─ 智能:数据驱动运维决策
数据链路:
设备采集 → 传输缓冲 → 解析存储 → 可视化告警
(Sensor) (gRPC) (Kafka) (InfluxDB + Grafana)
Telemetry + gNMI 是未来:
┌─ SNMP Polling → gNMI Subscribe
├─ 5 分钟 → 秒级
├─ 被动 → 主动
└─ 手动分析 → 智能告警
下篇预告:第314篇 — Python 自动化脚本集,将汇总网络自动化常用的 Python 脚本模板,涵盖配置备份、巡检、变更、合规检查等实战场景。