
凌晨 2 点,MySQL 主库宕机,Prometheus 告警触发。飞书收到了、钉钉没收到、企微干脆没配——第二天复盘时发现,团队 22 个人只有 12 个人看到告警。不是告警没发出来,是三套 IM 推送各自为政,有人只看飞书、有人只看钉钉,触达率只有 78%。 这就是我搭建统一 IM 网关的起因。
我负责的运维团队同时使用飞书、钉钉和企业微信三个 IM 平台:
AIOps 告警要同时推到三个平台,早期方案是每个平台各写一套推送逻辑:

问题 | 表现 |
|---|---|
维护成本高 | 三套代码独立维护,改一个告警模板要改三处 |
消息格式不统一 | 飞书用卡片、钉钉用 Markdown、企微用文本 |
会话记忆断裂 | 同一故障在飞书讨论了,钉钉群毫无感知 |
触达率低 | 有人没装某个 IM 客户端,告警就漏掉了 |
核心痛点:三平台各自推送,触达率仅 78%(有人只看一个 IM),且跨平台上下文完全断裂。

统一 IM 网关架构,单进程挂载飞书/钉钉/企微三个适配器,Redis 存储跨平台会话
特性 | 飞书 | 钉钉 | 企业微信 |
|---|---|---|---|
消息格式 | 交互卡片 | Markdown | Markdown/文本 |
Webhook 方式 | 自建应用 + 事件订阅 | 机器人 Webhook | 群机器人 Webhook |
交互能力 | 按钮、表单、选择器 | 按钮(ActionCard) | 仅文本链接 |
速率限制 | 5 条/秒/应用 | 20 条/分钟/群 | 20 条/分钟/群 |
文件上传 | 支持(需 v4 接口) | 支持(需 Media 接口) | 支持(需临时素材) |
鉴权方式 | App ID + App Secret | App Key + App Secret | Webhook Token |
推荐场景 | 交互式告警处理 | 标准告警通知 | 简单告警推送 |
为什么定义内部消息格式?屏蔽各平台差异,上层业务只需关心"发什么",不需要关心"怎么发"。
"""统一告警消息结构"""
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional
class Severity(Enum):
P0 = "P0" # 紧急:电话 + 全平台
P1 = "P1" # 严重:全平台推送
P2 = "P2" # 一般:按团队路由
P3 = "P3" # 提示:仅飞书
class Platform(Enum):
FEISHU = "feishu"
DINGTALK = "dingtalk"
WECOM = "wecom"
@dataclass
class AlertMessage:
"""统一告警消息"""
alert_id: str # 告警标识 ID
title: str # 告警标题
severity: Severity # 严重等级
summary: str # 摘要(一行描述)
detail: str # 详细描述
runbook_url: Optional[str] = None # Runbook 链接
dashboard_url: Optional[str] = None # 监控面板链接
labels: dict = field(default_factory=dict) # 自定义标签
targets: list[Platform] = field( # 目标平台
default_factory=lambda: [Platform.FEISHU, Platform.DINGTALK, Platform.WECOM]
)

