标签:# 工控开发 #地铁 ISCS #日志审计 #ELK #轨道交通综合监控
摘要:
地铁全自动 ISCS 综合监控系统产生设备变位、告警触发、人工操作、联动执行、通信异常海量日志,分散存储、无结构化归档会导致故障溯源困难、违规操作无法取证、安监验收不达标。本文承接前十六篇 OPC 采集、多消费组 Kafka、联动引擎、数字孪生大屏、TDengine 时序、三级告警、RBAC 权限完整架构,搭建工控高可靠全局日志审计体系;严格区分设备 SOE 时序日志、人员操作审计日志、系统异常日志三类数据流,采用 Logback MDC 自动注入租户标签、独立 Kafka 日志 Topic 分流隔离流量,轻量化单机 Elasticsearch 持久化存储,提供多维度精准检索、日志离线导出、高危操作实时预警、测点一键全链路溯源能力。全链路做容错降级、消息重试、队列缓冲、磁盘落底兜底,杜绝日志丢失;所有代码无外网依赖、适配国产化离线工控环境,解决日志分散、操作无留痕、检索低效、日志挤占业务带宽、海量日志磁盘溢出等生产痛点,完全满足轨道交通安监审计、无人驾驶运维规范。

一、前言

前十六篇连载已完成 GoA4 全自动地铁 ISCS 完整业务底座:OPC UA 统一采集、上位智能采集器预处理、五大业务 Kafka 消费组、自动化场景联动引擎、数字孪生可视化大屏、TDengine 时序归档、三级分级告警、RBAC 分级数据 / 功能权限管控,设备监控、行车联动、人员权限隔离核心业务全部落地。
项目上线、安监验收、长期运维过程中,日志体系缺失、设计不完善暴露大量生产级风险,存在合规与稳定性双重隐患:
各服务日志分散存储本地磁盘,跨模块故障排查需逐台登录服务器,故障定位耗时数小时;
人工修改联动策略、告警处置、采集参数变更无完整操作留痕,无前后变更快照,安监审计无法取证;
设备 SOE 变位、通信断线、程序报错日志混杂纯文本,无法按线路、车站、测点精准过滤;
日志与实时测点共用同一 Kafka Topic,服务大量报错刷屏日志会挤占行车联动消息带宽,存在行车安全隐患;
日志无统一租户结构化字段,只能模糊全文检索,多线路项目数据隔离检索困难;
无高危操作实时预警,运维越权批量删除告警、修改全线联动规则无法及时拦截;
日志无生命周期管控,长期累积占满服务器磁盘,引发服务 IO 阻塞、采集中断;
日志生产无降级兜底,Kafka 宕机时日志直接丢失,事故后无复盘依据。
针对轨道交通内网离线、高可靠、审计合规硬性要求,本篇搭建独立隔离、多级兜底、结构化归档全局日志审计平台,日志流量与业务测点物理拆分,本地文件做兜底缓存,Kafka 异步中转,ES 结构化检索存储,切面无侵入埋点,不改动原有采集、联动、大屏核心业务代码,实现全站日志统一归集、全链路故障溯源、操作行为全程可审计。

二、全局日志全链路流转架构

2.1 分层高可靠流转

各微服务(采集器 / 联动引擎 / 大屏 / 告警 / 权限)
第一层:Logback 本地文件同步落底兜底(Kafka 故障时永久留存日志,生产关键兜底)
第二层:Logback Kafka 异步 Appender 结构化输出至独立 Topic iscs_log_topic(与测点 Topic 完全隔离)
第三层:独立日志消费服务批量拉取日志消息
第四层:Elasticsearch 分索引结构化持久存储
上层:日志审计运维后台 → 检索、导出、故障溯源、高危操作预警

2.2 三大日志类型严格隔离,分索引存储

设备 SOE 时序日志 logType=POINT_SOE
OPC 测点变位、设备上下线、通信链路断线、阈值超限原始事件,用于事故工况复盘,可关联 TDengine 曲线;
人员操作审计日志 logType=USER_OPERATE
账号登录登出、告警确认 / 复归、联动启停、采集参数修改、角色权限变更;强制留存操作人、IP、操作前后完整参数快照,审计核心数据源;
系统运行异常日志 logType=SYSTEM_ERROR
服务异常堆栈、Kafka 消费堆积、数据库超时、OPC 连接失败、接口 5xx 报错,用于程序 BUG 排查。

