ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

基于Redis Streams的航空乘客安全监测与实时预警系统

基于Redis Streams的航空乘客安全监测与实时预警系统 航空乘客安全是一个典型的实时数据融合场景。值机系统、安检系统、登机口广播和客舱服务记录中可能同时出现同一名乘客或同一航班的多个安全相关事件。由于系统隔离处置人员很难在第一时间看到完整上下文更难以在事件刚发生时就触发联动处置。本文以航空乘客安全监测与预警系统的最小可运行版本为例说明如何用 Python、Redis Streams 和 MySQL 构建一条从事件采集、规则判定到实时告警的数据链路并给出可复现的代码、验证方法和排错路径。读者可以把这个系统当作理解航班安全数据平台的基础骨架后续再按生产需求扩展。1. 先拆解航空乘客安全系统的业务场景与核心难点1.1 业务场景中的数据孤岛现象航空乘客安全涉及的业务环节通常包括值机、安检、登机、客舱运行和到达后的处置。每个环节都有自己的业务系统数据格式、更新频率和实时性差别很大。以一次异常事件为例乘客在短时间内反复改签值机系统会留下操作记录安检系统可能上报疑似违禁品登机口广播系统可能记录乘客未登机客舱服务系统可能记录乘客状态异常。这些记录如果只停留在各自系统里后续的调度员和安全处置团队就无法形成完整事件视图。下面的表格列出了常见环节、数据来源、安全事件示例和实时性要求。业务环节典型数据源安全事件示例实时性要求值机值机系统、离港系统短时间内多次取消或改签秒级到分钟级安检安检管理系统、违禁品上报查获疑似违禁品秒级登机登机口系统、广播系统视频或人工上报异常行为秒级客舱客舱服务终端、机组上报乘客状态异常、干扰机组秒级到分钟级到达后地面服务系统旅客未按时离开候机区分钟级不同来源的数据如果只是简单接入数据库最多只能用于事后查询。真正困难的地方在于事件发生到告警触发的延迟要足够短同时要把同一乘客、同一航班的多个事件关联起来形成对安全态势的判断。1.2 技术难点的四个维度单看任意一个安全事件规则判断并不复杂。真正影响系统质量的是四个维度。第一是异构数据。不同来源的事件字段不同有的有航班号有的只有旅客编号有的时间格式还不同。系统需要在入口处统一格式否则规则引擎难以处理。第二是实时性。安全告警的价值随时间衰减从事件产生到处置人员看到告警通常要求秒级完成。演示系统可以用定时轮询生产环境则需要用消息队列加流式计算确保吞吐量和延迟可控。第三是准确性。告警过多会让处置人员疲劳漏报又会导致真实风险被忽略。准确性不仅依赖规则阈值还依赖事件关联能力。例如单次取消改签可能正常但同一旅客在十分钟内反复操作就需要提高关注级别。第四是可追溯性。告警产生后需要知道是哪个事件触发的、命中哪条规则、当时的原始数据是什么、处置状态如何。因此系统不能只保存结果还需要保存原始事件和规则版本。1.3 最小可行系统的边界本文要实现的 FlightGuard 演示系统只覆盖核心数据链路接收安全事件、写入消息队列、规则判定、生成告警、推送前端、落库保存。它不包含人脸识别、视频分析、硬件设备对接和人工任务调度也不涉及具体机场的私有协议。这样设计的好处是让读者先理解数据流再根据真实场景扩展。演示系统选择了 Python 技术栈使用 FastAPI 提供接口Redis Streams 作为轻量消息队列MySQL 保存事件和告警前端用 WebSocket 接收实时告警。它的目标是跑通“一条模拟安全事件从进入系统到展示在页面上”的最小闭环。2. 系统总体设计数据从哪来、到哪里去2.1 数据链路总览整个系统的数据流可以描述为安全事件源通过 REST 接口上报到采集服务采集服务将事件写入 Redis Stream后台规则消费服务从 Stream 中拉取新消息执行规则判定命中的规则生成告警记录写入 MySQL并通过 WebSocket 推送给监控页面未命中的事件只落库保存作为后续分析的原始数据。[安全事件源] - [采集接口] - [Redis Stream] - [规则消费服务] - [MySQL] | v [WebSocket] - [监控页面]这个链路的关键点是引入了 Redis Stream。事件不直接写入业务表而是先进入消息通道。这样即使短时间涌入大量事件采集接口也能快速响应消费服务可以根据能力逐步处理避免数据库被瞬时写入打爆。2.2 技术选型与理由组件示例版本用途选型说明Python3.10开发语言生态完善FastAPI 对异步支持好FastAPI0.115REST 接口和 WebSocket自带 OpenAPI 文档异步性能强Redis7.xStream 消息队列安装简单支持消费组和消息确认MySQL8.x结构化数据存储支持 JSON 字段适合保存异构事件APScheduler3.10后台定时任务用于演示阶段消费 StreamSQLAlchemy2.0数据库访问屏蔽 SQL 差异便于迁移ECharts5.x前端图表实时曲线和统计图成熟有人在演示阶段会直接选择 Kafka理由是生产环境最终要用 Kafka。但 Redis Streams 的消费组机制和 Kafka 有相似之处安装和维护成本却低得多。用 Redis Streams 跑通业务逻辑后续切成 Kafka 时主要改动集中在消费端和序列化层。对于学习项目来说这是更平滑的路径。2.3 统一事件格式约定为了让规则引擎不关心数据来源采集接口需要把所有来源的事件转换成统一 JSON 格式。示例事件如下{ event_id: evt_20240920_001, event_time: 2024-09-20T10:15:3008:00, source: checkin, flight_no: CA1234, passenger_id: P000123, event_type: abnormal_checkin, severity: medium, detail: { checkin_count: 4, interval_minutes: 8 } }字段含义如下字段类型说明event_idstring事件唯一标识建议调用方生成event_timestringISO8601 时间带时区sourcestring来源系统如 checkin、securityflight_nostring航班号可空passenger_idstring旅客标识可空event_typestring事件类型如 abnormal_checkinseveritystring事件初始等级如 low、medium、highdetailobject事件详情格式随事件类型变化detail 字段使用对象而非固定字段是为了兼容不同来源的扩展属性。规则引擎读取 detail 时需要判断字段是否存在不能假定所有事件都有相同结构。2.4 模块划分与职责模块职责关键技术点采集接口接收外部事件校验格式写入 Stream请求限流、字段校验消费服务从 Stream 拉取消息执行规则写告警消费组、消息确认、幂等规则引擎根据事件类型和阈值产生告警规则可配置、阈值外置WebSocket 服务向监控页面推送实时告警连接管理、断线重连存储层保存原始事件和告警记录MySQL JSON 字段、索引设计模块之间的依赖要尽量单向。采集接口不直接调用规则引擎规则引擎也不反向请求采集接口所有数据流动都通过 Redis Stream 完成。这样当事件量增大时可以单独扩展消费服务的实例数量。3. 环境准备与项目初始化3.1 本地环境要求开始编码前需要确认本机已经有下面这些软件。版本号是示例环境实际请以官方稳定版为准。软件版本建议用途检查命令Python3.10 或更高运行示例代码python --versionRedis7.x消息队列redis-server --versionMySQL8.x数据存储mysql --versionDocker可选快速启动 Redis 和 MySQLdocker --version推荐在虚拟环境中安装 Python 依赖避免与系统 Python 环境冲突。如果本机没有 Redis 和 MySQL使用 Docker 启动是成本最低的方式。3.2 创建项目目录结构项目目录结构如下模块按职责拆分后续增加功能时比较容易定位。flightguard/ ├── app/ │ ├── main.py │ ├── config.py │ ├── models.py │ ├── collectors.py │ ├── rules.py │ ├── alerting.py │ └── websocket.py ├── static/ │ └── index.html ├── sql/ │ └── init.sql ├── requirements.txt └── config.yaml创建目录后先写 requirements.txt 和 config.yaml再逐个实现 Python 模块。这样能让每一阶段的目标更清晰。3.3 安装 Python 依赖在 requirements.txt 中写入以下内容fastapi0.115.0 uvicorn[standard]0.30.6 redis5.0.7 PyMySQL1.1.1 SQLAlchemy2.0.32 PyYAML6.0.2 APScheduler3.10.4 pydantic2.8.2 websockets12.0然后执行安装命令python -m venv .venv source .venv/bin/activate pip install -r requirements.txt这一步常见的错误是直接使用系统 Python 安装导致权限问题或污染全局环境。使用虚拟环境后即使依赖版本冲突也可以随时删除重建。3.4 启动 Redis 和 MySQL使用 Docker 启动两个依赖服务docker run -d --name flightguard-redis -p 6379:6379 redis:7-alpine docker run -d --name flightguard-mysql \ -p 3306:3306 \ -e MYSQL_ROOT_PASSWORDroot \ -e MYSQL_DATABASEflightguard \ mysql:8.0启动后可以用下面命令确认服务状态docker ps redis-cli ping mysql -h127.0.0.1 -uroot -proot -e select version();这里要注意MySQL 容器首次启动需要初始化刚执行完 docker run 后立刻连接可能失败。等待 10 秒左右再检查是常见操作。如果使用非 Docker 环境需要确保本地 Redis 和 MySQL 已安装并启动并且数据库 flightguard 已创建。4. 核心实现从事件采集到实时告警4.1 MySQL 数据模型设计需要两张核心表safety_event 保存原始安全事件alert_record 保存规则判定后生成的告警。safety_event 表用于追溯和分析alert_record 表用于处置闭环。在 sql/init.sql 中写入CREATE TABLE IF NOT EXISTS safety_event ( id BIGINT PRIMARY KEY AUTO_INCREMENT, event_id VARCHAR(64) NOT NULL UNIQUE, event_time DATETIME NOT NULL, source VARCHAR(32) NOT NULL, flight_no VARCHAR(16) DEFAULT NULL, passenger_id VARCHAR(32) DEFAULT NULL, event_type VARCHAR(64) NOT NULL, severity VARCHAR(16) DEFAULT low, detail JSON, raw_data JSON, created_at DATETIME DEFAULT CURRENT_TIMESTAMP, KEY idx_flight_time (flight_no, event_time), KEY idx_event_time (event_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; CREATE TABLE IF NOT EXISTS alert_record ( id BIGINT PRIMARY KEY AUTO_INCREMENT, alert_id VARCHAR(64) NOT NULL UNIQUE, event_id VARCHAR(64) NOT NULL, rule_code VARCHAR(64) NOT NULL, alert_level VARCHAR(16) NOT NULL, alert_message VARCHAR(512) DEFAULT NULL, flight_no VARCHAR(16) DEFAULT NULL, passenger_id VARCHAR(32) DEFAULT NULL, status VARCHAR(16) DEFAULT pending, created_at DATETIME DEFAULT CURRENT_TIMESTAMP, handled_at DATETIME DEFAULT NULL, handler VARCHAR(64) DEFAULT NULL, KEY idx_alert_status (status), KEY idx_alert_time (created_at) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;detail 和 raw_data 使用 JSON 类型是因为安全事件的 detail 结构不固定。MySQL 8 的 JSON 类型既支持索引也支持通过函数提取字段比把所有 detail 字段都展开成列更灵活。event_id 和 alert_id 都加了唯一约束这是实现幂等的关键。4.2 配置文件规则阈值与连接参数在 config.yaml 中维护数据库连接、Redis 连接和规则阈值。把规则阈值放在配置文件里而不是硬编码到代码中是为了后续调整规则时不需要重新发布应用。server: host: 0.0.0.0 port: 8000 redis: host: 127.0.0.1 port: 6379 stream_key: safety:events consumer_group: alert_group mysql: host: 127.0.0.1 port: 3306 user: root password: root database: flightguard rules: abnormal_checkin: enabled: true max_count: 3 window_minutes: 10 alert_level: high security_contraband: enabled: true alert_level: criticalabnormal_checkin 规则的含义是同一旅客在 10 分钟内值机操作次数达到 3 次时触发 high 级告警。security_contraband 规则更简单只要事件类型是安检违禁品就触发 critical 级告警。生产环境中这些参数应该放到配置中心并保留规则变更记录。4.3 数据采集接口事件写入 Redis Stream采集接口使用 FastAPI 实现接收 JSON 事件后写入 Redis Stream。示例代码中省略了认证和限流但生产环境必须有。import uuid import yaml import redis.asyncio as aioredis from fastapi import FastAPI, HTTPException from pydantic import BaseModel, Field with open(config.yaml, r, encodingutf-8) as f: config yaml.safe_load(f) app FastAPI(titleFlightGuard Collector) redis_client aioredis.from_url( fredis://{config[redis][host]}:{config[redis][port]} ) class SafetyEvent(BaseModel): event_id: str Field(default_factorylambda: fevt_{uuid.uuid4().hex[:12]}) event_time: str source: str flight_no: str passenger_id: str event_type: str severity: str low detail: dict {} app.post(/api/v1/events) async def receive_event(event: SafetyEvent): try: payload event.model_dump_json() await redis_client.xadd( config[redis][stream_key], {payload: payload}, id* ) return {code: 0, message: ok, event_id: event.event_id} except Exception as e: raise HTTPException(status_code500, detailfwrite stream failed: {e})这里没有直接把事件写入 MySQL而是先写入 Redis Stream。原因是采集接口的响应速度不应该受数据库写入速度影响。当瞬时事件量很大时Stream 充当缓冲消费服务可以按自己的节奏处理。xadd 的 id 参数使用*让 Redis 自动生成消息 ID保证消息顺序性。4.4 规则引擎与消费服务从 Stream 到告警消费服务的主要工作是从 Redis Stream 读取新消息、解析事件、调用规则函数、把命中的告警写入 MySQL最后确认消息已被处理。import asyncio import json import uuid import yaml import redis.asyncio as aioredis from sqlalchemy import create_engine, text with open(config.yaml, r, encodingutf-8) as f: config yaml.safe_load(f) redis_client aioredis.from_url( fredis://{config[redis][host]}:{config[redis][port]} ) engine create_engine( fmysqlpymysql://{config[mysql][user]}:{config[mysql][password]} f{config[mysql][host]}:{config[mysql][port]}/ f{config[mysql][database]}?charsetutf8mb4 ) def evaluate_rules(event: dict) - list: alerts [] event_type event.get(event_type) detail event.get(detail, {}) rule config[rules].get(abnormal_checkin) if rule and rule.get(enabled) and event_type abnormal_checkin: if detail.get(checkin_count, 0) rule.get(max_count, 3): alerts.append(build_alert( event, rule_codeRULE_ABNORMAL_CHECKIN, levelrule.get(alert_level, high) )) rule config[rules].get(security_contraband) if rule and rule.get(enabled) and event_type security_contraband: alerts.append(build_alert( event, rule_codeRULE_CONTRABAND, levelrule.get(alert_level, critical) )) return alerts def build_alert(event: dict, rule_code: str, level: str) - dict: return { alert_id: falert_{uuid.uuid4().hex[:12]}, event_id: event[event_id], rule_code: rule_code, alert_level: level, alert_message: f{rule_code} triggered by {event.get(event_type)}, flight_no: event.get(flight_no), passenger_id: event.get(passenger_id), status: pending } def save_alerts_sync(alerts: list) - None: if not alerts: return with engine.begin() as conn: for alert in alerts: conn.execute( text( INSERT INTO alert_record (alert_id, event_id, rule_code, alert_level, alert_message, flight_no, passenger_id, status) VALUES (:alert_id, :event_id, :rule_code, :alert_level, :alert_message, :flight_no, :passenger_id, :status) ON DUPLICATE KEY UPDATE alert_id alert_id ), alert ) async def consume_events(): stream_key config[redis][stream_key] group_name config[redis][consumer_group] consumer_name fworker-{uuid.uuid4().hex[:8]} try: await redis_client.xgroup_create( stream_key, group_name, id0, mkstreamTrue ) except Exception: pass entries await redis_client.xreadgroup( group_namegroup_name, consumer_nameconsumer_name, streams{stream_key: }, count20, block1000 ) for stream, messages in entries: for msg_id, fields in messages: raw fields[bpayload].decode() event json.loads(raw) alerts evaluate_rules(event) if alerts: await asyncio.to_thread(save_alerts_sync, alerts) await redis_client.xack(stream_key, group_name, msg_id)重点是消息确认。xreadgroup 读取消息后如果直接处理但忘记 xackRedis 会把这些消息留在 pending 列表中。当消费服务重启时会再次收到这些消息容易造成重复告警。使用 ON DUPLICATE KEY UPDATE 更新 alert_id 是幂等保护手段之一但不能完全替代 xack 的正常执行。在 FastAPI 主进程里通过 APScheduler 启动消费任务from apscheduler.schedulers.asyncio import AsyncIOScheduler scheduler AsyncIOScheduler() scheduler.add_job(consume_events, interval, seconds2, idconsume_events) scheduler.start()演示项目把采集接口和消费任务放在同一个进程中好处是启动简单。生产环境建议拆成两个服务采集服务只负责写 Stream消费服务独立部署这样可以根据消息积压情况单独扩容消费实例。4.5 WebSocket 推送与前端监控页面告警写入 MySQL 后还需要主动推送到监控页面。这里使用 FastAPI 的 WebSocket 接口维护一个在线连接列表。from fastapi import WebSocket, WebSocketDisconnect class ConnectionManager: def __init__(self): self.active_connections [] async def connect(self, websocket: WebSocket): await websocket.accept() self.active_connections.append(websocket) def disconnect(self, websocket: WebSocket): if websocket in self.active_connections: self.active_connections.remove(websocket) async def broadcast(self, message: str): for conn in self.active_connections[:]: try: await conn.send_text(message) except Exception: self.disconnect(conn) manager ConnectionManager() app.websocket(/ws/alerts) async def websocket_endpoint(websocket: WebSocket): await manager.connect(websocket) try: while True: await websocket.receive_text() except WebSocketDisconnect: manager.disconnect(websocket)在消费服务保存告警后调用 manager.broadcast 将告警 JSON 推送到前端。前端页面放在 static/index.html使用 WebSocket 接收消息并展示。!DOCTYPE html html head meta charsetutf-8 / titleFlightGuard 安全监控/title script srchttps://cdn.jsdelivr.net/npm/echarts5/dist/echarts.min.js/script /head body div idchart stylewidth:100%;height:300px;/div ul idalerts/ul script const chart echarts.init(document.getElementById(chart)); chart.setOption({ xAxis: { type: time }, yAxis: { type: value, name: 告警数 }, series: [{ name: alerts, type: line, data: [] }] }); const ws new WebSocket(ws://localhost:8000/ws/alerts); ws.onmessage function(evt) { const msg JSON.parse(evt.data); const li document.createElement(li); li.textContent msg.alert_id msg.flight_no msg.alert_level msg.alert_message; document.getElementById(alerts).prepend(li); }; /script /body /html这个页面只实现了最基本的列表展示。生产环境中告警列表需要分页、过滤、确认按钮和处置记录图表也需要支持按航班、事件类型、告警等级聚合统计。5. 运行验证从模拟数据到告警效果5.1 初始化数据库并启动服务先导入表结构mysql -h127.0.0.1 -uroot -proot flightguard sql/init.sql然后启动 FastAPI 服务uvicorn app.main:app --host 0.0.0.0 --port 8000启动后访问 http://127.0.0.1:8000/docs 可以看到接口文档。如果页面正常返回说明 FastAPI 服务和路由已经加载成功。5.2 构造并发送模拟安全事件发送一条命中 abnormal_checkin 规则的事件curl -X POST http://127.0.0.1:8000/api/v1/events \ -H Content-Type: application/json \ -d { event_time: 2024-09-20T10:15:3008:00, source: checkin, flight_no: CA1234, passenger_id: P000123, event_type: abnormal_checkin, severity: medium, detail: { checkin_count: 4, interval_minutes: 8 } }响应应该类似{code:0,message:ok,event_id:evt_xxxx}如果 event_id 没有传FastAPI 的 default_factory 会生成一个。示例中为了便于追踪可以自己传一个 event_id。5.3 检查 Redis Stream 和 MySQL 告警表查看 Stream 长度和消费组状态redis-cli XLEN safety:events redis-cli XINFO GROUPS safety:events查看消息确认情况redis-cli XPENDING safety:events alert_group正常情况下消费成功后 pending 数量会归零。如果看到大量 pending 消息说明消费处理中存在异常或没有执行 xack。查询告警表mysql -h127.0.0.1 -uroot -proot flightguard \ -e select alert_id, rule_code, alert_level, flight_no, passenger_id, status from alert_record order by id desc limit 10;如果模拟事件正常命中规则会看到一条 RULE_ABNORMAL_CHECKIN 的告警。5.4 浏览器端实时告警验证浏览器打开 http://127.0.0.1:8000/static/index.html。再通过 curl 发送一条新的告警事件页面列表应该自动增加新记录。如果列表没有变化优先检查浏览器控制台是否出现 WebSocket 连接错误以及服务端是否执行了 broadcast。5.5 预期结果对照表输入事件命中的规则预期结果abnormal_checkincheckin_count4RULE_ABNORMAL_CHECKIN产生 high 告警abnormal_checkincheckin_count2无不产生告警security_contrabandRULE_CONTRABAND产生 critical 告警未知 event_type无只落库不产生告警6. 常见问题排查从现象定位根因6.1 事件写入成功但告警没有产生现象是 curl 返回 code 0事件也确实进入了 Redis Stream但 MySQL 里没有对应的告警记录。优先检查这几个点消费任务是否启动配置里的规则是否 enabled事件里的 event_type 是否为规则配置的类型detail 字段中的数字是否达到阈值。redis-cli XPENDING safety:events alert_group如果 pending 里有消息说明消费任务读取后没有 xack很可能是 save_alerts 抛异常了。检查服务端日志里是否有 SQLAlchemy 或 MySQL 错误比如表不存在、字段长度溢出、JSON 格式不对。问题现象可能原因检查方式处理建议告警未产生消费任务未启动查看服务启动日志确认 scheduler.start() 已执行告警未产生规则被禁用检查 config.yaml 中 enabled改为 true 后重启告警未产生pending 持续增长xpending 查看待确认消息修复代码并手工 ack 或重置告警未产生MySQL 写入异常查看日志堆栈根据报错修复表结构或连接6.2 WebSocket 页面收不到消息页面能正常打开但发送新事件后页面没有更新。检查 WebSocket 地址是否写错。如果页面通过 8000 端口访问WebSocket 地址通常应该使用ws://localhost:8000/ws/alerts。如果 FastAPI 服务配置了 CORS还需要确认是否允许前端的来源。浏览器控制台如果出现 403 或跨域错误需要添加对应的中间件。另一个常见问题是消费服务和 WebSocket 广播不在同一个进程。
返回列表