图:Gateway 配置文件编辑界面,展示飞书、钉钉、企微三个平台的配置项与连接矩阵
为什么用 YAML 配置而不是代码硬编码?平台密钥、路由策略需要频繁调整,配置化让运维人员无需改代码即可调整。
# gateway-config.yaml
server:
host: "0.0.0.0"
port: 8080
workers: 4
redis:
host: "127.0.0.1"
port: 6379
db: 0
session_ttl: 86400 # 会话记忆保留 24 小时
platforms:
feishu:
app_id: "${FEISHU_APP_ID}"
app_secret: "${FEISHU_APP_SECRET}"
encrypt_key: "${FEISHU_ENCRYPT_KEY}"
verification_token: "${FEISHU_VERIFY_TOKEN}"
rate_limit: 5 # 5 条/秒
enabled: true
dingtalk:
webhook: "${DINGTALK_WEBHOOK}"
secret: "${DINGTALK_SECRET}"
keyword: "【AIOps】" # 安全关键词
rate_limit: 20 # 20 条/分钟
enabled: true
wecom:
webhook: "${WECOM_WEBHOOK}"
rate_limit: 20 # 20 条/分钟
enabled: true
routing:
P0:
platforms: [ feishu, dingtalk, wecom ]
mention_all: true
P1:
platforms: [ feishu, dingtalk, wecom ]
mention_all: false
P2:
platforms: [ feishu, dingtalk ]
P3:
platforms: [ feishu ]
session:
enabled: true
cross_platform_sync: true # 跨平台会话同步
context_key_template: "alert:{alert_id}:context"
为什么用调度器模式?告警等级不同,目标平台和推送策略不同,调度器统一管理分发逻辑。
"""网关分发调度器"""
import asyncio
import logging
from typing import Any
logger = logging.getLogger("gateway-dispatcher")
class GatewayDispatcher:
"""告警分发调度器"""
def __init__(self, config: dict, adapters: dict):
self.config = config
self.adapters = adapters # {platform: adapter_instance}
self.routing = config.get("routing", {})
async def dispatch(self, message: "AlertMessage") -> dict:
"""分发告警到目标平台"""
severity = message.severity.value
route = self.routing.get(severity, {})
target_platforms = route.get("platforms", ["feishu"])
mention_all = route.get("mention_all", False)
results = {}
tasks = []
for platform_name in target_platforms:
adapter = self.adapters.get(platform_name)
if adapter and adapter.enabled:
tasks.append(
self._safe_send(adapter, message, mention_all)
)
# 并发推送到各平台
send_results = await asyncio.gather(*tasks, return_exceptions=True)
for i, platform_name in enumerate(target_platforms):
result = send_results[i]
if isinstance(result, Exception):
logger.error(f"推送 {platform_name} 失败: {result}")
results[platform_name] = {"status": "failed", "error": str(result)}
else:
results[platform_name] = {"status": "ok", "msg_id": result}
# 保存会话上下文到 Redis
if self.config.get("session", {}).get("enabled"):
await self._save_context(message)
return results
async def _safe_send(self, adapter: Any, message: "AlertMessage",
mention_all: bool) -> str:
"""安全发送,带重试"""
for attempt in range(3):
try:
msg_id = await adapter.send(message, mention_all=mention_all)
return msg_id
except Exception as e:
if attempt < 2:
await asyncio.sleep(1 * (attempt + 1))
else:
raise
async def _save_context(self, message: "AlertMessage"):
"""保存跨平台会话上下文"""
import redis.asyncio as aioredis
import json
r = aioredis.from_url("redis://127.0.0.1:6379/0")
key = f"alert:{message.alert_id}:context"
context = {
"alert_id": message.alert_id,
"title": message.title,
"severity": message.severity.value,
"summary": message.summary,
"discussions": [], # 跨平台讨论记录
}
ttl = self.config.get("session", {}).get("session_ttl", 86400)
await r.setex(key, ttl, json.dumps(context, ensure_ascii=False))
await r.close()
为什么飞书用交互卡片而非纯文本?卡片支持按钮操作,值班人员可以直接在消息上"确认告警"或"触发自愈",减少上下文切换。
"""飞书平台适配器 - 交互卡片"""
import httpx
import logging
logger = logging.getLogger("adapter-feishu")
class FeishuAdapter:
"""飞书消息适配器"""
def __init__(self, config: dict):
self.app_id = config["app_id"]
self.app_secret = config["app_secret"]
self.enabled = config.get("enabled", True)
self._token = None
async def _get_token(self) -> str:
"""获取飞书 tenant_access_token"""
async with httpx.AsyncClient() as client:
resp = await client.post(
"https://open.feishu.cn/open-apis/auth/v3/tenant_access_token/internal",
json={"app_id": self.app_id, "app_secret": self.app_secret}
)
return resp.json()["tenant_access_token"]
async def send(self, message: "AlertMessage", **kwargs) -> str:
"""发送交互卡片消息"""
token = await self._get_token()
severity_emoji = {"P0": "🔴", "P1": "🟠", "P2": "🟡", "P3": "🟢"}
card = {
"elements": [
{"tag": "div", "text": {
"tag": "lark_md",
"content": f"{severity_emoji.get(message.severity.value, '⚪')} "
f"**{message.title}**\n{message.summary}"
}},
{"tag": "hr"},
{"tag": "div", "text": {
"tag": "lark_md",
"content": message.detail
}},
{"tag": "action", "actions": [
{"tag": "button", "text": {"tag": "plain_text", "content": "确认告警"},
"value": {"action": "ack", "alert_id": message.alert_id},
"type": "primary"},
{"tag": "button", "text": {"tag": "plain_text", "content": "触发自愈"},
"value": {"action": "heal", "alert_id": message.alert_id},
"type": "danger"}
]}
]
}
async with httpx.AsyncClient() as client:
resp = await client.post(
"https://open.feishu.cn/open-apis/im/v1/messages",
headers={"Authorization": f"Bearer {token}"},
params={"receive_id": "oc_xxx", "msg_type": "interactive"},
json={"content": str(card)}
)
data = resp.json()
return data.get("data", {}).get("message_id", "unknown")

图:飞书与钉钉跨平台续聊效果对比,同一 alert_id 在两端的消息自动双向同步