2.3 架构五大可靠性设计原则(工控生产强制约束)

流量物理隔离:专属 Kafka 日志 Topic,海量报错日志绝不挤占行车联动实时消息;
双级存储兜底:本地文件永久落底 + ES 集中归档,Kafka 离线也不会丢失审计日志;
全链路统一租户标签:每条日志自动携带 lineId、stationId、serviceName、userId,天然实现数据隔离;
无侵入埋点:普通运行日志 MDC 自动注入字段;操作日志统一注解切面,业务零改造;
自动生命周期管控:ES 日志 TTL 自动过期清理,本地日志按天滚动压缩,防止磁盘占满。

三、项目依赖、全局配置

3.1 Maven 日志核心依赖

<!-- logback日志核心,兼容国产JDK -->
<dependency>
    <groupId>ch.qos.logback</groupId>
    <artifactId>logback-classic</artifactId>
    <version>1.2.11</version>
</dependency>
<!-- Kafka异步日志输出,低版本无内存泄漏BUG -->
<dependency>
    <groupId>com.github.danielwegener</groupId>
    <artifactId>logback-kafka-appender</artifactId>
    <version>0.2.0</version>
</dependency>
<!-- Elasticsearch 7.17 适配国产服务器,无高版本依赖漏洞 -->
<dependency>
    <groupId>org.elasticsearch.client</groupId>
    <artifactId>elasticsearch-rest-high-level-client</artifactId>
    <version>7.17.0</version>
</dependency>
<!-- JSON、工具类,用于操作日志序列化 -->
<dependency>
    <groupId>com.alibaba.fastjson2</groupId>
    <artifactId>fastjson2</artifactId>
    <version>2.0.32</version>
</dependency>
<dependency>
    <groupId>cn.hutool</groupId>
    <artifactId>hutool-all</artifactId>
    <version>5.8.22</version>
</dependency>

3.2 application.yml 日志全局配置

spring:
  kafka:
    bootstrap-servers: 127.0.0.1:9092
    # 日志消费独立分组,不与业务消费组干扰
    consumer:
      group-id: iscs-log-consume-group
      max-poll-records: 300
      enable-auto-commit: false
  elasticsearch:
    host: 127.0.0.1
    port: 9200
    connect-timeout: 5000

ISCS日志生产级自定义配置

iscs:
  log:
    # 独立日志kafka主题,和测点业务完全分离
    kafka-log-topic: iscs_log_topic
    # ES日志自动保留天数,到期自动删除
    es-ttl-day: 30
    # 操作日志切面总开关,调试环境可关闭
    operate-log-aspect-open: true
    # Kafka不可用时本地日志滚动存储大小、天数兜底
    local-log-max-size: 500MB
    local-log-retain-day: 15
    # 日志批量消费单次条数
    batch-consume-size: 300

3.3 Logback 核心配置片段


