基于Redis Streams的航空乘客安全监测与实时预警系统
2026/8/28 10:28:18 网站建设 项目流程

航空乘客安全是一个典型的实时数据融合场景。值机系统、安检系统、登机口广播和客舱服务记录中,可能同时出现同一名乘客或同一航班的多个安全相关事件。由于系统隔离,处置人员很难在第一时间看到完整上下文,更难以在事件刚发生时就触发联动处置。本文以航空乘客安全监测与预警系统的最小可运行版本为例,说明如何用 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后台定时任务用于演示阶段消费 Stream
SQLAlchemy2.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:30+08: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、security
flight_nostring航班号,可空
passenger_idstring旅客标识,可空
event_typestring事件类型,如 abnormal_checkin
severitystring事件初始等级,如 low、medium、high
detailobject事件详情,格式随事件类型变化

detail 字段使用对象而非固定字段,是为了兼容不同来源的扩展属性。规则引擎读取 detail 时,需要判断字段是否存在,不能假定所有事件都有相同结构。

2.4 模块划分与职责

模块职责关键技术点
采集接口接收外部事件,校验格式,写入 Stream请求限流、字段校验
消费服务从 Stream 拉取消息,执行规则,写告警消费组、消息确认、幂等
规则引擎根据事件类型和阈值产生告警规则可配置、阈值外置
WebSocket 服务向监控页面推送实时告警连接管理、断线重连
存储层保存原始事件和告警记录MySQL JSON 字段、索引设计

模块之间的依赖要尽量单向。采集接口不直接调用规则引擎,规则引擎也不反向请求采集接口,所有数据流动都通过 Redis Stream 完成。这样当事件量增大时,可以单独扩展消费服务的实例数量。

3. 环境准备与项目初始化

3.1 本地环境要求

开始编码前,需要确认本机已经有下面这些软件。版本号是示例环境,实际请以官方稳定版为准。

软件版本建议用途检查命令
Python3.10 或更高运行示例代码python --version
Redis7.x消息队列redis-server --version
MySQL8.x数据存储mysql --version
Docker可选快速启动 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 中写入以下内容:

fastapi==0.115.0 uvicorn[standard]==0.30.6 redis==5.0.7 PyMySQL==1.1.1 SQLAlchemy==2.0.32 PyYAML==6.0.2 APScheduler==3.10.4 pydantic==2.8.2 websockets==12.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_PASSWORD=root \ -e MYSQL_DATABASE=flightguard \ 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) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; 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) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

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: critical

abnormal_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", encoding="utf-8") as f: config = yaml.safe_load(f) app = FastAPI(title="FlightGuard Collector") redis_client = aioredis.from_url( f"redis://{config['redis']['host']}:{config['redis']['port']}" ) class SafetyEvent(BaseModel): event_id: str = Field(default_factory=lambda: f"evt_{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_code=500, detail=f"write 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", encoding="utf-8") as f: config = yaml.safe_load(f) redis_client = aioredis.from_url( f"redis://{config['redis']['host']}:{config['redis']['port']}" ) engine = create_engine( f"mysql+pymysql://{config['mysql']['user']}:{config['mysql']['password']}@" f"{config['mysql']['host']}:{config['mysql']['port']}/" f"{config['mysql']['database']}?charset=utf8mb4" ) 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_code="RULE_ABNORMAL_CHECKIN", level=rule.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_code="RULE_CONTRABAND", level=rule.get("alert_level", "critical") )) return alerts def build_alert(event: dict, rule_code: str, level: str) -> dict: return { "alert_id": f"alert_{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 = f"worker-{uuid.uuid4().hex[:8]}" try: await redis_client.xgroup_create( stream_key, group_name, id="0", mkstream=True ) except Exception: pass entries = await redis_client.xreadgroup( group_name=group_name, consumer_name=consumer_name, streams={stream_key: ">"}, count=20, block=1000 ) for stream, messages in entries: for msg_id, fields in messages: raw = fields[b"payload"].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 读取消息后,如果直接处理但忘记 xack,Redis 会把这些消息留在 pending 列表中。当消费服务重启时会再次收到这些消息,容易造成重复告警。使用 ON DUPLICATE KEY UPDATE 更新 alert_id 是幂等保护手段之一,但不能完全替代 xack 的正常执行。

在 FastAPI 主进程里,通过 APScheduler 启动消费任务:

from apscheduler.schedulers.asyncio import AsyncIOScheduler scheduler = AsyncIOScheduler() scheduler.add_job(consume_events, "interval", seconds=2, id="consume_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 charset="utf-8" /> <title>FlightGuard 安全监控</title> <script src="https://cdn.jsdelivr.net/npm/echarts@5/dist/echarts.min.js"></script> </head> <body> <div id="chart" style="width:100%;height:300px;"></div> <ul id="alerts"></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:30+08: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_checkin,checkin_count=4RULE_ABNORMAL_CHECKIN产生 high 告警
abnormal_checkin,checkin_count=2不产生告警
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 广播不在同一个进程。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询