☰
基于Docker的实时监控系统:Filebeat+Kafka+Flink全链路实战
2026/10/1 11:32:37 网站建设 项目流程

简介:基于Docker的实时监控系统项目资料,完整覆盖大数据实时链路:Filebeat负责日志采集,经Kafka与Zookeeper传输,由Flink完成实时计算,后端SpringBoot提供接口并用Redis缓存,前端Vue配合Ant Design与Echarts完成可视化展示。资源共98个文件,压缩包约329KB,主要包括23个Shell部署脚本、10个Java后端类、6个Scala Flink作业、13个Vue前端组件及XML/JS/配置文件,另附README与授权说明,目录按前端、后端、Flink任务和Linux脚本清晰分区。已有46人学习浏览,适合计算机相关专业学生用于毕业设计、课程设计或初期项目演示。文档中包含在线人数统计、消息量统计、队列统计等典型实时场景,代码均已运行验证,可在其基础上修改扩展,也可直接用于毕设、课设与实验环境搭建,帮助快速理解容器化实时监控系统的工程落地。

1. 基于 Docker 的实时监控系统:单机跑通 Filebeat、Kafka、Flink 全链路

实时监控系统听起来像重型工程,得堆 Hadoop 集群、配专职数据团队、再买一套商业可视化平台才能转起来。但这套基于 Docker 的实时监控系统实际跑通之后,你会发现 Filebeat + Kafka + Zookeeper + Flink 这四件套在单台机器上就能完成从日志采集、消息中转、实时统计到前端展示的完整闭环。它解决的是一个很具体的问题:业务日志进来后,在线人数、消息吞吐、接口调用量、区域分布这些指标,如何按分钟级别更新到网页图表上。源码包里是 SpringBoot 后端、Vue + Ant Design + ECharts 前端,以及六个 Flink 统计模块,还带完整部署文档。适合拿它做大数据课设、毕设演示,或者在团队里快速验证一套实时监控方案。

2. 架构与四层链路:从日志采集到前端渲染,每个组件解决什么问题

2.1 数据流全景:一条用户日志走到图表的完整路径

先说数据流。这套系统的源头是各业务服务产生的日志,通常是 Nginx 访问日志,或者后端应用自己打印的运行日志。Filebeat 以轻量代理的方式部署在日志所在的位置,它不做复杂清洗,只负责把新增的日志行读出来,推给 Kafka。Kafka 在这条链路里承担缓冲层角色:实时监控场景的流量峰值不可控,Flink 的计算速度不一定追得上突发写入,Kafka 的多分区机制把生产者和消费者解耦,数据先堆在队列里,Flink 按自己的节奏消费。

Kafka 的元数据管理交给 Zookeeper,比如 broker 注册、topic 分区的 leader 选举、消费者组协调。然后是 Flink:它消费 Kafka 里指定 topic 的数据,按时间窗口做聚合统计。典型场景是每秒钟有几千条登录或支付日志流进来,Flink 把它们按用户维度或区域维度分桶,在窗口结束时聚合出结果。统计结果不直接写数据库,而是写进 Redis,因为 Redis 读写是亚毫秒级,后端查询接口实时读 Redis 不会给链路造成压力。

后端这一层是 SpringBoot,暴露 REST API,从 Redis 里取统计结果组装成 JSON 返回给前端。前端是 Vue 单页应用,Ant Design 负责页面布局,ECharts 渲染折线图、柱状图、漏斗图。用户看到的每分钟在线人数曲线,本质上是 Flink 每过一分钟把窗口聚合结果更新到 Redis,前端每隔几秒调用一次后端接口刷新图表。

2.2 组件选型:Filebeat 为什么比 Logstash 轻,Flink 为什么比 Spark Streaming 合适

有人会问,日志采集为什么不直接用 Logstash?Logstash 功能确实强大,但它基于 JVM,内存占用通常是 Filebeat 的十倍以上。监控系统采集端要部署在业务机器上,不可能每台机器都放一个吃 1GB 内存的进程。Filebeat 是 Go 写的二进制,运行内存占用只有几十 MB,配置就是 input 和 output 两块。它做不了复杂清洗,但监控场景本身不需要清洗,采集、传输、中转就够了。