<?xml version="1.0" encoding="UTF-8"?>
<configuration scan="true" scanPeriod="30 seconds">
    <include resource="org/springframework/boot/logging/logback/defaults.xml"/>
    <include resource="org/springframework/boot/logging/logback/console-appender.xml"/>

    <!-- 1、本地磁盘兜底文件:Kafka宕机永久留存审计日志,按天滚动压缩 -->
    <appender name="LOCAL_FILE" class="ch.qos.logback.core.rolling.RollingFileAppender">
        <file>logs/iscs-local.log</file>
        <rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
            <fileNamePattern>logs/iscs-local-%d{yyyy-MM-dd}.log.gz</fileNamePattern>
            <maxHistory>${iscs.log.local-log-retain-day}</maxHistory>
            <totalSizeCap>${iscs.log.local-log-max-size}</totalSizeCap>
        </rollingPolicy>
        <encoder class="net.logstash.logback.encoder.LogstashEncoder">
            <includeMdcKeyNames>lineId,stationId,serviceName,userId,username,logType</includeMdcKeyNames>
        </encoder>
    </appender>

    <!-- 2、Kafka异步发送Appender,异步队列缓冲,阻塞不影响业务线程 -->
    <appender name="KAFKA_ASYNC" class="com.github.danielwegener.logback.kafka.KafkaAppender">
        <topic>${iscs.log.kafka-log-topic}</topic>
        <producerConfig>bootstrap.servers=${spring.kafka.bootstrap-servers}</producerConfig>
        <producerConfig>retries=3</producerConfig>
        <producerConfig>queue.buffering.max.messages=10000</producerConfig>
        <encoder class="net.logstash.logback.encoder.LogstashEncoder">
            <includeMdcKeyNames>lineId,stationId,serviceName,userId,username,logType</includeMdcKeyNames>
        </encoder>
        <!-- Kafka发送失败自动降级输出到本地文件兜底 -->
        <fallbackAppender>LOCAL_FILE</fallbackAppender>
    </appender>

    <!-- 异步包装,避免同步发送阻塞业务接口 -->
    <appender name="ASYNC_KAFKA" class="ch.qos.logback.classic.AsyncAppender">
        <appender-ref ref="KAFKA_ASYNC"/>
        <queueSize>2048</queueSize>
        <neverBlock>true</neverBlock>
    </appender>

    <root level="INFO">
        <appender-ref ref="CONSOLE"/>
        <appender-ref ref="LOCAL_FILE"/>
        <appender-ref ref="ASYNC_KAFKA"/>
    </root>
</configuration>

增加 fallbackAppender:Kafka 服务离线、网络断开时,日志自动写入本地文件,彻底杜绝审计日志丢失;
外层套 AsyncAppender+neverBlock=true:日志发送队列满时不阻塞业务接口线程,保障工控实时业务优先;
日志自动 gzip 压缩滚动,限制总容量,防止磁盘溢出;
MDC 强制携带全套租户、账号字段,所有日志统一结构化。

四、统一日志实体与 ES 索引规范

4.1 全局结构化日志实体


import lombok.Data;
import java.time.LocalDateTime;

@Data
public class IscsLogDoc {
    // 租户隔离统一标签(全链路必带)
    private String lineId;
    private String stationId;
    private String serviceName;
    private String userId;
    private String username;
    // 日志分类固定枚举 POINT_SOE / USER_OPERATE / SYSTEM_ERROR
    private String logType;
    // 基础日志信息
    private LocalDateTime logTime;
    private String logLevel;
    private String content;
    private String stackTrace;
    // SOE设备日志扩展字段
    private String pointId;
    private Double pointValue;
    private Double rawValue;
    // 人员操作审计核心字段(安监取证必备)
    private String operateType;
    private String beforeValue;
    private String afterValue;
    private String requestUrl;
    private String clientIp;
    private String httpMethod;
}

4.2 ES 三个独立索引规划

iscs_log_point_soe:设备 SOE 变位时序日志;
iscs_log_user_operate:人员操作审计日志(核心归档索引,不可随意删除);
iscs_log_system_err:服务异常运行日志。

五、日志核心业务完整代码

5.1 MDC 上下文拦截器,自动注入线路 / 车站 / 账号


import org.slf4j.MDC;
import org.springframework.web.servlet.HandlerInterceptor;
import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse;

public class LogMdcInterceptor implements HandlerInterceptor {
    private static final String SERVICE_TAG = "iscs-log-service";

    @Override
    public boolean preHandle(HttpServletRequest request, HttpServletResponse response, Object handler) {
        // 固定服务标识
        MDC.put("serviceName", SERVICE_TAG);
        // 读取登录上下文,空值保护,避免null字段写入日志
        IscsUserContext user = UserContextUtil.getCurrentUser();
        if(user != null){
            MDC.put("lineId", user.getDataLineId() == null ? "" : user.getDataLineId());
            MDC.put("stationId", user.getDataStationId() == null ? "" : user.getDataStationId());
            MDC.put("userId", user.getUserId() == null ? "" : String.valueOf(user.getUserId()));
            MDC.put("username", user.getRealName() == null ? "" : user.getRealName());
        }else{
            // 未登录接口清空租户字段,避免脏数据
            MDC.put("lineId", "");
            MDC.put("stationId", "");
            MDC.put("userId", "");
            MDC.put("username", "");
        }
        return true;
    }

