尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Flink实时监控系统:Kafka+Filebeat+Redis全链路调优实战

Flink实时监控系统:Kafka+Filebeat+Redis全链路调优实战 简介本资源是一套基于Docker构建的实时监控系统完整工程包面向计算机相关专业学生、教师及初级大数据开发人员解决多组件协同下的日志采集、流式处理与可视化监控落地问题适用于毕业设计、课程设计及Flink实战入门。压缩包共98个文件含23个Shell脚本用于Docker环境编排与服务启停、13个Vue前端页面基于Ant Design与Echarts实现动态图表展示、10个Java后端模块SpringBoot核心业务逻辑、6个Scala Flink作业代码含用户统计、消息队列、区域分析等实时计算任务以及YML/Properties配置、XML依赖定义和Markdown文档等整体仅329KB轻量易部署。已有46人学习下载资源附带详细技术文档、全链路组件说明及已验证可运行的源码涵盖Filebeat→Kafka→Zookeeper→Flink→Redis→Vue完整数据流结构清晰、注释充分特别适合从零理解实时数仓监控架构并快速二次开发。1. 这不是又一个“Docker 跑通就完事”的监控 Demo它用 Flink 做毫秒级窗口聚合Kafka 吞吐压测实测 120 万条/分钟Filebeat 采集日志延迟稳定在 80ms 内——适合拿去交毕设、跑课设、搭企业级演示环境的完整闭环系统你肯定见过太多标着“实时监控”的 GitHub 项目Docker Compose 一跑前端页面弹出几个假数据柱状图README 里写着“支持高并发”但连 Kafka Producer 都是用for i in {1..100}; do echo test$i | kafka-console-producer...硬塞进去的。这种项目答辩时老师问一句“如果每秒涌入 5000 条 Nginx access 日志Flink 的 watermark 怎么设Kafka 分区数和消费者组 offset 提交策略怎么配”当场哑火。而这个资源不一样——它从 Filebeat 的close_inactive: 5m和harvester_buffer_size: 16384开始抠起Zookeeper 配了tickTime2000initLimit10syncLimit5三重心跳保障Kafka Server.properties 显式禁用auto.create.topics.enabletrueFlink Job 用TumblingEventTimeWindows.of(Time.seconds(10))搭配AllowedLateness.of(Time.seconds(5))处理乱序后端 SpringBoot 的Async线程池大小按 CPU 核数 × 2.5 动态计算Redis 缓存 key 全部带业务前缀过期时间如cache:online:user:20240520:12h。它不是教你怎么装 Docker Desktop而是告诉你当docker-compose up -d执行完docker ps列出 7 个容器后下一步该kubectl get pods -n monitoring吗不是立刻curl -X POST http://localhost:8080/api/v1/metrics/trigger-simulate拉起 10 万条模拟日志流然后盯住 Flink Web UI 的numRecordsInPerSecond指标是否稳定在 1650±30——这才是真实场景下“实时监控”该有的肌肉记忆。如果你正卡在毕设开题要画技术架构图、课设要交可运行 demo、或者需要一套能直接嵌入企业内部培训的 Flink 实战案例这份资料就是你不用再拼凑 5 个 GitHub 仓库、不用重写 3 个配置文件、不用反复调试 Zookeeper Session Timeout 的那个“最后一块拼图”。2. 从零启动用 docker-compose.yml 拉起全链路组件关键参数逐行解析与本地适配技巧这套系统之所以能“下载即跑”核心在于其docker-compose.yml不是简单堆砌镜像而是针对开发机资源做了精细裁剪并预埋了生产级调优参数。下面我带你一行行拆解重点不是“怎么复制粘贴”而是理解每个字段背后的取舍逻辑——比如为什么 Kafka 的KAFKA_ADVERTISED_LISTENERS必须写成PLAINTEXT://host.docker.internal:9092而不是localhost为什么 Flink 的jobmanager.memory.process.size设为1600m而不是2g。2.1 docker-compose.yml 核心服务定义与参数含义我们先看最关键的services区块。注意这不是标准模板所有参数都经过本机Mac M1 Pro / Windows WSL2 / Ubuntu 22.04实测验证不要直接套用官方镜像默认值。version: 3.8 services: zookeeper: image: confluentinc/cp-zookeeper:7.3.2 environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 ZOOKEEPER_INIT_LIMIT: 10 ZOOKEEPER_SYNC_LIMIT: 5 ZOOKEEPER_SERVER_ID: 1 ZOOKEEPER_SERVERS: 0.0.0.0:2888:3888;localhost:2888:3888 ports: - 2181:2181 networks: - monitoring-net kafka: image: confluentinc/cp-kafka:7.3.2 depends_on: - zookeeper environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_LISTENERS: PLAINTEXT://:9092 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://host.docker.internal:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 KAFKA_AUTO_CREATE_TOPICS_ENABLE: false # 关键禁用自动建 Topic KAFKA_LOG_RETENTION_HOURS: 1 ports: - 9092:9092 networks: - monitoring-net filebeat: image: docker.elastic.co/beats/filebeat:8.10.2 volumes: - ./logs:/usr/share/filebeat/logs:ro - ./filebeat.yml:/usr/share/filebeat/filebeat.yml:ro - /var/lib/docker/containers:/var/lib/docker/containers:ro depends_on: - kafka networks: - monitoring-net flink-jobmanager: image: flink:1.17.1-scala_2.12-java11 environment: FLINK_PROPERTIES: | jobmanager.rpc.address: flink-jobmanager taskmanager.numberOfTaskSlots: 4 parallelism.default: 2 state.backend: filesystem state.checkpoints.dir: file:///opt/flink/checkpoints state.savepoints.dir: file:///opt/flink/savepoints web.upload.dir: /opt/flink/usrlib jobmanager.memory.process.size: 1600m # 关键避免 OOM taskmanager.memory.process.size: 2048m ports: - 8081:8081 command: jobmanager volumes: - ./flink-jobs:/opt/flink/usrlib:ro - ./flink-checkpoints:/opt/flink/checkpoints - ./flink-savepoints:/opt/flink/savepoints networks: - monitoring-net flink-taskmanager: image: flink:1.17.1-scala_2.12-java11 depends_on: - flink-jobmanager environment: FLINK_PROPERTIES: | jobmanager.rpc.address: flink-jobmanager taskmanager.numberOfTaskSlots: 4 parallelism.default: 2 state.backend: filesystem state.checkpoints.dir: file:///opt/flink/checkpoints state.savepoints.dir: file:///opt/flink/savepoints jobmanager.memory.process.size: 1600m taskmanager.memory.process.size: 2048m command: taskmanager volumes: - ./flink-jobs:/opt/flink/usrlib:ro - ./flink-checkpoints:/opt/flink/checkpoints - ./flink-savepoints:/opt/flink/savepoints networks: - monitoring-net redis: image: redis:7.0-alpine command: redis-server --appendonly yes --maxmemory 512mb --maxmemory-policy allkeys-lru ports: - 6379:6379 networks: - monitoring-net backend: build: ./realtimebackend environment: SPRING_PROFILES_ACTIVE: docker REDIS_HOST: redis REDIS_PORT: 6379 KAFKA_BOOTSTRAP_SERVERS: kafka:9092 FLINK_JOBMANAGER_HOST: flink-jobmanager FLINK_JOBMANAGER_PORT: 6123 depends_on: - redis - kafka - flink-jobmanager ports: - 8080:8080 networks: - monitoring-net frontend: build: ./realtime-vue ports: - 8088:80 depends_on: - backend networks: - monitoring-net提示host.docker.internal是关键破局点在 Windows/macOS 上Docker Desktop 自动注入host.docker.internal指向宿主机但在 Linux如 Ubuntu上需手动添加extra_hosts: - host.docker.internal:host-gateway到kafka服务下否则 Filebeat 无法反向连接宿主机的 Kafka。这是新手最常翻车的第一步——别急着查 Kafka 连接超时先docker exec -it kafka-container ping host.docker.internal看通不通。2.2 Filebeat 日志采集配置为什么close_inactive必须设为5mFilebeat 的filebeat.yml不是默认模板它针对监控日志的“短生命周期、高频率写入”特性做了专项优化。核心参数如下filebeat.inputs: - type: filestream enabled: true paths: - /usr/share/filebeat/logs/*.log close_inactive: 5m # 关键避免小文件频繁 reopen 导致 Kafka 生产者阻塞 harvester_buffer_size: 16384 # 单次读取缓冲区提升吞吐 scan_frequency: 10s # 每 10 秒扫描新文件平衡实时性与 CPU fields: log_type: access_log output.kafka: hosts: [kafka:9092] topic: monitoring-logs partition.round_robin: reachable_only: false required_acks: 1 # 不用 -1避免网络抖动导致采集停滞 compression: gzip max_message_bytes: 1000000 # 1MB匹配 Kafka server.properties 的 message.max.bytes参数逻辑说明close_inactive: 5m当一个日志文件 5 分钟内无新内容写入Filebeat 主动关闭该文件句柄。若设为1s很多教程误抄会导致每秒创建/销毁数千个文件句柄Kafka Producer 因频繁send()调用而 CPU 占用飙升至 90%最终numRecordsOutPerSecond断崖下跌。实测5m是吞吐与资源消耗的黄金平衡点。harvester_buffer_size: 16384默认1638416KB已足够增大到65536反而因内存拷贝增加延迟减小到4096会因频繁系统调用降低吞吐。required_acks: 1只要 Leader 副本写入成功即返回不等 ISR 同步完成。监控场景允许极低概率丢日志0.001%但绝不能容忍采集延迟 1s——这是用一致性换实时性的明确取舍。2.3 Kafka Topic 创建脚本为什么必须显式指定分区数与副本因子光靠docker-compose up启动 Kafka 并不能自动创建 Topic。项目中scripts/create-topics.sh是刚需内容如下#!/bin/bash # 进入 Kafka 容器执行 Topic 创建 docker exec -it kafka kafka-topics.sh \ --create \ --bootstrap-server localhost:9092 \ --topic monitoring-logs \ --partitions 6 \ # 关键分区数 Flink 并行度 × 1.5预留扩容 --replication-factor 1 \ --config retention.ms3600000 \ --config segment.bytes1073741824 docker exec -it kafka kafka-topics.sh \ --create \ --bootstrap-server localhost:9092 \ --topic flink-results \ --partitions 3 \ --replication-factor 1为什么分区数必须是6Flink Job 的并行度设为2见flink-jobs/src/main/java/com/example/RealTimeJob.java但 Kafka 分区数必须 ≥ Flink 并行度否则部分 TaskManager 无法消费数据。设为6是为后续横向扩展 Flink TaskManager 预留空间6 ÷ 2 3单节点可扩至 3 个 TM。若设为2后期加机器时需kafka-reassign-partitions.sh重分配分区耗时且易出错。这是血泪经验答辩前夜扩容结果kafka-topics.sh --describe显示Under Replicated Partitions: 0但kafka-consumer-groups.sh --group flink-group --describe却显示LAG持续增长——根源就是分区数不足新加入的 TM 拿不到分区分配。3. Flink 实时计算核心从 DataStream API 到状态后端配置10 行代码实现 10 秒在线人数统计Flink 是这套系统的“大脑”所有实时指标在线人数、接口调用量、消息队列积压都由它计算。项目中的flink-jobs模块不是玩具代码而是严格遵循 Flink 最佳实践编写的生产级 Job。我们以最典型的OnlineNumStatistics为例拆解其设计哲学与避坑细节。3.1 核心统计逻辑EventTime 窗口 允许延迟 状态后端OnlineNumStatistics.java的核心逻辑仅 10 行但每行都直击实时计算痛点public class OnlineNumStatistics { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); // 强制使用事件时间 env.enableCheckpointing(5000); // 5秒一次 Checkpoint env.setStateBackend(new FsStateBackend(file:///opt/flink/checkpoints)); // 文件系统状态后端 DataStreamString source env.addSource(new FlinkKafkaConsumer(monitoring-logs, new SimpleStringSchema(), kafkaProps)); // 关键提取事件时间戳并设置 Watermark DataStreamLogEvent logStream source .map(line - JSON.parseObject(line, LogEvent.class)) .assignTimestampsAndWatermarks( WatermarkStrategy.LogEventforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) - event.getTimestamp()) // 从日志 JSON 提取时间戳 ); // 关键10秒滚动窗口 允许5秒延迟 DataStreamTuple2String, Long onlineCount logStream .filter(event - user_login.equals(event.getEventType())) .keyBy(event - event.getUserId()) .window(TumblingEventTimeWindows.of(Time.seconds(10))) .allowedLateness(Time.seconds(5)) // 允许迟到5秒的数据参与计算 .process(new ProcessWindowFunctionLogEvent, Tuple2String, Long, String, TimeWindow() { Override public void process(String userId, Context context, IterableLogEvent elements, CollectorTuple2String, Long out) { long count StreamSupport.stream(elements.spliterator(), false).count(); out.collect(Tuple2.of(userId, count)); } }); // 输出到 Kafka onlineCount.map(tuple - JSON.toJSONString(tuple)) .addSink(new FlinkKafkaProducer(flink-results, new SimpleStringSchema(), kafkaProps)); env.execute(Online User Count Job); } }关键参数解读forBoundedOutOfOrderness(Duration.ofSeconds(5))假设日志最大乱序 5 秒。若设为0则任何晚于当前 Watermark 的数据都会被丢弃设为10s则窗口触发延迟 10 秒违背“实时”初衷。5s是经压测确定的平衡值——Nginx 日志从生成到 Filebeat 采集平均耗时 80msKafka 传输 200msFlink 网络延迟 100ms总延迟 500ms留 4.5s 余量应对突发抖动。allowedLateness(Time.seconds(5))允许迟到 5 秒的数据触发窗口重计算。没有它凌晨 00:00:05 的登录日志因服务器时钟漂移晚到会被丢弃导致在线人数统计失真。开启后Flink 会为每个窗口维护一个“迟到数据缓存区”成本可控。FsStateBackend使用文件系统而非 RocksDB作为状态后端。原因本项目状态量小用户 ID 计数RocksDB 的 JNI 调用开销反而更高且file:///opt/flink/checkpoints已挂载到宿主机目录Checkpoint 持久化可靠。若换成千亿级用户画像场景则必须切 RocksDB。3.2 Flink Web UI 监控要点3 个必看指标定位性能瓶颈启动flink-jobmanager后访问http://localhost:8081进入 Web UI。不要只盯着“Job Status”绿灯以下 3 个指标才是判断系统健康的核心指标位置指标名称健康阈值异常现象根本原因Job Overview → MetricsnumRecordsInPerSecond≥ 1500 (模拟负载下) 500 且持续下降Kafka Producer 阻塞Filebeat 配置不当或磁盘 IO 瓶颈Task Managers → [TM] → Metricsbuffers.outputQueueLength 100 500 且持续上涨TaskManager 输出缓冲区满下游 Kafka 吞吐不足检查 Kafka 分区数、网络带宽Job → [Operator] → Metricslatency(P95) 200ms 1000msWindow 计算逻辑复杂或状态访问慢检查ProcessWindowFunction是否含 DB 查询注意Latency 指标必须开启默认 Flink 不收集延迟指标。需在flink-conf.yaml中添加metrics.latency.interval: 3000030秒采样metrics.latency.granularity: operator算子级粒度否则latency列为空你永远不知道是 Kafka 慢还是 Flink 计算慢。4. SpringBoot 后端与 Vue 前端联调REST API 设计、Redis 缓存穿透防护、ECharts 动态刷新后端realtimebackend和前端realtime-vue的交互不是简单的axios.get(/api/stats)而是围绕“实时性”与“抗压性”深度协同。这里没有“万能通用方案”所有设计都服务于一个目标当 Kafka 每秒涌入 1 万条flink-results消息时前端图表每 2 秒刷新一次后端 API 响应时间仍稳定在 15ms 内。4.1 后端 API 设计为什么/api/v1/metrics/online返回的是 Redis Hash 而非实时 SQL 查询OnlineController.java中的关键接口RestController RequestMapping(/api/v1/metrics) public class OnlineController { Autowired private RedisTemplateString, Object redisTemplate; // 关键不查数据库直接从 Redis Hash 读取 GetMapping(/online) public ResponseEntityMapString, Object getOnlineStats() { // Redis Key: stats:online:20240520 (日期分区) String key stats:online: LocalDate.now().toString(); HashOperationsString, String, Object hashOps redisTemplate.opsForHash(); MapString, Object data hashOps.entries(key); return ResponseEntity.ok(data); } // 关键Flink 结果写入 Redis 的方式 KafkaListener(topics flink-results, groupId backend-group) public void listenFlinkResults(String message) { try { JSONObject json JSON.parseObject(message); String userId json.getString(f0); // Tuple2.f0 是 userId Long count json.getLong(f1); // Tuple2.f1 是 count String key stats:online: LocalDate.now().toString(); // 使用 HINCRBY 原子操作避免并发覆盖 redisTemplate.opsForHash().increment(key, userId, count); // 设置过期时间防止冷数据堆积 redisTemplate.expire(key, Duration.ofHours(24)); } catch (Exception e) { log.error(Failed to process flink result, e); } } }设计逻辑绝不走 MySQL实时监控指标要求亚秒级响应MySQL 即使加索引10 万行GROUP BY user_id也需 200ms。Redis Hash 的HGETALL是 O(N) 但 N≤1000当日活跃用户实测 5ms。Key 设计带日期分区stats:online:20240520。避免单个 Key 过大如stats:online:all导致 Redis RDB 持久化阻塞主线程。每日新建 Key旧 Key 自动过期。HINCRBY原子操作Flink 可能并发写入同一userIdHSET会覆盖HINCRBY则累加确保计数准确。这是解决“缓存与 DB 一致性”的最简方案——既然指标本身就是近似值10秒窗口累加比强一致更合理。4.2 Vue 前端 ECharts 配置如何让折线图每 2 秒平滑追加新点而不卡顿realtime-vue/src/views/OnlineChart.vue的核心逻辑export default { data() { return { chart: null, option: { tooltip: { trigger: axis }, grid: { left: 3%, right: 4%, bottom: 3%, containLabel: true }, xAxis: { type: category, data: [] }, // 时间轴动态更新 yAxis: { type: value }, series: [{ name: 在线人数, type: line, smooth: true, symbol: none, // 关键去掉标记点减少渲染压力 data: [] }] } } }, mounted() { this.initChart() this.startPolling() }, methods: { initChart() { this.chart echarts.init(this.$refs.chartDom) this.chart.setOption(this.option) // 关键启用 Canvas 渲染禁用 SVG大数据量下 SVG 卡顿 this.chart.setOption({ renderer: canvas }) }, startPolling() { // 关键使用 setTimeout 而非 setInterval避免请求堆积 const poll () { axios.get(/api/v1/metrics/online) .then(res { const now new Date().toLocaleTimeString([], {hour: 2-digit, minute:2-digit}) // 只保留最近 60 个点30 分钟 * 2秒/点 if (this.option.xAxis.data.length 60) { this.option.xAxis.data.shift() this.option.series[0].data.shift() } this.option.xAxis.data.push(now) this.option.series[0].data.push(res.data.total || 0) this.chart.setOption(this.option, true) // true 表示不合并强制重绘 }) .finally(() setTimeout(poll, 2000)) // 下次请求在本次完成后 2 秒发起 } poll() } } }性能关键点symbol: noneECharts 默认在每个数据点画圆圈100 个点就是 100 个 DOM 元素。设为none后纯线条渲染CPU 占用下降 70%。renderer: canvasVue 项目默认用 SVG但 SVG 是 DOM 操作数据量大时重排重绘极慢Canvas 是位图绘制性能碾压 SVG。setTimeout替代setInterval若某次 API 请求耗时 3 秒setInterval会立即触发下一次导致请求队列堆积、内存暴涨。setTimeout确保“上一次完成才开始下一次”稳如老狗。5. 避坑指南5 个真实踩过的坑从 Docker 网络不通到 Flink Checkpoint 失败这套系统在 M1 Mac、Windows 11 WSL2、Ubuntu 22.04 上均实测通过但过程中踩过不少深坑。以下是 5 个最高频、最隐蔽、最容易让你花 3 小时 debug 却只改一行配置的问题。每条都按“现象 → 原因 → 解决”给出可立即执行的方案。5.1 现象docker-compose up后filebeat容器日志疯狂刷Failed to connect to broker但kafka容器docker logs显示正常启动原因Filebeat 容器内 DNS 解析失败无法将kafka服务名解析为 IP。根本原因是 Docker 默认 DNS 服务器8.8.8.8在某些网络环境下被干扰而docker-compose的自定义网络monitoring-net未配置 DNS。解决在docker-compose.yml的filebeat服务下添加dns配置并重启filebeat: # ... 其他配置 dns: - 127.0.0.11 # Docker 内置 DNS - 8.8.8.8然后执行docker-compose down docker-compose up -d # 验证docker exec -it filebeat nslookup kafka5.2 现象Flink Web UI 显示 Job 运行中但flink-resultsTopic 无任何消息backend日志报Connection refused连接 Kafka原因backend服务的application-docker.yml中spring.kafka.bootstrap-servers配置为localhost:9092但localhost在容器内指向容器自身而非宿主机的 Kafka。Docker 容器间通信必须用服务名。解决修改realtimebackend/src/main/resources/application-docker.ymlspring: kafka: bootstrap-servers: kafka:9092 # 改为服务名 kafka不是 localhost重新构建镜像cd realtimebackend docker build -t realtime-backend .5.3 现象前端http://localhost:8088页面空白浏览器控制台报Failed to fetch但curl http://localhost:8080/api/v1/metrics/online返回正常原因Vue 开发服务器vue-cli-service serve默认开启跨域代理但生产构建npm run build后静态文件由 Nginx 托管而realtime-vue/nginx.conf中proxy_pass指向了http://backend:8080但backend服务名在 Nginx 容器内不可达Nginx 未加入monitoring-net网络。解决修改realtime-vue/Dockerfile在COPY nginx.conf /etc/nginx/conf.d/default.conf后添加网络声明FROM nginx:alpine COPY nginx.conf /etc/nginx/conf.d/default.conf COPY dist/ /usr/share/nginx/html/ # 关键让 Nginx 容器加入 monitoring-net # 此行由 docker-compose 自动处理无需改 Dockerfile只需确保 docker-compose.yml 中 frontend 有 networks确认docker-compose.yml中frontend服务包含frontend: # ... 其他 networks: - monitoring-net # 必须有5.4 现象Flink Job 运行 10 分钟后自动 CancelWeb UI 报Checkpoint failed: Could not materialize checkpointflink-checkpoints目录为空原因flink-jobmanager和flink-taskmanager的state.checkpoints.dir都指向file:///opt/flink/checkpoints但该路径在两个容器内是隔离的。JobManager 写入的 Checkpoint 文件TaskManager 根本读不到导致状态恢复失败。解决将 Checkpoint 目录挂载为共享卷。修改docker-compose.ymlvolumes: flink-checkpoints: # 新增命名卷 services: flink-jobmanager: # ... 其他 volumes: - ./flink-jobs:/opt/flink/usrlib:ro - flink-checkpoints:/opt/flink/checkpoints # 改为命名卷 - ./flink-savepoints:/opt/flink/savepoints flink-taskmanager: # ... 其他 volumes: - ./flink-jobs:/opt/flink/usrlib:ro - flink-checkpoints:/opt/flink/checkpoints # 同样挂载命名卷 - ./flink-savepoints:/opt/flink/savepoints然后执行docker volume rm flink-checkpoints docker-compose up -d5.5 现象redis容器启动后docker logs redis显示1:M 20 May 12:34:56.789 # Cant handle RDB format version 10随后崩溃退出原因宿主机/var/lib/docker/volumes/下残留了旧版 Redis 的 RDB 文件如从 Redis 6 升级到 7新版 Redis 7 无法读取旧格式。解决彻底清理 Redis 数据卷# 查看所有 redis 相关卷 docker volume ls | grep redis # 删除谨慎确保无重要数据 docker volume rm redis-volume-name # 或暴力清空开发环境安全 sudo rm -rf /var/lib/docker/volumes/*redis* docker-compose up -d redis6. 进阶技巧用docker stats实时监控容器资源3 行命令定位 Kafka 吞吐瓶颈当你把系统跑起来真正考验功力的不是“能不能跑”而是“怎么知道它跑得好不好”。我不会教你用top看宿主机 CPU而是给你一套精准到容器、到进程、到 Kafka 分区的诊断组合拳。这套方法我在三次毕设答辩现场救场帮同学从“老师说你这系统不实时”逆袭到“这个监控维度很专业”。6.1 第一步用docker stats锁定高负载容器在终端执行docker stats --format table {{.Name}}\t{{.CPUPerc}}\t{{.MemUsage}}\t{{.NetIO}} --no-stream你会看到类似输出NAME CPU % MEM USAGE / LIMIT NET I/O kafka 82.34% 1.2GiB / 3.8GiB 12.4MB / 8.2MB flink-taskmanager 45.12% 1.8GiB / 3.8GiB 5.6MB / 15.3MB filebeat 12.78% 245MiB / 3.8GiB 8.9MB / 2.1MB关键洞察如果kafkaCPU 80%而filebeatCPU 15%说明瓶颈在 Kafka 本身如磁盘 IO 或网络不是 Filebeat 采集慢如果flink-taskmanagerMEM USAGE 接近 LIMIT但kafkaCPU 正常说明 Flink 状态过大需调taskmanager.memory.process.sizeNET I/O的OUT发送远大于IN接收说明 Kafka 正在向外输送大量数据Flink 消费正常反之则可能是 Flink 消费慢数据在 Kafka 积压。6.2 第二步进入 Kafka 容器用kafka-run-class.sh查分区级 Lag当docker stats怀疑 Kafka 积压时深入容器docker exec -it kafka bash # 查看 consumer group 的 lag关键 kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group flink-group --describe输出中重点关注TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID monitoring-logs 0 125000 125000 0 flink-12345 monitoring-logs 1 124998 125000 2 p a hrefhttps://download.csdn.net/download/weixin_49376454/90173132 stylecolor:#ec7500;font-size:14px; 本文还有配套的精品资源点击获取 /a img altmenu-r.4af5f7ec.gif srchttps://csdnimg.cn/release/wenkucmsfe/public/img/menu-r.4af5f7ec.gif stylewidth:16px;margin-left:4px;vertical-align:text-bottom;cursor:text; /p
返回列表