实时计算引擎选了 Flink 而不是 Spark Streaming,核心原因是延迟语义不同。Spark Streaming 是微批处理,一批数据的间隔通常设几秒,监控大屏上曲线有明显阶梯感。Flink 是纯流式计算,事件进来就逐条处理,窗口结束立即输出,延迟能压到毫秒级。这套系统做在线人数、消息量这类对延迟敏感的指标,Flink 更合适。再加上 Flink 内置的窗口机制和状态管理,写统计逻辑比自己在 Spark 上搭 DStream 省事很多。

至于 Redis 在 Flink 和 SpringBoot 之间做中转,很多人不理解为什么不直接落 MySQL。实时统计结果的写入频率很高,如果每秒钟写几十次窗口结果,磁盘 IO 和锁竞争会成为瓶颈。Redis 的 Key 过期机制还能顺带解决历史数据清理问题。这一点第 4 章讲 Key 设计时展开。

2.3 源码包结构:flinkdev 六个统计模块怎么对应业务指标

解压后顶层是 real-time-monitoring-system-master,分成三大块:realtime-vue 是前端工程,realtimebackend 是 SpringBoot 后端,flinkdev 是 Flink 实时计算模块。flinkdev 下面按统计指标拆成六个子模块,建议按这个索引去读代码:

模块目录统计指标结果在 Redis 里的 Key 前缀
onlineNumStatistics在线人数stat:online:
msgStatistics消息吞吐量stat:msg:
queueStatistics队列堆积情况stat:queue:
areaStatistics区域用户分布stat:area:
interfaceStatistics接口调用频次stat:interface:
userStatistics用户维度画像stat:user:

我一般按这个顺序读源码:先看 flinkdev 里任意一个统计模块的主类,搞清楚它消费哪个 topic、用什么窗口、结果写到 Redis 哪个 Key。然后看 realtimebackend 的 Controller,看后端从哪些 Key 读数据。最后才碰前端,因为图表只是把接口返回的数组渲染出来,真正的业务逻辑都在 Flink 和 SpringBoot 里。Linux 目录下是部署脚本,images 目录是运行效果截图,系统跑起来以后可以对照截图验证界面是否一致。

3. Docker 部署实战:用 Compose 编排整套环境并提交 Flink 任务

3.1 docker-compose 定义服务:依赖顺序、端口映射与健康检查

我会把 Zookeeper、Kafka、Redis、Flink 的 JobManager 和 TaskManager 统一写进一个 docker-compose.yml。单独 docker run 一个个拉起来太容易漏环境变量,出问题也不好复现。下面是这套系统的基础编排文件。

version: '3.8' services: zookeeper: image: zookeeper:3.6 container_name: monitor-zk restart: always ports: - "2181:2181" environment: ZOO_MY_ID: 1 kafka: image: bitnami/kafka:3.2 container_name: monitor-kafka restart: always ports: - "9092:9092" - "9093:9093" # 宿主机访问 Kafka 的外部端口 depends_on: - zookeeper environment: KAFKA_CFG_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_CFG_LISTENERS: INTERNAL://:9092,EXTERNAL://:9093 KAFKA_CFG_ADVERTISED_LISTENERS: INTERNAL://kafka:9092,EXTERNAL://localhost:9093 KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT KAFKA_CFG_INTER_BROKER_LISTENER_NAME: INTERNAL redis: image: redis:6.2 container_name: monitor-redis restart: always ports: - "6379:6379" flink-jobmanager: image: flink:1.14 container_name: monitor-flink-jm ports: - "8081:8081" # Flink Web UI command: jobmanager environment: - JOB_MANAGER_RPC_ADDRESS=monitor-flink-jm flink-taskmanager: image: flink:1.14 container_name: monitor-flink-tm command: taskmanager depends_on: - flink-jobmanager environment: - JOB_MANAGER_RPC_ADDRESS=monitor-flink-jm