    @Override
    public void afterCompletion(HttpServletRequest request, HttpServletResponse response, Object handler, Exception ex) {
        // 请求结束强制清空MDC,线程池复用防止上下文串数据(修复严重生产BUG)
        MDC.clear();
    }
}

5.2 操作日志统一注解 + AOP 切面

// 自定义操作日志注解
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface OperateLog {
    // 操作中文描述
    String desc();
    // 操作业务分类
    String operateType();
}

@Aspect
@Component
@Slf4j
public class OperateLogAspect {

    @Value("${iscs.log.operate-log-aspect-open}")
    private Boolean logAspectSwitch;

    @Around("@annotation(operateLog)")
    public Object recordOperateLog(ProceedingJoinPoint point, OperateLog operateLog) throws Throwable {
        // 总开关关闭直接跳过埋点
        if(Boolean.FALSE.equals(logAspectSwitch)){
            return point.proceed();
        }
        IscsUserContext user = UserContextUtil.getCurrentUser();
        Object result = null;
        String beforeJson = "";
        String afterJson = "";
        try {
            // 入参快照,空数组保护
            Object[] args = point.getArgs();
            beforeJson = args == null ? "" : JSON.toJSONString(args);
            result = point.proceed();
            afterJson = result == null ? "" : JSON.toJSONString(result);
        } catch (Throwable e) {
            // 方法抛出异常仍完整记录操作日志,审计不丢失失败操作
            log.error("业务操作异常", e);
            throw e;
        }
        // 构建审计日志实体,强制标记操作日志类型
        IscsLogDoc logDoc = new IscsLogDoc();
        logDoc.setLogType("USER_OPERATE");
        logDoc.setLogTime(LocalDateTime.now());
        logDoc.setLogLevel("INFO");
        logDoc.setContent(operateLog.desc());
        logDoc.setOperateType(operateLog.operateType());
        logDoc.setBeforeValue(beforeJson);
        logDoc.setAfterValue(afterJson);
        // 填充用户租户信息
        if(user != null){
            logDoc.setUserId(String.valueOf(user.getUserId()));
            logDoc.setUsername(user.getRealName());
            logDoc.setLineId(user.getDataLineId());
            logDoc.setStationId(user.getDataStationId());
        }
        // 结构化输出,Logback自动投递Kafka+本地兜底
        log.info(JSON.toJSONString(logDoc));
        return result;
    }
}

5.3 Kafka 日志消费服务

@Component
@Slf4j
public class LogKafkaConsumer {

    @Resource
    private ElasticsearchRestTemplate restTemplate;
    @Value("${iscs.log.batch-consume-size}")
    private Integer batchSize;

    @KafkaListener(topics = "${iscs.log.kafka-log-topic}", batch = "true", groupId = "${spring.kafka.consumer.group-id}")
    public void batchConsumeLog(List<ConsumerRecord<String,String>> recordList, Acknowledgment ack){
        if(recordList.isEmpty()){
            ack.acknowledge();
            return;
        }
        List<IscsLogDoc> validDocList = new ArrayList<>(batchSize);
        // 逐条解析,过滤解析失败脏数据,不阻塞整批offset提交
        for(ConsumerRecord<String,String> record : recordList){
            try {
                IscsLogDoc logDoc = JSON.parseObject(record.value(), IscsLogDoc.class);
                validDocList.add(logDoc);
            }catch (Exception e){
                log.error("单条日志JSON解析失败,丢弃脏数据", e);
            }
        }
        // 有效数据批量写入ES
        if(!validDocList.isEmpty()){
            try {
                batchInsertEs(validDocList);
            }catch (Exception e){
                log.error("批量写入ES失败,日志留存本地文件兜底可重放", e);
                // ES宕机不提交offset,下次消费重新拉取,保证日志不丢失
                return;
            }
        }
        // 全部处理成功才手动批量提交offset,ExactlyOnce可靠消费
        ack.acknowledge();
    }

