项目一:基于 Apache Flink 的实时日志分析系统
主要功能:
实时采集 Web 服务日志(如 Nginx、应用日志等),通过 Flink 接入 Kafka 实现日志的高吞吐传输。
使用 Flink 的 DataStream API 对日志进行实时清洗、字段提取、格式转换等操作。
实现关键字段的聚合统计,例如每分钟请求量、IP 访问频率、错误率统计等,并通过自定义窗口进行时间滑动统计。
构建异常检测模块,对出现异常请求(如大量 5xx 状态码)进行实时检测并发送预警到邮件。
处理后的结果数据输出到 MySQL 和 Elasticsearch,用于后续数据存储与可视化展示。
个人职责:
负责 Flink 作业的编写与优化,设计数据处理流程和业务逻辑。
编写 Kafka 消费逻辑,解析 JSON 日志并封装为自定义数据模型。
实现窗口聚合统计与异常检测模块,使用 processFunction、自定义 watermark 等特性处理乱序数据。
搭建测试环境,参与系统部署与日志调试,确保系统在高并发下的稳定性和低延迟。