这段配置里最容易出错的是 Kafka 的 listener 配置。INTERNAL 监听 9092 端口,给 Docker 网络内的组件互相访问用;EXTERNAL 监听 9093 端口,给宿主机上的进程用。advertised.listeners 的作用是告诉客户端“你来连我这个地址”,Filebeat 如果在容器内,拿到的地址必须是它实际能访问的地址,否则就会出现能连上端口但立刻断开的现象,这个坑第 5 章单独讲。

compose 文件的保存位置建议放在项目根目录,和 Linux 目录下的部署脚本平级。执行docker-compose up -d以后,用docker-compose ps看容器状态,用docker-compose logs kafka看具体日志。depends_on 只能保证容器创建顺序,不能保证服务内部真正就绪,所以脚本里我会加一个轮询 Zookeeper 端口的等待循环:

for i in {1..30}; do if docker exec monitor-zk zkServer.sh status | grep -q "standalone\|leader"; then echo "zookeeper is ready" break fi sleep 2 done

等 Zookeeper 真正进入可用状态再让 Kafka 启动,能避免 Kafka 客户端反复重试连不上元数据的问题。

3.2 Filebeat 配置:采集路径、自定义字段与 Kafka 输出

Filebeat 是整个链路的第一公里。它读取日志文件,把每一行包装成 JSON 发到 Kafka。下面是一份典型配置。