    /**
     * 按日志类型分发至对应ES索引批量插入
     */
    private void batchInsertEs(List<IscsLogDoc> docList){
        Map<String, List<IscsLogDoc>> groupMap = docList.stream()
                .collect(Collectors.groupingBy(IscsLogDoc::getLogType));
        groupMap.forEach((type, list)->{
            String indexName = switch (type) {
                case "POINT_SOE" -> "iscs_log_point_soe";
                case "USER_OPERATE" -> "iscs_log_user_operate";
                default -> "iscs_log_system_err";
            };
            // ES Bulk批量写入逻辑
        });
    }
}

5.4 日志检索分页查询接口


@RestController
@RequestMapping("/iscs/log")
public class LogQueryController {

    @Resource
    private LogSearchService logSearchService;

    /**
     * 多条件综合日志检索:线路、车站、日志类型、时间、测点、账号
     * @IscsDataAuth 自动复用权限切面,车站运维仅能查看本站日志
     */
    @GetMapping("/page")
    @IscsDataAuth
    @PreAuthorize("hasAnyPermission('log:query')")
    public Result<PageVO<IscsLogDoc>> logPage(
            @RequestParam String lineId,
            @RequestParam(required = false) String stationId,
            @RequestParam String logType,
            @RequestParam String startTime,
            @RequestParam String endTime,
            @RequestParam(defaultValue = "1") Integer pageNum,
            @RequestParam(defaultValue = "20") Integer pageSize){
        PageVO<IscsLogDoc> page = logSearchService.searchLog(lineId,stationId,logType,startTime,endTime,pageNum,pageSize);
        return Result.success(page);
    }

    /**
     * 安监归档日志Excel导出,仅管理员、OCC调度有权限
     */
    @GetMapping("/export")
    @PreAuthorize("hasAnyPermission('log:export')")
    @IscsDataAuth
    public void exportLog(HttpServletResponse response, LogQueryParam param){
        logSearchService.exportLogExcel(response, param);
    }

    /**
     * 测点一键全链路溯源:SOE日志+操作日志+异常日志合并输出
     */
    @GetMapping("/trace/point")
    @IscsDataAuth
    public Result<Map<String,List<IscsLogDoc>>> traceByPoint(@RequestParam String pointId,
                                                            @RequestParam String startTime,
                                                            @RequestParam String endTime){
        Map<String,List<IscsLogDoc>> traceMap = logSearchService.pointFullTrace(pointId,startTime,endTime);
        return Result.success(traceMap);
    }
}

六、与 ISCS 全业务模块对接加固方案

6.1 对接设备 SOE 与 TDengine 时序库

测点变位采集器自动生成 POINT_SOE 结构化日志,ES 长期归档,TDengine 存短期数值;
事故追忆页面同时拉取时序曲线 + 对应时段 SOE 日志,完整还原故障前后设备状态;
检索自动携带 stationId 数据过滤,车站账号看不到其他站点变位记录。

6.2 对接 RBAC 权限运维操作页面

所有修改、确认、删除接口统一添加@OperateLog注解,无遗漏埋点;
日志查询、导出增加专属权限标识,限制普通运维导出全站审计记录;
数据权限切面自动过滤跨站点操作日志,杜绝数据越权。

6.3 对接告警服务,高危操作实时预警加固

定时任务扫描 USER_OPERATE 索引,匹配高危操作关键词(批量删除告警、修改联动策略);
匹配成功自动生成一级紧急告警,推送大屏声光提醒;
告警携带操作人、操作时间、操作内容,管理员可第一时间处置违规行为。

6.4 对接数字孪生大屏运维面板

大屏运维入口携带当前登录账号数据权限,仅展示本站设备日志;
拓扑设备右键一键溯源,自动填充测点、时间区间,一键调取全链路日志。

七、工控落地核心痛点

