plantpulse-cep(复合事件处理)
角色
Complex Event Processing 引擎。实时检测多个传感器事件的时间和逻辑模式,执行报警、Flow 触发和自动化。
| 项目 | 值 |
|---|---|
| 模块名 | plantpulse-cep |
| 安装路径 | /opt/kopens/plantpulse-platform/plantpulse-cep/ |
| 端口 | 7400(HTTP)· 7401(TLS) |
| 引擎 | Esper CEP(内置 Tomcat 上的 Web 应用) |
pd 服务 | cep — 仅限 MASTER |
| API 认证 | X-API-Key: <PP_CEP_API_KEY> |
| 控制台登录 | PP_CEP_WEB_USER / PP_CEP_WEB_PASSWORD — 为空则仅关闭控制台 |
使用示例
"温度连续 5 秒以上超过 80°C,且同时压力上升超过平均值的 15%,则发送危险报警"
用 EQL(Event Query Language)定义这类复合条件,CEP 引擎会实时匹配并发布结果。
架构
目录结构
plantpulse-cep/
├── config/plantpulse-cep.properties # 생성물 ← /etc/kopens/conf/plantpulse-cep.properties.template
├── bin/start.sh stop.sh
├── logs/
└── server/ # 내장 Tomcat — webapps/ROOT.war · logs/catalina.<date>.log
主要配置
来自模板 /etc/kopens/conf/plantpulse-cep.properties.template 的关键项目。占位符在渲染时填充。
# 캐시 · HA 스냅샷 (Valkey)
cache.host=${PP_REDIS_HOST}
cache.port=${PP_REDIS_PORT}
cache.password=${PP_REDIS_PASSWORD}
cep.ha.enabled = true # 재시작 때 엔진 상태(윈도 · 테이블) 보존
cep.ha.snapshot.interval_ms = 1000
# 이벤트 저장 (Cassandra)
storage.host=${PP_CASSANDRA_HOST}
storage.port=${PP_CASSANDRA_PORT:9042}
storage.keyspace=${PP_KEYSPACE}
cep.storage.ttl_seconds=864000 # 10일
# 처리 스레드
cep.consumer.thread_count=16
cep.stream.shard_count=16
cep.async.max_pool_size=256
# 인증
cep.api.key=${PP_CEP_API_KEY} # 관리 API 의 X-API-Key
cep.web.user=${PP_CEP_WEB_USER} # 콘솔 로그인
cep.web.password=${PP_CEP_WEB_PASSWORD}
修改值的过程请参阅 如何修改配置 — 模板,API 密钥轮换请参阅 修改密码 · API 密钥。
认证
管理 API 通过 X-API-Key 标头进行认证。必须与服务器控制台的 PP_CEP_API_KEY 相同,只修改其中一个会立即返回 401 — passwd.sh PP_CEP_API_KEY 同时修改两处。
curl -H "X-API-Key: ${PP_CEP_API_KEY}" http://127.0.0.1:7400/api/v1/status
EQL 示例
在服务器控制台的 AUTOMATION > CEP 菜单或 CEP 控制台 https://<server-ip>:7401/ 中注册并部署 EQL。
-- 5초 윈도에서 평균 온도가 80도 초과
@Name('high-temp')
SELECT tagId, avg(value) AS avg_temp
FROM TagPoint(tagId = 'TEMP_REACTOR_01').win:time(5 sec)
GROUP BY tagId
HAVING avg(value) > 80
OUTPUT EVERY 1 sec;
-- 두 태그의 시간 정렬 패턴
@Name('temp-pressure-correlation')
SELECT a.tagId AS temp_id, b.tagId AS pressure_id
FROM pattern [
every a = TagPoint(tagId = 'TEMP_01' AND value > 80)
-> b = TagPoint(tagId = 'PRESS_01' AND value > 100) WHERE timer:within(10 sec)
];
语法参考 EQL 帮助。
操作命令
docker exec plantpulse-datalake pd status cep
docker exec plantpulse-datalake pd restart cep
docker exec plantpulse-datalake pd logs --lines 100 cep
# 표준 헬스체크 (익명, readiness — 엔진 + 컨슈머 + 복구 완료 · Valkey 연결 시에만 UP)
curl -fsS http://127.0.0.1:7400/api/health
# 200 {"status":"UP","service":"plantpulse-cep-server-web",…,"checks":{"engine":true,"redis":true}}
# 503 {"status":"STARTING"} 기동 중 / 503 {"status":"DEGRADED"} 의존성 다운
常见问题
| 现象 | 原因 | 处理方案 |
|---|---|---|
| 规则不工作 | 未部署 | 在控制台中点击「部署」 |
| 处理延迟 | 窗口内存不足 | 减小 EQL 窗口或增大堆 |
| Kafka lag 增加 | 处理线程不足 | 增加 cep.consumer.thread_count(模板中配置) |
| API 401 | API 密钥不匹配 | passwd.sh PP_CEP_API_KEY |
| 控制台登录被拒 | PP_CEP_WEB_PASSWORD 为空 | Web 控制台登录账户 |