filebeat.inputs: - type: log enabled: true paths: - /data/logs/*.log fields: source: web # 业务来源标记,Flink 会用它做过滤 env: prod # 环境标记 fields_under_root: true output.kafka: hosts: ["kafka:9092"] topic: "monitor-log" required_acks: 1 compression: gzip max_message_bytes: 1048576

paths 定义了要监视的日志路径,支持通配符,/data/logs/*.log 会把该目录下所有 .log 文件纳入监控。fields 给日志打标签,fields_under_root: true 让标签直接出现在 JSON 顶层而不是嵌套的 fields 字段里,Flink 解析时少一层结构。output.kafka 里 topic 指定写入目标,required_acks: 1 表示 leader 写入成功即返回,兼顾吞吐和可靠性。compression: gzip 压缩网络传输,日志量大时能省不少带宽。

启动 Filebeat 我用容器挂载方式,把配置和日志目录都暴露进容器:

docker run -d --name monitor-filebeat \ -v /data/logs:/data/logs \ -v /data/filebeat/filebeat.yml:/usr/share/filebeat/filebeat.yml:ro \ -v /data/filebeat/data:/usr/share/filebeat/data \ --network host \ docker.elastic.co/beats/filebeat:7.17

--network host 让 Filebeat 直接用宿主机网络栈,它访问 Kafka 时的地址就和宿主机一致。如果 Kafka 跑在 compose 网络里,需要把 hosts 里的地址换成 kafka:9092,并且确保 Filebeat 容器加入了 compose 创建的同一个网络。

调试阶段用 dry-run 模式,不实际发送消息,只打印解析结果:

docker exec -it monitor-filebeat filebeat -e -c /usr/share/filebeat/filebeat.yml -d "publish"

-e 日志输出到终端,-d "publish" 打开 publish 模块调试日志。看到类似 Successfully published events 的输出,说明采集和序列化都正常,再去排查 Kafka 侧。

3.3 提交 Flink 任务:jar 包、入口类、参数与日志验证

flinkdev 工程用 Maven 构建,先打包成可提交的 jar:

cd flinkdev mvn clean package -DskipTests ls target/

然后 docker cp 把 jar 传进 JobManager 容器:

docker cp target/mainflink-1.0.0.jar monitor-flink-jm:/opt/flink/jars/

提交作业用 docker exec 进入容器执行 flink run:

docker exec -it monitor-flink-jm ./bin/flink run \ -d \ -m monitor-flink-jobmanager:8081 \ -c com.realtime.statistics.OnlineNumStatistics \ /opt/flink/jars/mainflink-1.0.0.jar \ --kafka.bootstrap.servers kafka:9092 \ --kafka.topic monitor-log \ --redis.host redis \ --window.seconds 60

-m 指定 JobManager 的 RPC 地址,-c 指定作业入口类。后面四个参数是统计作业自定义的运行参数:Kafka 地址、订阅的 topic、Redis 地址、窗口时长。实际运行时类名以你工程里编译出的为准,不确定就可以先执行jar tf mainflink-1.0.0.jar | grep Statistics看类路径。

提交成功后终端会打印作业 ID,JobManager 的 Web 页面 8081 端口能看到作业状态和吞吐指标。作业如果一直卡在 RESTARTING,先检查 JobManager 和 TaskManager 是不是同一个网络,再确认 -m 地址能否从容器内解析。

4. 后端与前端联通:Redis 里的实时数据如何变成浏览器图表

4.1 SpringBoot 接口:从 Redis 读统计结果的代码骨架

Flink 作业算完的结果实时更新到 Redis,SpringBoot 负责把结果捞出来暴露成接口。这是一个在线人数接口的实现骨架。

@RestController @RequestMapping("/api/monitor") public class OnlineStatController { private final StringRedisTemplate redisTemplate; public OnlineStatController(StringRedisTemplate redisTemplate) { this.redisTemplate = redisTemplate; } @GetMapping("/online") public Result<OnlineStatDTO> onlineStat() { // Flink 按分钟写入,Key 形如 stat:online:202501161430 String key = "stat:online:" + LocalDateTime.now() .format(DateTimeFormatter.ofPattern("yyyyMMddHHmm")); String value = redisTemplate.opsForValue().get(key); OnlineStatDTO dto = new OnlineStatDTO(); dto.setTimestamp(key.substring(key.length() - 12)); dto.setOnlineCount(value == null ? 0 : Long.parseLong(value)); return Result.ok(dto); } }

Key 格式 stat:online:202501161430 的含义是统计类型为 online,后缀是分钟级时间戳。后端取数时用当前时间拼接出 Key,读出字符串再转数字。这里有个细节:Redis 读出的值可能是 null,因为 Flink 第一个窗口还没跑完,所以接口对 null 做了兜底,返回 0 而不是报空指针。

SpringBoot 集成 Redis 引入 spring-boot-starter-data-redis 即可。配置文件里的连接地址,容器化部署时用 redis,宿主机直接跑时用 localhost。项目里 redisTemplate 用的是 StringRedisTemplate,它只处理字符串,正好匹配 Flink 写入的字符串结果,省去自定义序列化器。

4.2 Redis 的 Key 设计与过期策略:统计结果不能无限堆积

监控系统每秒钟都在产生新的窗口数据,如果每个 Key 都不设置过期时间,Redis 内存很快会被历史数据撑爆。Flink 写入时必须给 Key 设置过期时间,保留最近一小时的数据就够图表展示。

// Flink 统计完成后的写 Redis 逻辑 String key = "stat:" + statType + ":" + minuteTime; redisClient.setex(key, 3600, String.valueOf(count));

setex 第二个参数 3600 表示这个 Key 一小时后自动过期。后端查询当前分钟或最近十分钟的数据,一小时以内的历史足够支撑折线图。如果产品要求看全天趋势,把过期时间改成 86400,但内存消耗会线性增长,需要先估算统计项数量和每秒写入次数,再确定保留周期。

另一个容易忽略的问题是过期时间只对单条 Key 生效,一分钟写一条,一小时内存里还是会累计六十条记录。高频指标我一般加离线归档任务,把旧 Key 的数据采样写进 MySQL,Redis 只保留热数据。

4.3 前端 Vue 组件:Axios 轮询加 ECharts 渲染

前端是 Vue 单页应用,Ant Design 负责布局和表格,ECharts 负责图表。核心逻辑是一个定时器循环调用后端接口,把返回的时间序列数据交给 ECharts 的 setOption。

export default { name: 'OnlineChart', data() { return { chartInstance: null, timeWindow: [], countWindow: [], timer: null, }; }, mounted() { this.chartInstance = echarts.init(this.$refs.chartDom); this.startPolling(); }, beforeDestroy() { clearInterval(this.timer); this.chartInstance.dispose(); }, methods: { async startPolling() { this.fetchData(); this.timer = setInterval(this.fetchData, 5000); }, async fetchData() { const { data } = await axios.get('/api/monitor/online/history?minutes=10'); this.timeWindow = data.data.map(it => it.timestamp); this.countWindow = data.data.map(it => it.onlineCount); this.chartInstance.setOption({ xAxis: { type: 'category', data: this.timeWindow }, yAxis: { type: 'value' }, series: [{ type: 'line', data: this.countWindow, smooth: true }], }); }, }, };

轮询间隔 5 秒,窗口取最近 10 分钟数据,图表上总是一条带十个点的折线在滚动。setOption 时 xAxis 和 series 的 data 长度要一致,如果某个时间点接口没返回数据,后端先补 0,前端不做特殊处理,避免图表断点。

前端构建后放 Nginx,反向代理把 /api 请求转发到 SpringBoot。开发环境的跨域问题靠 vue.config.js 里的 proxy 解决:

module.exports = { devServer: { port: 9528, proxy: { '/api': { target: 'http://localhost:8080', changeOrigin: true, }, }, }, };

changeOrigin 必须为 true,否则后端收到的 Host 头是前端域名,某些基于 Host 做校验的过滤器会直接拒掉。pathRewrite 按需配置,后端接口如果本来就在 /api 前缀下可以省略。

5. 避坑排查:Docker 部署这套系统最容易翻车的五个位置

5.1 Docker Desktop 启动失败:虚拟化与 WSL2 问题

现象:Windows 上装完 Docker Desktop 双击启动,几秒后提示 virtualization not supported 或者 failed to start because virtualization is disabled,Docker 引擎一直起不来。

原因:Docker Desktop 的 Windows 容器模式依赖 Hyper-V 或 WSL2 后端,这两者都要求 CPU 的虚拟化扩展在 BIOS 里开启。很多办公电脑出厂时 BIOS 里的虚拟化技术是关闭的,Windows 检测不到虚拟化能力,Docker 自然无法创建虚拟机。

解决:进 BIOS 找到 Intel Virtualization Technology 或 SVM(AMD 平台)开启,保存重启。进入系统后确认“适用于 Linux 的 Windows 子系统”和“虚拟机平台”两个 Windows 功能已勾选,执行wsl --set-default-version 2把默认版本设成 WSL2。如果配置了 WSL2 还是失败,打开任务管理器性能页签,看“虚拟化”一栏是不是“已启用”,不是就说明问题还在 BIOS 层。

5.2 容器间网络不通:服务名解析不了,连接被拒绝

现象:docker-compose up 全部正常,浏览器能打开 Flink Web UI,但 Flink 作业消费 Kafka 时一直报 timeout,日志里出现 Connection refused。

原因:服务之间不在同一个自定义桥接网络里。compose 默认会为项目创建网络,但如果你用单独的 docker run 启动 Kafka,又用 compose 启动 Flink,两个容器各自属于不同网络,容器名在对方网络里解析不到。另外,Kafka 的 advertised.listeners 如果写的是 localhost,Flink 在容器里访问 localhost 指向的是自己,不是 Kafka。

解决:在 compose 文件底部显式声明统一网络,所有服务加入同一个网络。单独用 docker run 启动的容器,加--network指定加入 compose 创建的网络。验证命令是在 Flink 容器里执行docker exec -it monitor-flink-jm ping kafka,能通再提交作业。

5.3 镜像拉取缓慢卡死:Docker Hub 访问不通

现象:docker-compose up 执行到 pulling 阶段,进度条长时间不动,或者报 dial tcp: i/o timeout,等一小时也拉不完 Zookeeper 镜像。

原因:Docker Hub 服务器在海外,国内网络直连不稳定,镜像层又多,大镜像中途断流后 Docker 客户端虽然支持断点续传,但重试逻辑不一定每次都有效。

解决:给 Docker daemon 配置镜像加速源。Linux 下修改 /etc/docker/daemon.json,Windows 下在 Docker Desktop 的 Settings 里改 Docker Engine 配置:

{ "registry-mirrors": [ "https://docker.m.daocloud.io", "https://dockerproxy.com" ] }

保存后重启 Docker 再拉镜像。加速源可用性会随网络环境变化,哪个不通就换哪个,别死磕一个地址。离线内网环境就在有网的机器上先 docker pull,再 docker save 成 tar 传输,目标机器 docker load 导入。

5.4 Kafka 客户端连不上:advertised.listeners 引发的怪问题

现象:Filebeat 或 Flink 能联通 Kafka 的 9092 端口,但真正发送消息时立刻报错,然后自动重试,日志里出现 Connection to node -1 could not be established 或 Disconnected from Kafka。

原因:Kafka 的 broker 把内部地址告诉给了客户端,客户端按这个地址去连接,却发现路由不通。配置里如果只写了KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,这个 kafka 主机名只有 Docker 内部网络能解析。宿主机上的进程拿到 kafka:9092 后尝试解析 kafka 这个域名,解析失败或者走了错误路由,连接建立不起来。

解决:像第 3 章 compose 里那样区分 INTERNAL 和 EXTERNAL 两套 listener。容器内组件用 kafka:9092,宿主机进程用 localhost:9093 或宿主机的真实 IP。客户端用哪个端口连,advertised 里对应的地址就必须是客户端可达的地址。血泪经验:我每次重装这套环境都要在这里折一次,后来干脆把外部访问统一走宿主机 IP:9093,内部组件全部走 9092,两套互不干扰。

5.5 Flink 作业在消费但统计结果不动:窗口和 Watermark 的问题

现象:Flink 作业已提交,Kafka topic 消费位点在前进,但 Redis 里对应统计 Key 一直没更新,网页图表一直是 0。

原因:事件时间窗口需要 Watermark 触发,而 Watermark 从事件里抽取时间戳生成。如果写入 Kafka 的日志没有时间戳字段,代码里却用事件时间语义,Watermark 永远不会推进,窗口也就永远不关闭。另一种可能是窗口时长设得太大,比如 10 分钟窗口,当前数据只攒了 2 分钟,自然看不到输出。

解决:监控统计场景对乱序数据容忍度有限,直接把时间语义改成 ProcessingTime,每 60 秒窗口结束立刻计算输出,看效果最快。也可以给日志补一个 timestamp 字段,在 Flink 里实现 AssignerWithPeriodicWatermarks 再用事件时间。排错时先docker exec -it monitor-redis redis-cli keys 'stat:*'看 Redis 有没有 Key,有 Key 说明 Flink 写成功了,问题在查询接口或前端;没有 Key 再回头查 Flink 日志。

6. 验证与进阶:部署完第一件事是压测数据链路

6.1 三步验证:Kafka 主题、Redis Key、前端图表

部署完成后,我固定按三步验证链路。先看 Kafka 的 topic 有没有数据流入,在 Kafka 容器里执行:

docker exec -it monitor-kafka kafka-console-consumer.sh \ --bootstrap-server localhost:9092 \ --topic monitor-log \ --from-beginning --timeout-ms 5000

有日志行不断输出,说明 Filebeat 到 Kafka 这一段正常。然后看 Redis 里 Flink 写入的统计结果:

docker exec -it monitor-redis redis-cli keys 'stat:*'

最后打开前端页面,看图表是否随轮询刷新。三步走完,整条链路才算验收通过。

6.2 制造测试日志:模拟访问验证端到端

没有现成业务日志时,手动造一条 JSON 日志触发整个链路:

echo "{\"timestamp\":\"$(date +%s)\",\"user_id\":10001,\"area\":\"beijing\",\"action\":\"login\"}" >> /data/logs/web.log

Filebeat 采集这条日志推到 Kafka,Flink 消费后计入窗口的在线人数和区域分布。等一分钟再看 Redis 和前端,数据应该已经反映在图表上。每一条日志走完这个流程,就能把每个环节的可疑点一一排除。

后面接真实生产环境,重点放在水平扩展上:Kafka 的 topic 分区数要匹配 Flink 并行度,TaskManager 数量按数据量动态调整,Redis 考虑哨兵集群模式而不是单点。从那以后,我每次给别人交付这套系统,都会强制走一遍上面的三步验证加一条测试日志,确认 Kafka 消费、Redis 写入、前端渲染都活着,才敢说环境是好的。这套资源拿回去以后,建议先按第 3 章的 compose 把环境拉起来,再逐步替换成自己业务的字段,比从零开始写要快得多。希望帮到你。

本文还有配套的精品资源,点击获取

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

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

立即咨询