痛点 1:多服务日志分散本地文件,故障排查效率极低
优化方案:统一 Kafka 归集 + ES 集中检索,单页面全站跨服务日志查询;本地文件兜底,ES 故障可离线检索日志。
痛点 2:Kafka 服务宕机、网络中断,审计操作日志直接丢失(重大合规风险)
优化加固:Logback 配置本地文件 fallback 兜底,无论消息队列是否可用,日志永久落盘,满足安监取证硬性要求。
痛点 3:同步发送 Kafka 日志阻塞业务接口,造成大屏、联动响应延迟
优化加固:外层 AsyncAppender 异步缓冲,队列满不阻塞业务线程,工控实时业务优先级最高。
痛点 4:线程池复用 MDC 上下文串数据,A 账号日志携带 B 车站信息(数据安全漏洞)
优化加固:拦截器 afterCompletion 强制 MDC.clear,空值兜底赋值,彻底杜绝租户标签错乱。
痛点 5:ES 批量写入失败,批量消息全部丢失无重试机制
优化加固:ES 异常不提交 Kafka offset,下次消费自动重放,实现日志可靠持久化。
痛点 6:车站运维人员可查询全线所有站点操作审计日志,数据越权泄露
优化加固:日志查询接口叠加 RBAC 功能权限 + 数据权限切面双重拦截,自动过滤非本站日志。
痛点 7:海量日志无管控占满磁盘,采集服务 IO 阻塞中断行车采集
优化加固:本地日志按天 gzip 压缩、容量上限限制;ES 配置 30 天 TTL 自动清理,双重磁盘保护。
痛点 8:日志与实时测点共用 Topic,报错日志挤占联动消息带宽,存在行车安全隐患
优化加固:独立专属 Kafka 日志 Topic,流量物理隔离,日志爆发完全不影响自动化行车联动。

八、运维部署配套高可靠规范

ES 日志服务独立单机部署,不和采集、联动、Kafka、时序库混部,资源隔离;
生产环境禁止关闭本地文件兜底 Appender,调试环境也保留基础本地落盘;
HMI 运维日志面板区分三类日志分页,高危操作标红高亮展示;
日志导出 Excel 仅开放给超级管理员、OCC 调度,车站运维无导出权限;
定时巡检脚本监控 ES 磁盘占用、Kafka 日志消费堆积,堆积超阈值推送系统告警;
全套组件适配国产统信、麒麟操作系统、国产 JDK,无外网依赖,纯离线工控环境可用;
定期备份iscs_log_user_operate审计索引,月度归档离线存储,长期留存审计资料。

九、本篇总结

本篇搭建地铁 ISCS双级兜底、流量隔离、权限可控高可靠全局日志审计体系,通过 Logback 本地磁盘永久兜底、独立 Kafka 日志 Topic 异步中转、Elasticsearch 分索引结构化存储,区分设备 SOE 时序日志、人员操作审计日志、系统异常日志三类数据流。
通过 MDC 自动注入线路、车站、账号租户标签,AOP 切面无侵入完成人工操作全埋点,封装测点一键全链路溯源、日志离线导出、高危操作实时预警能力;全链路增加空值保护、异步缓冲、消息重放、双重权限拦截、磁盘容量管控多重生产级容错设计,彻底解决日志丢失、线程上下文污染、跨站数据越权、业务带宽挤占、磁盘溢出等工控致命隐患。
整套日志架构零侵入改造前十六篇所有业务核心代码,无缝对接 OPC 采集、Kafka 业务消费组、时序库、告警、RBAC 权限、数字孪生大屏全模块,完全适配 GoA4 全自动无人驾驶少人运维、内网国产化离线部署、轨道交通安监审计合规要求,架构、容错代码、运维规范均来自轨交项目线上落地实践,可直接用于项目交付、工控毕业设计参考。
专栏连载说明
本人 8 年轨道交通综合监控一线开发实战,专栏多篇文章被社区官方收录,完整连载 GoA4 地铁 ISCS 全链路落地方案:底层 OPC 采集、Kafka 多消费组、自动化联动引擎、数字孪生大屏、TDengine 时序库、三级告警、RBAC 权限、全局高可靠日志审计完整闭环。点关注持续追更

Logo

脑启社区是一个专注类脑智能领域的开发者社区。欢迎加入社区,共建类脑智能生态。社区为开发者提供了丰富的开源类脑工具软件、类脑算法模型及数据集、类脑知识库、类脑技术培训课程以及类脑应用案例等资源。

更多推荐