
1. 项目概述为什么需要Arkime的实时通知如果你用过Arkime前身是Moloch进行全流量包捕获和分析你肯定遇到过这样的场景你设置了一个复杂的搜索规则用来监控某个特定IP的异常行为或者追踪某个敏感数据包的流向。然后呢你只能一遍遍地手动刷新页面或者设置一个定时任务去轮询查询结果。这不仅效率低下更关键的是你可能会错过那些需要立即响应的“黄金时间窗口”。一个攻击告警晚几分钟看到后果可能完全不同。这就是我决定折腾Arkime实时通知系统的原因。Arkime本身是一个强大的工具但其Web界面的交互本质上是“拉取”模式缺乏主动“推送”能力。我们需要一种机制让Arkime在捕获到符合特定条件的数据包、会话结束时或者系统状态发生变化时能主动、实时地通知我们。想象一下当有可疑的SSH爆破尝试、内部服务器对外发起异常连接或者某个关键服务的流量突然归零时你的手机、Slack频道或者内部告警平台能立刻收到一条消息——这才是真正意义上的态势感知。我选择的方案是WebSocket。相比传统的轮询Polling或长轮询Long PollingWebSocket提供了全双工、低延迟的通信通道。一旦连接建立服务器可以随时向客户端推送消息客户端也可以随时发送请求非常适合Arkime这种事件驱动、需要实时反馈的场景。整个配置过程从零开始到收到第一条实时通知我花了不到5分钟。下面我就把这套经过实战验证的WebSocket事件推送配置机制拆解给你看。2. 核心思路与架构设计在动手之前我们先理清整个通知系统的数据流和组件角色。Arkime本身并不直接内置一个功能齐全的“通知中心”它的强大之处在于其开放的API和灵活的架构。我们的目标是在不修改Arkime核心代码的前提下为其增加一个实时事件推送层。2.1 整体架构拆解整个系统可以看作由三个核心部分组成事件源Arkime这是数据的生产者。Arkime在捕获包、解析会话、更新数据库时会产生各种事件。我们需要一个“钩子”来捕获这些事件。事件处理与推送中间件我们的WebSocket服务这是系统的中枢。它负责监听Arkime产生的事件根据预定义的规则进行过滤、格式化然后通过WebSocket连接推送给所有在线的客户端。事件消费者客户端这可以是任何能连接WebSocket并解析JSON消息的程序。比如一个简单的HTMLJavaScript网页用于在监控大屏上实时显示告警。一个Python脚本接收到事件后调用企业微信、钉钉或Slack的机器人API发送消息。一个Go语言的后台服务将事件存入时序数据库如InfluxDB用于后续分析。那么如何从Arkime“钩”出事件呢这里有几个备选方案监听Arkime的日志文件Arkime会输出详细的操作日志。我们可以用tail -f或者类似logstash的工具去解析日志但这种方式比较“脏”日志格式可能变化且解析复杂事件如一个完整会话的元数据很困难。利用Arkime的APIArkime提供了丰富的RESTful API。我们可以定时轮询API来获取新数据但这又回到了老路上不是真正的“推送”。监听数据库变化推荐Arkime将所有捕获的会话Sessions和包Packets的元数据存储在Elasticsearch或OpenSearch中。当新数据写入时数据库本身会产生“变化”。我们可以监听这个变化。我选择的是第三种方案因为它最直接、最可靠。具体来说我使用了一个轻量级的工具Elasticsearch的“滚动查询”配合一个简单的Node.js WebSocket服务。为什么不直接用Elasticsearch的官方“Alerting”功能或Watcher因为我们需要的是高度定制化、低延迟的推送并且希望通知逻辑与我们的其他系统如前端展示、自定义机器人紧密集成一个自建的WebSocket服务提供了最大的灵活性。2.2 技术选型与工具清单基于以上思路我选用了以下工具栈它们都是轻量级、易部署的代表后端服务WebSocket ServerNode.js ws库。Node.js的异步非阻塞特性非常适合处理大量并发的WebSocket连接ws库是Node.js领域最成熟、性能最好的WebSocket实现。数据库监听使用Elasticsearch的“滚动查询”Scroll API或更现代的“Point in Time”PITAPI进行增量数据查询。我们不会使用消耗较大的持续监听如_changesAPIElasticsearch本身不直接提供类似功能而是采用一种“智能轮询”的方式。Arkime环境假设你已经有一个正常运行的Arkime集群包括Capture节点和Viewer节点并且数据后端是Elasticsearch/OpenSearch。这个方案的优势在于无侵入性。我们不需要修改Arkime的任何配置或代码只需要它有写入数据库的权限和我们有读取数据库的权限即可。整个WebSocket服务是独立部署的即使它宕机也不会影响Arkime核心的抓包和分析功能。3. 五分钟极速配置实战接下来我们进入实操环节。请确保你有一个可以安装Node.js的环境Linux/Mac/Windows WSL均可。3.1 第一步初始化项目与安装依赖1分钟打开终端创建一个新的项目目录并进入。mkdir arkime-websocket-notifier cd arkime-websocket-notifier npm init -y然后安装我们需要的核心依赖ws用于创建WebSocket服务器axios用于向Arkime/Elasticsearch发送HTTP请求dotenv用于管理环境变量。npm install ws axios dotenv3.2 第二步创建WebSocket服务器核心文件2分钟在项目根目录下创建两个文件.env和server.js。首先编辑.env文件配置你的环境变量。请务必根据你的实际环境修改这些值。# Arkime Viewer API 地址用于获取最新会话 ARKIME_API_BASEhttp://your-arkime-viewer-ip:8005 # Elasticsearch 地址用于直接查询如果API不满足需求 ES_HOSThttp://your-es-ip:9200 # WebSocket 服务器监听的端口 WS_PORT8080 # 轮询间隔毫秒例如5000表示每5秒检查一次新数据 POLL_INTERVAL5000 # 要监控的Arkime数据库索引模式通常为 sessions3-* ES_INDEX_PATTERNsessions3-*注意这里提供了两种数据源选择。Arkime Viewer的API更友好但可能无法满足非常复杂的过滤条件。Elasticsearch直接查询更强大但需要你熟悉其查询语法。我们后续以Arkime API为例因为它更简单通用。接下来创建server.js这是服务端的核心逻辑。const WebSocket require(ws); const axios require(axios); require(dotenv).config(); const wss new WebSocket.Server({ port: process.env.WS_PORT || 8080 }); console.log(WebSocket 服务器已启动监听端口: ${process.env.WS_PORT || 8080}); // 存储所有连接的客户端 let clients new Set(); // 存储上一次查询到的最新会话时间戳 let lastTimestamp Date.now() - 60000; // 默认从1分钟前开始查 wss.on(connection, (ws) { console.log(新的客户端连接); clients.add(ws); ws.on(close, () { console.log(客户端断开连接); clients.delete(ws); }); // 可以立即发送一条欢迎消息或当前状态 ws.send(JSON.stringify({ type: system, message: 已连接到Arkime实时通知服务, timestamp: new Date().toISOString() })); }); /** * 轮询Arkime API检查新的会话 */ async function pollNewSessions() { try { // 构造Arkime的搜索API URL // 这里我们搜索从 lastTimestamp 之后开始的新会话并按时间倒序排列取最前的几个 const query { query: { bool: { filter: [ { range: { firstPacket: { gt: lastTimestamp } } } ] } }, sort: [{ firstPacket: desc }], size: 10 // 每次最多取10条新会话避免数据量过大 }; const response await axios.post( ${process.env.ARKIME_API_BASE}/api/sessions?date-1, query, { headers: { Content-Type: application/json } } ); const sessions response.data?.data || []; if (sessions.length 0) { // 更新最后时间戳为最新会话的时间 lastTimestamp sessions[0].firstPacket; // 构建推送消息 const notification { type: new_sessions, count: sessions.length, sessions: sessions.map(s ({ id: s.id, srcIp: s.srcIp, dstIp: s.dstIp, srcPort: s.srcPort, dstPort: s.dstPort, protocol: s.protocol, firstPacket: new Date(s.firstPacket).toLocaleString(), lastPacket: new Date(s.lastPacket).toLocaleString(), bytes: s.bytes, packets: s.packets, node: s.node })), timestamp: new Date().toISOString() }; // 广播给所有连接的客户端 broadcast(notification); console.log(发现 ${sessions.length} 条新会话已推送。); } } catch (error) { console.error(轮询Arkime API时出错:, error.message); } } /** * 广播消息给所有客户端 */ function broadcast(data) { const message JSON.stringify(data); clients.forEach(client { if (client.readyState WebSocket.OPEN) { client.send(message); } }); } // 启动定时轮询 setInterval(pollNewSessions, process.env.POLL_INTERVAL || 5000); console.log(开始轮询间隔: ${process.env.POLL_INTERVAL || 5000}ms);这段代码做了以下几件事创建了一个WebSocket服务器。维护了一个clients集合来管理所有活跃连接。定义了一个pollNewSessions函数它定期向Arkime的sessionsAPI发送请求查询自上次检查以来出现的新会话。如果发现新会话就将其格式化然后通过broadcast函数推送给所有客户端。使用setInterval启动定时轮询。3.3 第三步创建测试客户端页面1分钟为了快速验证服务是否工作我们创建一个简单的HTML客户端。在项目根目录创建client.html。!DOCTYPE html html langzh-CN head meta charsetUTF-8 titleArkime 实时通知监控/title style body { font-family: sans-serif; margin: 20px; } #messages { border: 1px solid #ccc; padding: 10px; height: 400px; overflow-y: auto; margin-top: 10px; } .session { border-bottom: 1px dashed #eee; padding: 5px; margin: 5px 0; } .ip { font-weight: bold; color: #007bff; } .protocol { background-color: #e9ecef; padding: 2px 5px; border-radius: 3px; font-size: 0.9em; } /style /head body h1Arkime 实时会话通知/h1 div连接状态: span idstatus未连接/span/div button onclickconnectWebSocket()连接/button button onclickdisconnectWebSocket()断开/button hr div idmessages/div script let ws; const messagesDiv document.getElementById(messages); const statusSpan document.getElementById(status); // 修改这里的地址为你的WebSocket服务器地址 const wsUrl ws://localhost:8080; function connectWebSocket() { if (ws ws.readyState WebSocket.OPEN) { logMessage(系统, 已经连接了。); return; } ws new WebSocket(wsUrl); statusSpan.textContent 连接中...; statusSpan.style.color orange; ws.onopen () { statusSpan.textContent 已连接; statusSpan.style.color green; logMessage(系统, WebSocket连接已建立。); }; ws.onmessage (event) { const data JSON.parse(event.data); handleNotification(data); }; ws.onerror (error) { statusSpan.textContent 连接错误; statusSpan.style.color red; logMessage(系统, 连接错误: ${error.message}); }; ws.onclose () { statusSpan.textContent 已断开; statusSpan.style.color gray; logMessage(系统, WebSocket连接已关闭。); }; } function disconnectWebSocket() { if (ws) { ws.close(); ws null; } } function handleNotification(notification) { switch (notification.type) { case system: logMessage(系统, notification.message); break; case new_sessions: logMessage(新会话, 捕获到 ${notification.count} 条新会话); notification.sessions.forEach(session { const sessionEl document.createElement(div); sessionEl.className session; sessionEl.innerHTML strong${session.firstPacket}/strongbr span classip${session.srcIp}:${session.srcPort}/span - span classip${session.dstIp}:${session.dstPort}/span span classprotocol${session.protocol}/spanbr 数据包: ${session.packets}, 字节: ${session.bytes} (节点: ${session.node}) ; messagesDiv.appendChild(sessionEl); }); // 自动滚动到底部 messagesDiv.scrollTop messagesDiv.scrollHeight; break; default: logMessage(未知, JSON.stringify(notification)); } } function logMessage(sender, message) { const msgEl document.createElement(div); msgEl.innerHTML strong[${new Date().toLocaleTimeString()}] ${sender}:/strong ${message}; messagesDiv.appendChild(msgEl); } // 页面加载后自动连接可选 window.onload connectWebSocket; /script /body /html3.4 第四步启动与测试1分钟启动WebSocket服务器node server.js如果看到WebSocket 服务器已启动监听端口: 8080和开始轮询间隔: 5000ms的输出说明服务端启动成功。启动Arkime流量捕获确保你的Arkime Capture节点正在运行并捕获流量。打开客户端页面用浏览器直接打开client.html文件或通过一个简单的HTTP服务器如python3 -m http.server 3000然后访问http://localhost:3000/client.html。点击“连接”按钮。触发流量并观察在你的网络环境中产生一些新的网络流量例如访问一个网页ping一个地址。等待几秒不超过轮询间隔5秒你应该能在网页的“消息”区域看到新的会话信息被实时推送并显示出来。至此一个最基础的Arkime实时通知系统就搭建完成了。从创建项目到看到第一条推送确实可以在5分钟内完成。4. 从基础到生产高级配置与优化上面的例子只是一个起点它推送了所有新会话。在实际生产环境中我们需要更精细的控制、更稳定的服务和更丰富的事件类型。4.1 事件过滤与规则引擎推送给所有新会话会产生大量噪音。我们需要一个规则引擎来过滤。修改server.js中的pollNewSessions函数在构造查询时加入过滤条件。例如只推送源IP或目的IP为特定内网网段的会话const query { query: { bool: { filter: [ { range: { firstPacket: { gt: lastTimestamp } } }, { bool: { should: [ { prefix: { srcIp: 10.0.0. } }, // 源IP是10.0.0.0/24网段 { prefix: { dstIp: 192.168.1. } } // 目的IP是192.168.1.0/24网段 ], minimum_should_match: 1 } } ] } }, sort: [{ firstPacket: desc }], size: 20 };更进一步我们可以将规则配置化。创建一个rules.json文件[ { name: 监控内部服务器对外连接, conditions: { srcIp: [10.0.1.10, 10.0.1.11], dstIp_not_in: [10.0.0.0/8, 192.168.0.0/16] }, action: notify, channel: critical_alert }, { name: 检测高危端口扫描, conditions: { dstPort: [22, 3389, 445], packets: { lt: 3 } }, action: notify, channel: security_warning } ]然后在服务器代码中加载这个规则文件对查询到的每个会话进行匹配只有匹配规则的会话才会被推送并且可以指定推送到不同的“频道”对应不同的客户端订阅。4.2 支持多种事件类型除了“新会话”Arkime还有很多其他有价值的事件可以推送字段值统计变化某个特定标签tags的会话数突然激增。SPI View变化用户保存的视图SPI View中有新的匹配项。系统状态Capture节点离线、磁盘空间不足。这些事件的获取方式不同。例如获取字段统计可能需要调用Arkime的api/stats接口监听SPI View可能需要定期查询api/spiview。我们可以为每种事件类型编写独立的轮询函数并统一通过WebSocket服务管理。在server.js中可以扩展为多个定时任务// 引入规则 const rules require(./rules.json); // 轮询新会话 setInterval(() pollAndFilter(new_sessions, rules), POLL_INTERVAL_SESSIONS); // 轮询系统状态间隔可以长一些 setInterval(pollSystemHealth, POLL_INTERVAL_HEALTH); // 轮询特定SPI View setInterval(() pollSpiview(my_watchlist), POLL_INTERVAL_SPIVIEW);4.3 连接管理与心跳机制在生产环境中需要更健壮的连接管理。心跳保活WebSocket连接可能因为网络问题或代理超时而断开。需要在客户端和服务端实现心跳机制Ping/Pong。服务端server.js可以定期向每个客户端发送Ping帧。客户端client.html中的JavaScript需要监听onmessage事件识别Ping帧并回复Pong。断线重连客户端应该监测连接状态在断开后尝试指数退避重连。连接认证在建立WebSocket连接时可以要求客户端提供Token或进行HTTP Basic认证WebSocket握手阶段支持HTTP头。一个简单的心跳实现可以在服务端添加// 在 connection 事件内 ws.isAlive true; ws.on(pong, () { ws.isAlive true; }); // 全局定时器每隔30秒检查一次所有客户端连接 setInterval(() { wss.clients.forEach((client) { if (client.isAlive false) { console.log(终止无响应客户端连接); return client.terminate(); } client.isAlive false; client.ping(); // 发送Ping帧 }); }, 30000);4.4 性能优化与扩展使用Elasticsearch PITSearch After替代Scroll API如果直接查询ES对于大量数据Scroll API会消耗大量资源。使用Point-in-Time (PIT) 和search_after参数进行高效的深度分页查询更适合持续增量获取的场景。消息队列解耦当通知规则非常复杂或客户端数量极多时轮询逻辑和推送逻辑可以解耦。轮询服务将事件投递到消息队列如Redis Pub/Sub, RabbitMQ, Kafka再由多个推送服务实例消费队列并处理WebSocket推送。这提高了系统的可扩展性和可靠性。客户端分组频道/房间不是所有客户端都需要所有消息。可以实现一个简单的“订阅”机制客户端在连接时可以发送一个subscribe消息指定它关心的频道如security,performance,all服务端只向订阅了相应频道的客户端推送消息。5. 常见问题与排查实录在实际部署和运行中我遇到并解决了一些典型问题这里记录下来供你参考。5.1 问题一收不到任何推送消息检查点1WebSocket服务器是否正常运行运行netstat -tlnp | grep :8080(Linux) 或Get-NetTCPConnection -LocalPort 8080(Windows PowerShell) 查看端口是否被监听。查看服务器控制台是否有错误日志。检查点2客户端连接是否正确在浏览器中按F12打开开发者工具查看“网络”(Network)选项卡过滤WS类型检查WebSocket连接状态是否为101 Switching Protocols。如果连接失败检查client.html中的wsUrl地址和端口是否正确以及服务器防火墙是否放行了该端口。检查点3轮询逻辑是否获取到数据在server.js的pollNewSessions函数中添加console.log(查询结果:, sessions.length);查看控制台输出。如果始终为0可能是lastTimestamp初始值太新或者查询条件太严格。尝试将lastTimestamp初始化为一个更早的时间如Date.now() - 3000005分钟前。检查点4Arkime API是否可访问在服务器上使用curl命令测试Arkime APIcurl -X POST http://your-arkime-viewer:8005/api/sessions?date-1 -H Content-Type: application/json -d {query:{match_all:{}}}。确保能返回JSON数据。5.2 问题二推送延迟高或不稳定可能原因1轮询间隔POLL_INTERVAL设置不当。解决根据你对实时性的要求调整。太短如1秒会给Arkime API和自身服务带来不必要的压力太长如30秒则延迟明显。对于一般监控5-10秒是一个平衡点。可能原因2Arkime Viewer API响应慢。解决优化查询。避免使用过于复杂的聚合查询确保查询的字段有索引。考虑直接使用Elasticsearch API并确保查询使用了合理的过滤条件限制返回大小size。可能原因3网络问题。解决确保WebSocket服务器、Arkime Viewer和客户端之间的网络延迟低且稳定。如果跨地域需要考虑网络优化。5.3 问题三客户端连接频繁断开可能原因1缺乏心跳机制被中间网络设备如负载均衡器、代理超时断开。解决按照前面“连接管理与心跳机制”一节实现服务端的Ping/Pong保活。可能原因2客户端页面被浏览器休眠或切到后台。解决在客户端代码中监听visibilitychange事件当页面从后台切回前台时检查WebSocket连接状态必要时重新连接。可能原因3服务器端资源耗尽。解决监控Node.js进程的内存和CPU使用情况。确保及时清理断开的客户端连接代码中已有on(close)事件处理。对于大量连接考虑使用ws库的permessage-deflate扩展开启压缩或升级服务器配置。5.4 问题四如何将通知发送到其他平台如钉钉、企业微信我们的WebSocket服务推送的是标准化JSON消息。你可以编写一个“桥接”客户端它连接到WebSocket服务器收到消息后再调用第三方平台的API。例如一个简单的Python桥接脚本使用websockets库和requests库import asyncio import websockets import json import requests async def listen_and_forward(): uri ws://your-websocket-server:8080 async with websockets.connect(uri) as websocket: while True: message await websocket.recv() data json.loads(message) if data.get(type) new_sessions: # 格式化消息 text fArkime告警发现{data[count]}条新会话\n for s in data[sessions][:3]: # 只取前3条 text f- {s[srcIp]}:{s[srcPort]} - {s[dstIp]}:{s[dstPort]} ({s[protocol]})\n # 调用钉钉机器人 dingtalk_webhook https://oapi.dingtalk.com/robot/send?access_tokenYOUR_TOKEN payload { msgtype: text, text: {content: text} } requests.post(dingtalk_webhook, jsonpayload) asyncio.run(listen_and_forward())这个脚本作为一个常驻进程运行就能把WebSocket通知转化为钉钉群消息。整个配置过程的核心在于理解Arkime的数据产出方式并选择一个轻量、灵活的中介WebSocket服务来桥接数据生产者和消费者。五分钟搭建的只是一个原型但基于这个原型你可以根据实际监控需求扩展出非常强大和个性化的实时网络安全态势感知系统。