
简介这是一套面向工业现场的可扩展Andon系统实现方案采用TypeScript编写旨在帮助制造企业优化异常响应流程、平滑过渡至预测性维护模式并通过实时告警与可视化看板提前暴露生产隐患。资源共158个文件压缩包仅710KB核心包括41个ts与29个tsx前端及业务逻辑源码、32个vtl页面模板、26个json配置文件另附js脚本、GraphQL接口定义、测试配置和部署shell脚本目录结构清晰适合需要快速搭建或二次开发Andon系统的开发与运维人员。当前已有112人学习使用。从内容预览可见系统内置完整的GraphQL交互模块涵盖查询、变更与订阅三类操作并配有架构说明图与看板入口可帮助读者理解端到端数据流同时包含工程化测试与版本忽略配置便于直接进入开发调试。整体上这份资源兼顾业务设计与工程落地既能支撑现场流程优化也可作为预测性维护系统前端与接口层的参考实现。1. 让 Andon 系统从“拉绳报警”长出 TypeScript 的可扩展骨架Andon 在很多传统工厂里还是那根挂着的绳子员工一拉灯亮、音乐响班组长跑过去。这个机制解决了“异常被看见”但解决不了“异常被记录、被分析、被提前阻止”。现代 Andon 系统的核心是报警之后的事件链谁报的、什么类型、多久被响应、多久被关闭、同一台设备一周报了三次——这些数据一旦结构化就同时服务于流程优化与预测性维护。标题里的“可扩展”分两层一是接入规模从一条产线到多个工厂二是业务扩展从被动的响应工具迁移到主动的设备健康预警。TypeScript 在这里的价值不只是补类型而是用判别联合把 Andon 状态机约束在编译期。这套方案适合做 MES、IIoT、工业应用集成的后端工程师也适合不想一上来就铺 Kafka、又想给后续数据分析留数据底座的团队。2. 用 TypeScript 事件模型兜住 Andon 的领域约束2.1 判别联合把“报警状态”从自由字符串变成编译期约束很多 Andon 项目从一张andon_events表起步status字段是字符串pending、open、resolved混着写。写着写着就出现Resolved和resolved并存的脏数据流程统计直接受影响。用 TypeScript 做领域层第一步是把状态的定义权从数据库收回到代码里。我一般会用判别联合discriminated union定义 Andon 的状态机export type AndonType equipment | material | quality | safety; export type AndonPhase | { phase: raised; at: string; stationId: string } | { phase: acknowledged; at: string; responderId: string } | { phase: escalated; at: string; level: 1 | 2 | 3 } | { phase: resolved; at: string; durationSec: number; closureCode: string }; export interface AndonCall { id: string; type: AndonType; lineId: string; stationId: string; history: AndonPhase[]; }phase字段就是判别字段每个阶段自带自己的上下文。level被约束成1 | 2 | 3意味着升级路径里不会出现level: 0或者level: 99这类让下游规则引擎困惑的值。history数组保存完整时间线后续统计响应时长、升级时长都不需要去查“操作日志表”直接从领域对象上读。这一点对外层代码的影响是任何处理AndonPhase的函数都必须先处理phase为raised的分支否则 TypeScript 编译器直接报错。写测试的时候也只需要针对这四种 phase 枚举输入不会出现“字段组合爆炸”。这里顺带说一个代码评审里常见的问题有人图省事把history声明为ArrayRecordstring, unknown或者初始化一个const events [{}]让 TypeScript 推断成{}[]之后所有写入都在any边缘游走。正确做法是显式声明const events: AndonPhase[] []让数组元素从一开始就被约束。2.2 为什么 Andon 适合事件表而不是只靠 CRUDAndon 系统的业务数据有强烈的审计需求几个小时后复盘一个异常要能说清楚谁在什么时候做了什么。如果只用 CRUD 更新andon_calls.status历史信息就丢了。所以这里我用 append-only 的事件表作为事实来源。它不是完整的 Event Sourcing而是一个中间形态业务表存当前状态事件表存所有历史两者用call_id关联。维度只维护一张当前状态表当前状态表 事件表查询当前报警单表读取快多一次关联量小无感历史审计丢失完整保留数据重放不支持可按 seq 重放统计实现成本最低两张表成本可控实际做的时候我选择的是后者因为 Andon 的下一步必然是数据分析没有历史序列预测性维护无从谈起。下面这种设计对查询和审计都比较友好CREATE TABLE andon_calls ( id text PRIMARY KEY, line_id text NOT NULL, station_id text NOT NULL, type text NOT NULL, status text NOT NULL, created_at timestamptz NOT NULL, updated_at timestamptz NOT NULL ); CREATE TABLE andon_events ( seq bigserial PRIMARY KEY, call_id text NOT NULL REFERENCES andon_calls(id), phase text NOT NULL, payload jsonb NOT NULL, emitted_at timestamptz NOT NULL DEFAULT now() ); CREATE INDEX idx_andonevents_call_time ON andon_events(call_id, seq);两张表的分工是andon_calls给前端列表页查询当前状态andon_events给统计任务和预测性维护提供原始序列。用bigserial作为seq可以保证同一call_id的排序稳定不依赖业务表里的时间字段。注意不要把payload里的时间当排序依据跨系统时钟不一致会产生乱序数据库序列才是可靠顺序。事件 ID 我建议用 ULID 而不是 UUID。ULID 在分布式环境下可以独立生成字典序又大体接近生成时间调试时在日志里能直观看出时间先后。如果要接入 Redis Streams 的分区策略ULID 也可以直接在消费者端参与排序。2.3 用纯函数实现“创建一次 Andon 报警”领域的创建动作不适合直接写在 Controller 里我一般会独立成一个纯函数。它的职责是校验输入、生成事件、返回供存储层落库的对象不直接碰 Redis 和数据库。import { ulid } from ulid; /** 创建一次 Andon 报警只生成领域对象不触碰数据库 */ export function createAndonCall( input: { lineId: string; stationId: string; type: AndonType; now?: string; } ): AndonCall { const now input.now ?? new Date().toISOString(); return { id: ulid(), type: input.type, lineId: input.lineId, stationId: input.stationId, history: [ { phase: raised, at: now, stationId: input.stationId }, ], }; }now参数做成可选是为了测试时传入固定时间避免测试结果依赖系统时钟。这个函数本身没有副作用单测只需要三行断言history长度为 1、phase是raised、id非空。这里不需要在创建时校验工位是否存在那属于应用层职责放在领域纯函数里会让它背上一堆基础设施依赖。如果在同一条产线的同一个工位短时间内对同类型异常重复触发业务上应该合并而不是重复建模。这个去重逻辑我放在 Redis 缓存层做以andon:dup:{stationId}:{type}为 keySETNX设置 5 分钟过期只有第一次调用才继续走创建流程。这样不会有并发窗口也不会拖慢领域层。3. 用 NestJS Redis Streams 搭可扩展 Andon 处理链路3.1 为什么是 Redis Streams 而不是 KafkaNestJS 是 TypeScript 后端服务里模块划分最清晰的一档配合装饰器与依赖注入适合 Andon 这种按业务域拆分的系统。Andon 系统的消息量级比较特殊中小工厂每天几千到几万条事件峰值并发也不高但每条事件的链路很长牵扯响应人通知、升级策略、统计写入。拿 Kafka 跑这个量级运维成本上不划算。Redis Streams 的优势是Redis 本来就在基础设施里字段和消费者组语义直观XREADGROUP可以做到类似 Kafka consumer group 的负载均衡。几种常见方案的对比按我自己的使用体验排列方案持久化消费者组部署成本适用阶段Postgres LISTEN/NOTIFY弱事件不落盘无低只适合内部触发Redis Streams可靠可配置 MAXLEN支持低单工厂、跨产线场景Kafka高可长期保存成熟高多工厂、大数据分析前置Redis Streams 也有做不了的事没有复杂的重平衡协议消费者崩溃后分区重分配需要依赖XCLAIM消息保留时长要手动搭配XTRIM。不过这些对一个 Andon 系统的体量来说不是问题。真正需要迁移到 Kafka 的信号是多工厂统一汇聚每天百万级事件下游同时有 5 个以上消费组在跑分析任务。3.2 生产端以 lineId 为分区键写入 Stream写入 Stream 的 key 我按产线拆andon:stream:{lineId}。这样同一条产线的事件顺序天然聚合在同一个流里统计该产线时直接读单条 Stream 即可。如果只有一个全局 Stream 也可以但后期想按产线隔离消费组或做权限控制时会后悔。/** 推送一条 Andon 事件到 Redis Streams */ async function emitAndonEvent(call: AndonCall): Promisestring { const streamKey andon:stream:${call.lineId}; return redis.xadd( streamKey, MAXLEN, ~, 20000, *, payload, JSON.stringify(call), ); }MAXLEN ~ 20000里的~表示近似裁剪Redis 在内存压力合适时才做清理性能比精确裁剪高。20000 条上限按单产线一天事件量的十倍留余量。写入成功后返回的 ID 形如1710000000000-0前半部分是毫秒时间戳后半部分是序列号。要注意的是多个服务实例同时往同一个 Stream 写入没有问题Redis 单线程保证每条追加有序。所以生产端可以随 NestJS 实例水平扩展不需要额外做锁。3.3 消费端消费者组、幂等与 ack消费端我建议做成独立的 NestJS service用nestjs/schedule的Interval驱动循环而不是在构造函数里开while(true)。后者会让 NestJS 的优雅退出失效。下面是标准消费循环Injectable() export class AndonEventConsumer { private readonly streamKey: string; private readonly group andon-main-group; private readonly consumerName worker-${process.env.HOSTNAME ?? local}; Interval(1000) async poll() { const result await redis.xreadgroup( GROUP, this.group, this.consumerName, COUNT, 20, BLOCK, 200, STREAMS, this.streamKey, , ); if (!result) return; for (const [, entries] of result) { for (const [id, fields] of entries) { const payload fields.find(([k]) k payload)?.[1]; if (!payload) continue; await this.handleEvent(JSON.parse(payload)); await redis.xack(this.streamKey, this.group, id); } } } }这段代码有几个参数需要解释。COUNT 20是单次拉取上限避免一次拉动几千条导致某条处理超时阻塞整个循环。BLOCK 200是阻塞毫秒数NestJS 的Interval(1000)本身已经有 1 秒节奏阻塞时间设太短会空转太长会让关闭变得迟钝。表示只读消费者组里未被投递的消息这是消费者组语义的核心用法。handleEvent里必须做幂等。消费者组机制保证消息至少投递一次但进程崩溃在xack之前就会导致同一消息被重复处理。我建议在 Redis 里放一个处理标记/** 幂等处理 Andon 事件 */ async handleEvent(call: AndonCall) { const dedupKey andon:handled:${call.id}; const ok await redis.set(dedupKey, 1, EX, 86400, NX); if (!ok) return; // 实际业务通知响应人、更新数据库状态、触发升级策略 }用set的NX原子地完成写入EX 86400是去重窗口只要大于重启重放的时间范围就不会出现漏处理。需要说明的是这个幂等标记会随 Redis 重启清空但 24 小时内的事件不会重复投递所以影响很小。如果确实要更稳可以在andon_events表的(call_id, phase)上加唯一索引。关于超时升级不要在内存里setTimeout进程重启后全部丢失。我把它放在一个延迟队列里数据结构选 Redis ZSETkey 为andon:ack:${call.id}score 是截止时间戳const deadline Date.now() 5 * 60 * 1000; await redis.zadd(andon:ack:${call.id}, deadline, call.id);定时任务每秒执行一次zrangebyscore取所有 score 小于当前时间戳的事件逐一执行升级。zadd天然支持同 key 覆盖同一事件重复创建不会产生脏数据。这个方案比定时扫数据库轻也比进程内定时器可靠。4. 部署到 Docker Compose 与 K8s环境变量、健康检查与扩缩容4.1 最小可部署的基础设施与 Docker Compose 编排一个可运行的 Andon 系统最小需要三个组件NestJS API 服务、Redis Streams、Postgres。不要把时序数据库在起步阶段就加进依赖链里事件表在 Postgres 里先跑起来等统计查询出现明显压力时再同步到时序库。我通常让 Postgres 同时承担业务状态表和事件表的存储两三天内的统计查询都可以直接走 SQL。新版 Docker Compose 可以省略version字段直接声明servicesservices: andon-api: build: . environment: NODE_ENV: production PORT: 3000 REDIS_URL: redis://redis:6379/0 DATABASE_URL: postgres://andon:andonpostgres:5432/andon depends_on: redis: condition: service_healthy postgres: condition: service_healthy ports: - 3000:3000 restart: unless-stopped redis: image: redis:7-alpine command: [redis-server, --appendonly, yes, --maxmemory, 512mb] healthcheck: test: [CMD-SHELL, redis-cli ping | grep PONG] interval: 5s timeout: 3s retries: 5depends_on配了condition: service_healthy这要求 Compose v2 和较新的 Docker 版本别用在老旧的 Docker Toolbox 环境。Redis 挂appendonly yes是为了 Streams 的持久化MAXLEN 裁剪后旧事件虽然没了但投递中的消息不会丢。maxmemory 512mb是给 Streams 和幂等标记留的总内存上限按实际事件量适当调大。Postgres 的健康检查建议用pg_isready不要用psql -c select 1后者会在数据库正常但网络抖动时误报失败。API 服务的重启策略设为unless-stopped适合开发环境和单机部署K8s 环境不需要这个字段交给 Deployment 的 restartPolicy 处理。4.2 让 NestJS 暴露聚合健康检查K8s 就绪检查探活会直接打 API 的/health这个端点不能只返回{ status: ok }它应该把外部依赖的状态聚合进去。用nestjs/terminus可以快速实现 Redis 和数据库探活Controller(health) export class HealthController { constructor( private health: HealthCheckService, private redis: RedisHealthIndicator, private db: TypeOrmHealthIndicator, ) {} Get() HealthCheck() check() { return this.health.check([ () this.redis.pingCheck(redis, { timeout: 1500 }), () this.db.pingCheck(database, { timeout: 1500 }), ]); } }liveness 和 readiness 探针要分开设计liveness 只检查进程是否活着readiness 才检查 Redis 和 DB。如果把 DB 检查挂进 liveness数据库抖动时 Pod 会被反复重启启动风暴反而拖垮依赖。我用一个原则readiness 挂全部依赖liveness 检查进程自身即可。K8s Deployment 里的配置片段readinessProbe: httpGet: path: /health port: 3000 initialDelaySeconds: 10 periodSeconds: 10 livenessProbe: httpGet: path: /health/live port: 3000 initialDelaySeconds: 20 periodSeconds: 30periodSeconds不要设太短健康检查本身会消耗 Redis 和 DB 连接频率过高对生产环境是种隐性压力。4.3 水平扩缩容时要同时调的 4 个参数NestJS 服务是无状态的理论上加副本就能水平扩展实际操作中有四个参数要同步调整否则副本增加只会让问题转移。参数位置扩缩容时的调整建议消费者组 consumer 后缀环境变量每副本唯一不能写死Postgres 连接池maxTypeORM/Prisma按副本数 * 10预留Redisxreadgroup的 COUNT消费者代码副本多了单次可以降到 10滚动更新的 maxUnavailableK8s Deployment设 1避免所有消费者同时重连消费者组里的consumerName我通常用${HOSTNAME}K8s 里 Pod 名天然唯一本地开发用local兜底。Postgres 连接池参数在 TypeORM 里对应extra.max很多项目在这一栏写死10扩副本后数据库连接却还占着连接池数在扩容前先按“峰值副本数 × 单副本连接数”算好上限。Redis Streams 在消费者组扩缩容时会自动重新分配条目但要注意同一条 Stream 的多个消费者在投递时是按条目轮询不是按 key 哈希。如果对乱序敏感就把 Stream 拆小到产线粒度。同产线的事件在同一个 Stream 里顺序不会被多个消费者打乱。5. 从 Andon 事件流走向预测性维护指标口径与规则预警5.1 定义 MTBF 和响应超时率的统计口径想要预测故障先把“设备过得怎么样”量化。Andon 事件表里最直接的统计维度是一段时间内每台设备/每个工位的报警次数、平均解决时长、超时响应占比。下面这个查询可以直接跑在andon_events上-- 近7天各工位报警次数 SELECT station_id, count(*) AS alarm_count FROM andon_events WHERE phase raised AND emitted_at now() - interval 7 days GROUP BY station_id ORDER BY alarm_count DESC LIMIT 20;对应到 TypeScript 侧我会定义一个指标结构让规则引擎和报表共用同一份类型export interface StationHealthMetric { stationId: string; windowDays: number; alarmCount: number; avgResolveSeconds: number; escalateRate: number; } export interface PredictedIssue { level: warning | critical; stationId: string; reason: string; }严格来说 MTBF 是“平均故障间隔时间”但 Andon 报警并不等于故障所以我这里用“无报警持续时长”近似等后续接入设备工单数据后再把 Andon 报警与工单关联计算真实 MTBF。口径必须写进统计文档否则不同人看同一张报表会得出不同结论。5.2 在 TypeScript 里写轻量规则引擎预测性维护的第一步不是机器学习而是稳定的规则基线。写一个函数输入最近若干天的统计输出预警对象export function evaluateStationRisk( m: StationHealthMetric ): PredictedIssue | null { const threshold Number(process.env.ALARM_THRESHOLD ?? 5); if (m.alarmCount threshold m.escalateRate 0.3) { return { level: warning, stationId: m.stationId, reason: 近 ${m.windowDays} 天报警 ${m.alarmCount} 次升级率 ${(m.escalateRate * 100).toFixed(0)}%, }; } return null; }阈值5和0.3是启动值建议作为环境变量注入。之后用历史数据反推采集两周事件算报警次数分布的均值与标准差取“均值 3σ”作为阈值。这个办法在规则引擎阶段足够稳定也容易向业务解释。规则触发后不要在内存里存写入andon_predictions表由另一个流程决定是否生成 Andon 报警。这就是“防止出现问题”的落地闭环设备还没坏统计上先露苗头维修团队提前一周收到训练级预警。5.3 验证规则引擎用脚本回放历史事件测试预测规则最直接的方式是从事件表里捞一段历史按时间顺序回放再把规则输出与人工标注对比。一个简单脚本npx ts-node scripts/replay-events.ts --days 30 --station station-l2-01replay-events.ts里做的事就是查询andon_events、按seq升序、逐条喂给evaluateStationRisk、打印触发记录。跑通后把同样的逻辑接成定时任务每天凌晨计算前一天的指标并生成预警规则体系就活了。这个验证步骤里最容易踩的坑是统计口径不一致线上报警计数基于 UTC产线的“一天”基于本地时区两者差 8 小时周报数字对不上。处理方式是所有统计 SQL 统一用emitted_at AT TIME ZONE Asia/Shanghai转成产线本地时间再按天聚合别在业务代码里手动加8 * 3600。生成预警的同时把抑制原因写进andon_predictions.suppress_reason周会复盘时打开统计表对比阈值调节前后的误报率比翻历史预警记录更高效。本文还有配套的精品资源点击获取