图:跨平台会话记忆同步流程,AIOps 网关 → Redis → 飞书/钉钉/企微三端同步
为什么需要跨平台会话同步?团队成员分散在不同 IM,如果各平台讨论互不可见,就会出现重复排查和沟通断层。
场景 | 优化前触达率 | 优化后触达率 | 提升幅度 |
|---|---|---|---|
P0 紧急告警 | 78%(单平台推送) | 99%(三平台 + @all) | ⬆️ 27% |
P1 严重告警 | 82%(双平台推送) | 97%(三平台推送) | ⬆️ 18% |
P2 一般告警 | 75%(仅飞书) | 93%(飞书 + 钉钉) | ⬆️ 24% |
跨平台协作 | 0%(无同步) | 85%(会话记忆同步) | ⬆️ 新能力 |
现象:钉钉群机器人配置了安全关键词,但网关推送的告警消息全部被静默丢弃,没有任何报错。
原因:钉钉群机器人的安全设置要求消息内容必须包含指定关键词,但网关组装的消息标题没有加关键词前缀,钉钉直接丢弃且不返回错误。
解决:在消息模板中统一添加 【AIOps】 前缀,确保触发安全关键词校验:
def format_dingtalk_title(message: "AlertMessage") -> str:
"""钉钉消息必须包含安全关键词"""
keyword = "【AIOps】"
if keyword not in message.title:
return f"{keyword}{message.title}"
return message.title
提醒:钉钉 Webhook 机器人的安全策略不会返回错误码,消息被拦截时完全静默,排查时一定要先检查安全设置。
现象:网关启动后前几分钟推送正常,运行一段时间后所有飞书推送开始报 401,重启后又恢复。
原因:飞书 tenant_access_token 有效期仅 2 小时,网关在启动时获取一次后就缓存使用,过期后没有刷新机制。
解决:添加 Token 自动刷新,在距离过期 5 分钟时主动续期:
import time
class FeishuAdapter:
def __init__(self, config):
self._token = None
self._token_expire_at = 0 # token 过期时间戳
async def _get_token(self) -> str:
"""带自动刷新的 Token 获取"""
if self._token and time.time() < self._token_expire_at - 300:
return self._token # 还有效,提前 5 分钟刷新
async with httpx.AsyncClient() as client:
resp = await client.post(
"https://open.feishu.cn/open-apis/auth/v3/tenant_access_token/internal",
json={"app_id": self.app_id, "app_secret": self.app_secret}
)
data = resp.json()
self._token = data["tenant_access_token"]
self._token_expire_at = time.time() + data.get("expire", 7200)
return self._token
提醒:所有 IM 平台的 Token 都有有效期,确保刷新逻辑在项目初期就规划好,否则上线后必出问题。
现象:告警风暴期间,企业微信群只收到前几条消息,后续消息全部丢失,但网关日志显示发送成功。
原因:企微 Webhook 限制 20 条/分钟,超出部分直接丢弃,HTTP 返回码仍是 200,但 errcode 字段非 0。网关只检查了 HTTP 状态码,没有检查响应体中的错误码。
解决:添加响应体错误码校验,并实现速率限制队列:
import asyncio
class WecomAdapter:
RATE_LIMIT = 20 # 20 条/分钟
_send_times: list[float] = []
async def send(self, message: "AlertMessage", **kwargs) -> str:
"""带速率限制的企微发送"""
await self._wait_rate_limit()
async with httpx.AsyncClient() as client:
resp = await client.post(self.webhook, json=payload)
data = resp.json()
if data.get("errcode", 0) != 0:
raise Exception(f"企微推送失败: {data}")
return data.get("msgid", "ok")
async def _wait_rate_limit(self):
"""速率限制等待"""
now = time.time()
self._send_times = [t for t in self._send_times if now - t < 60]
if len(self._send_times) >= self.RATE_LIMIT:
wait = 60 - (now - self._send_times[0]) + 1
await asyncio.sleep(wait)
self._send_times.append(time.time())
提醒:各 IM 平台的速率限制策略差异很大,企微和钉钉都是超限静默丢弃,必须自行实现速率控制。
跨平台 IM 网关的核心价值:
适用场景:多 IM 平台并存的运维团队、AIOps 告警分发、跨团队协作通知
💬 你们团队用哪个 IM 做运维通知?多平台推送踩过什么坑?评论区聊聊~
⭐️ 觉得有用?点个「在看」和「转发」,让更多运维兄弟告别告警漏推~
👇 扫码关注「行者架构谈」,每周五篇 AIOps 实战干货
📜 真实性声明 本文所有内容均基于作者在运维岗位的真实工作经验。触达率数据来自测试环境验证,技术细节保持完整和真实。为保护商业机密,部分敏感信息已做脱敏处理,但技术细节保持完整和真实。