学习目标
学完本章,你应该能够:
- 讲清楚 ELK 是什么、解决什么问题,以及它在"可观测性"里占据的位置(日志维度)。
- 配出一条 Logstash 管道(input/filter/output),并说清 Grok、Date、Mutate、Kv 各自干嘛。
- 解释为什么生产上 Filebeat + Logstash 要分工,而不是让 Logstash 直接蹲在每台机器上。
- 讲透 Canal 的本质:它怎么靠"伪装成 MySQL 从库"拿到 Binlog,以及 binlog 三种模式的取舍。
- 用 Canal 落地"监听 Binlog 异步更新缓存 / 数据校验",并讲清顺序性为什么要按主键 hash 分区。
前置知识(如果下面任一点生疏,先回看对应章):
- 第05章 缓存:知道 Redis 基本用法与"双写一致性"难题(Canal 更新缓存要用到)。
- 第07章 Kafka:知道 Topic / 分区 / 消费组(Canal 转发到 Kafka、消息顺序性要靠它)。
- 第17章 ES:知道索引与写入(Logstash 最终把日志写进 ES)。
- 基本 MySQL 主从概念(理解 Canal 伪装从库的前提)。
本章你会动手做的事:
- 用 Docker Compose 起 ELK + Filebeat,把自己 Go 服务的 JSON 日志跑进 Kibana 看板。
- 起一个 Canal + MySQL,改一行数据,在 Kafka 里看到那条 Binlog 消息。
- 写一个小消费者,监听 Canal 的 Binlog 去更新本地缓存,验证"改库 → 缓存自动变"。
一、ELK 介绍与应用
1.1 什么是 ELK
ELK 是一个强大的开源日志管理和分析平台,由三个核心组件组成:
| 组件 | 作用 | 类比 |
|---|---|---|
| Elasticsearch | 分布式搜索引擎,实时存储、搜索、分析 | 数据存储 + 检索引擎 |
| Logstash | 日志数据的收集、处理、传输 | 数据管道 / ETL 工具 |
| Kibana | 数据可视化工具,实时分析和交互式搜索 | 数据展示面板 |
核心理解(讲义原话):“这一切都可以总结为四个字:文本分析"。对程序员来说,最重要的两个功能是日志分析和实时监控。
白话类比:ELK 就像工厂的"监控室三件套”——Filebeat/Logstash 是巡线员(到处捡日志纸条),ES 是档案室(把纸条归档还能秒查),Kibana 是大屏(把档案画成图给老板看)。没有它,线上出问题你就像在黑屋子里找开关。
1.2 ELK 的应用领域
- 日志分析:追踪应用程序和系统的日志,帮助诊断问题、优化性能
- 实时监控:通过对实时数据的分析,及时发现和解决问题
- 安全分析:监测潜在的安全威胁和异常行为
- 业务智能:利用数据可视化分析,帮助业务决策
flowchart LR
A[应用日志] --> F[Filebeat 采集]
F --> L[Logstash 处理]
L --> E[(ES 存储)]
E --> K[Kibana 展示]这张图在讲:一条日志从产生到被看见,要走过"采集 → 处理 → 存储 → 展示"四段。后面所有组件都是这条链路上的节点。
二、Logstash
2.1 Logstash 简介
Logstash 是 ELK 中的数据处理引擎,负责日志数据的收集、过滤、转换和传输。
核心功能三段式:
Input(输入) → Filter(过滤) → Output(输出)
↓ ↓ ↓
从各种来源 解析、结构化 发送到目的地
接收数据 过滤数据 (如 ES)
flowchart LR
I[Input
从文件/Kafka/Beats 收] --> F[Filter
解析/结构化/清洗]
F --> O[Output
写 ES/Redis/文件]这张图在讲:Logstash 就是一根"数据管道",左边进原始日志,中间结构化,右边出到目的地。
2.2 Logstash 配置文件结构
配置文件就是三段式的声明:
# logstash.conf 示例
input {
file {
path => "/var/log/webook/*.log" # 日志文件位置
start_position => "beginning" # 从文件开头读取
}
}
filter {
# 使用 grok 插件解析日志消息,提取时间戳、日志级别和消息内容
grok {
match => { "message" => "%{TIMESTAMP_ISO8601:timestamp} %{LOGLEVEL:level} %{GREEDYDATA:msg}" }
}
}
output {
elasticsearch {
hosts => ["http://elasticsearch:9200"]
index => "webook-logs-%{+YYYY.MM.dd}" # 按天分索引
}
}
实践建议:配置不需要死记硬背,使用时查阅文档或问 GPT 即可。
2.3 Logstash 常用 Filter 插件
(1)Grok 插件
通过正则表达式解析非结构化日志,提取字段。
filter {
grok {
match => { "message" => "%{TIMESTAMP_ISO8601:timestamp} %{LOGLEVEL:level} %{GREEDYDATA:msg}" }
}
}
Grok 内置了大量模式(如 TIMESTAMP_ISO8601、IP、EMAIL),可以直接使用。
(2)Date 插件
将字符串转换为日期格式,通常与 Grok 配合使用,标准化时间戳字段。
filter {
date {
match => ["timestamp", "ISO8601"] # 解析 timestamp 字段,按 ISO8601 格式
target => "@timestamp" # 写入 @timestamp 字段(ES 默认时间字段)
}
}
(3)Mutate 插件
提供数据变换操作:重命名、拼接、删除字段等。
filter {
mutate {
add_field => { "new_field" => "Hello, World!" } # 添加字段
remove_field => ["unwanted_field"] # 删除字段
rename => { "old_field" => "new_field" } # 重命名字段
}
}
建议:不要在 Mutate 中写过于复杂的逻辑,复杂的数据处理应该通过良好的日志规范在源头解决,而不是在 Logstash 里强行补救。
(4)Kv 插件
从未结构化文本中提取键值对。
filter {
kv {
source => "message" # 从 message 字段提取
field_split => "," # 键值对之间用逗号分隔
}
}
在结构化日志(如 Go 的 logger.Field)场景下不太用得上,主要用于半结构化文本。
三、Kibana
3.1 Kibana 简介
Kibana 是 ELK 中的数据可视化工具,通过直观的用户界面帮助用户查询、分析和可视化 ES 中的数据。
核心功能:
- 仪表板(Dashboard):创建交互式仪表板,集成多个图表和可视化组件
- 搜索和过滤:在数据集中执行高级搜索和过滤,定位感兴趣的数据
- 图表和可视化:使用柱状图、折线图、地图等多种图表类型呈现数据
定位:讲义原话 “Kibana 是一个侧重于数据展示的框架,主打四个字 —— 花里胡哨"。
3.2 Kibana 应用场景
- 日志分析:搜索、过滤、可视化 ES 中的日志数据,实时监控系统运行
- 性能监控:通过仪表板展示系统性能指标,发现和解决性能问题
- 安全分析:可视化分析安全事件,提高对潜在威胁的识别和响应能力
3.3 Kibana vs Grafana
| 维度 | Kibana | Grafana |
|---|---|---|
| 设计目标 | 与 ES 深度集成,专注日志和指标可视化 | 通用仪表板,支持多种数据源 |
| 数据源 | 主要支持 ES | 支持 Prometheus、MySQL、InfluxDB 等多种 |
| 告警 | 较新且相对简单 | 强大成熟,支持邮件、Slack、Webhook 等多种通知渠道 |
| 适用场景 | 日志分析 | 通用监控、指标告警 |
实践建议:日志分析用 Kibana,指标监控和告警用 Grafana,两者常配合使用。
flowchart TD
E[(ES 日志)] --> K[Kibana
日志分析]
P[(Prometheus 指标)] --> G[Grafana
监控告警]这张图在讲:日志走 Kibana、指标走 Grafana,分工明确、各取所长。
四、部署 ELK
4.1 Docker Compose 部署
ES 已在前一章部署,这里只需额外部署 Logstash 和 Kibana:
# docker-compose.yml
services:
logstash:
image: docker.elastic.co/logstash/logstash:8.0.0
volumes:
- ./logstash.conf:/usr/share/logstash/pipeline/logstash.conf
ports:
- "5044:5044" # Filebeat 发送数据的端口
depends_on:
- elasticsearch
kibana:
image: docker.elastic.co/kibana/kibana:8.0.0
environment:
- ELASTICSEARCH_HOSTS=http://elasticsearch:9200 # 关键:配置 ES 地址
ports:
- "5601:5601"
depends_on:
- elasticsearch
4.2 在 Kibana 中配置 ES 数据源
- 浏览器访问
http://localhost:5601 - 选择 Elasticsearch logs 作为数据源
- 默认配置会引导安装 Filebeat,但默认是 Filebeat 直接送数据到 ES,与我们想经过 Logstash 的预期不符,需要修改配置
4.3 Filebeat 配置
# filebeat.yml
filebeat.inputs:
- type: log
enabled: true
paths:
- /var/log/webook/*.log # 日志文件路径,可用通配符
# 输出到 Logstash(不是直接到 ES)
output.logstash:
hosts: ["logstash:5044"]
# 也可以直接发送到 ES(不经过 Logstash):
# output.elasticsearch:
# hosts: ["elasticsearch:9200"]
注意:Windows 路径需要做对应修改。
4.4 Go 应用日志初始化
为了让日志写入对应目录,初始化日志时要指定好目录,并引入 lumberjack 库管理日志切片:
package logger
import (
"gopkg.in/natefinch/lumberjack.v2"
"go.uber.org/zap"
"go.uber.org/zap/zapcore"
)
// InitLogger 初始化日志,使用 lumberjack 管理切片
// 当日志文件很大或时间变化时,自动分割成多个文件
func InitLogger(filepath string) *zap.Logger {
// 步骤 1:配置 lumberjack 切片器(大小/备份数/保留天数)
lumberJackLogger := &lumberjack.Logger{
Filename: filepath, // 日志文件路径
MaxSize: 100, // 单文件最大 MB
MaxBackups: 5, // 保留旧文件数
MaxAge: 30, // 保留天数
Compress: true, // 是否压缩
}
// 步骤 2:用 JSON 编码器 + 写入 lumberjack 构造 zap core
core := zapcore.NewCore(
zapcore.NewJSONEncoder(zap.NewProductionEncoderConfig()),
zapcore.AddSync(lumberJackLogger),
zapcore.InfoLevel,
)
// 步骤 3:返回 logger
return zap.New(core)
}
lumberjack 的作用:日志文件超过阈值或时间变化时切片,便于 Filebeat 收集、归档管理。
4.5 Logstash 配置(处理 JSON 日志)
input {
beats {
port => 5044 # 接收 Filebeat 发送的数据
}
}
filter {
# 将 message 转化为 JSON 后作为一个 data 字段
# Filebeat 传递的数据,日志内容在 message 里
json {
source => "message"
target => "data"
}
}
output {
elasticsearch {
hosts => ["http://elasticsearch:9200"]
index => "webook-logs-%{+YYYY.MM.dd}"
}
}
4.6 Kibana 展示
不做处理时,Kibana Discover 展示的数据非常难读,可通过选择关心的列作为展示字段来优化。
五、为什么使用 Filebeat
实践中结合 Filebeat 是很常见的做法,主要原因:
| 优势 | 说明 |
|---|---|
| 轻量级高效 | 比 Logstash 轻量得多,适合部署在每台机器上采集日志 |
| 实时性 | 实时监测日志文件变化,迅速传输新日志 |
| 模块化配置 | 支持系统日志、NGINX、Apache 等多种日志格式 |
| 多输出支持 | 可发送到 Logstash、ES、Kafka 等多个目的地 |
| 结合 Logstash | Filebeat 采集 + Logstash 处理,分工协作 |
| 自动发现 | 支持自动发现新日志文件,标准化日志格式 |
| 容器环境友好 | 轻松集成到容器环境,支持容器化平台 |
架构理解:Filebeat 部署在每台业务机器上轻量采集,Logstash 集中处理过滤,ES 存储检索,Kibana 展示。这是经典的"采集 - 处理 - 存储 - 展示"四段式架构。
⚠️ 新手必踩的坑:别让 Logstash 蹲在每台机器上。Logstash 是 JVM 应用,吃内存吃 CPU,几十台机器各跑一个 Logstash 能把机器拖垮。正确做法是每台机器只跑轻量 Filebeat 采集,集中的 Logstash 做处理。这就是"采集轻、处理重"的分工。
5.1 完整数据流
应用日志 → lumberjack 切片 → Filebeat 采集 → Logstash 过滤 → ES 存储 → Kibana 展示
↓
也可直接发送到 ES(跳过 Logstash)
flowchart LR
APP[应用 + lumberjack] --> FB[Filebeat 轻量采集]
FB --> LS[Logstash 集中处理]
LS --> ES[(ES 存储)]
ES --> KB[Kibana 展示]
FB -. 也可直连 .-> ES这张图在讲:标准链路是"应用→Filebeat→Logstash→ES→Kibana”;Filebeat 也能跳过 Logstash 直写 ES,但那就失去集中处理能力了。
六、Canal 简介
6.1 Canal 是什么
Canal 是一款开源的数据库实时变更监控和数据同步工具,支持 MySQL、MariaDB、阿里云 RDS 等。它是一个典型的 CDC(Change Data Capture)工具,能够实时捕获数据库变更,提供高性能数据同步服务。
应用场景:数据仓库同步、实时分析、缓存更新、数据校验等。
优势:
- 实时性高:实时监控数据库变更
- 灵活性强:配置和使用相对简单
- 开源社区支持:活跃社区,及时更新
白话类比:Canal 就像一个坐在数据库旁边的"抄写员"。数据库每改一笔账(INSERT/UPDATE/DELETE),抄写员立刻把这笔变动抄成一张小纸条,送到你指定的地方(Kafka)。你不用改业务代码,就能"感知"到数据库发生的所有变化。
6.2 Canal 的作用
- 实时监控数据库变更:捕获 INSERT、UPDATE、DELETE 操作
- 数据同步:将一个数据库的变更同步到另一个数据库,保持一致性
- 支持实时分析:将数据及时传输到数据仓库 / 分析平台
- 解耦数据库系统:引入新数据库或更改结构时,不影响其他部分运作
6.3 Canal 的基本组成
| 组件 | 作用 |
|---|---|
| Canal Server | 核心组件,连接数据库并实时监控变更,捕获变更日志发送给客户端 |
| Canal Client | 与 Server 通信,接收并处理变更信息 |
| Binlog | 数据库二进制日志,Canal 实时监控的基础 |
| 数据格式转换器 | 将变更日志转换为 JSON、Avro 等格式 |
| Canal 配置文件 | 包含数据库连接、监控规则、数据格式等配置 |
| ZooKeeper(可选) | 分布式场景下用于服务协调和管理,提供高可用和容错 |
七、Binlog 基础
7.1 什么是 Binlog
Binlog 是 MySQL 中的二进制日志,记录数据库中的每个变更操作(INSERT、UPDATE、DELETE 的详细信息)。它是 Canal 实时捕获变更的重要基础。
Binlog 在主从同步中的角色:
1. 从库连上主库
2. 从库发起数据同步请求
3. 主库开启一个线程,将 Binlog 发送到从节点
4. 从节点收到 Binlog,先写到 Relay log,再逐步执行 Relay log 中的数据变更
关键理解:Canal 的原理就是伪装成 MySQL 从库,让主库把 Binlog 发送过来,从而实时获得数据变更。Canal 解析 Binlog 后转换为业务可读的消息格式。
flowchart TD
M[(MySQL 主库)] -->|推送 Binlog| C[Canal 伪装成从库]
C -->|解析变更| K[Kafka Topic]
K --> CON[消费者: 更新缓存/校验]这张图在讲:Canal 把自己打扮成"从库",主库就把 Binlog 推给它,它解析后丢进 Kafka,业务消费者据此更新缓存或做校验。业务代码全程无感知。
7.2 Binlog 的三种模式
| 模式 | 说明 | 优点 | 缺点 |
|---|---|---|---|
| Row-based(基于行) | 记录每行数据的变更 | 变更信息最详细 | 日志量大 |
| Statement-based(基于语句) | 记录 SQL 语句 | 日志量小 | 无法捕获复杂变更(如 NOW()) |
| Mixed(混合) | 自动选择行级或语句级 | 平衡详细信息和日志量 | 复杂度提高 |
讲义建议:“基于行的日志记录用起来比较方便”。Canal 支持解析这三种模式,根据实际情况选择。
⚠️ 新手必踩的坑:Statement 模式踩
NOW()/UUID()。基于语句的 Binlog 只记录"执行了UPDATE ... SET t=NOW()",但从库重放时NOW()取值和主库不同,导致主从数据不一致。所以 Canal 场景几乎都用 Row 模式——它记录的是"改完后的真实值",重放结果一定一致。
7.3 Canal 的配置分类
Canal 的配置比较复杂,大体分为三部分:
- Canal Server 本体配置:Canal Server 自身运行所需的配置
- 数据库连接配置:每个要连接的数据库都需要一份配置,包含连接信息、用户信息(用户需具备较高权限)
- 转发配置:Canal 收到 Binlog 后要转发到哪里(如 Kafka topic)
Canal Server 本体配置(关键片段)
# canal.properties
canal.serverMode = kafka # 使用 Kafka 作为转发模式
canal.mq.servers = kafka:9092 # Kafka 地址
数据库连接配置(instance 配置)
# example/instance.properties
canal.instance.master.address = mysql:3306
canal.instance.dbUsername = canal # 需要高权限用户
canal.instance.dbPassword = canal123
canal.mq.topic = webook_binlog # 发送到该 Kafka topic
Kafka 转发配置
canal.mq.servers = kafka:9092
canal.mq.partition = 0 # 默认分区
# 也可以按 hash 分区,保证同一主键的消息顺序
canal.mq.partitionHash = .*\\..*:$pk$ # 按 主键 hash 分区
7.4 docker compose 配置
services:
mysql:
image: mysql:8.0
command:
- --binlog-format=ROW # 关键:使用 Row 格式的 Binlog
- --binlog-row-image=FULL
environment:
MYSQL_ROOT_PASSWORD: root
volumes:
- ./init.sql:/docker-entrypoint-initdb.d/init.sql # 创建 canal 用户
canal:
image: canal/canal-server:latest
environment:
- canal.destinations=example
volumes:
- ./canal.properties:/home/admin/canal-server/conf/canal.properties
- ./example/instance.properties:/home/admin/canal-server/conf/example/instance.properties
depends_on:
- mysql
- kafka
实践建议(讲义原话):“不要深究配置问题,直接借助 docker compose 文件启动即可”。
7.5 启动验证
启动后用 Kafka 消费者工具(如 kafka-console-consumer)订阅 webook_binlog topic,更新任意表的任意数据,看到控制台输出 Binlog 消息即部署成功。
八、Canal 消息格式
8.1 插入语句消息
{
"data": [
{
"id": "1",
"name": "大明",
"email": "john@example.com"
}
],
"database": "webook",
"table": "users",
"type": "INSERT",
"ts": 1234567890
}
8.2 更新语句消息
更新语句比插入多了 old 字段,表示更新前的数据:
{
"data": [
{
"id": "1",
"name": "大明2",
"email": "john@example.com"
}
],
"old": [
{
"name": "大明"
}
],
"type": "UPDATE"
}
8.3 删除语句消息
data 字段:被删除的数据
old 字段:反而没有
type:DELETE
讲义吐槽(原话):“这种设计曾经导致我写过一堆垃圾代码。” 因为直觉上
old才应该存放被删除的数据,但实际是在data里。
⚠️ 新手必踩的坑:DELETE 的数据在
data不在old。直觉上"被删的东西应该在 old 里",但 Canal 把删除行的完整内容放在data。处理删除时务必从data取主键去删缓存,而不是去翻old(DELETE 消息里old是空的)。
九、Canal 使用案例
9.1 案例一:借助 Canal 更新缓存
正常使用缓存时,更新缓存与更新 DB 的并发问题是个经典难题(双写一致性)。Canal 提供了一种解耦方案:监听 Binlog 异步更新缓存。
两种做法:
| 做法 | 描述 | 优缺点 |
|---|---|---|
| 直接用 Canal 数据更新 | 用 Binlog 中的数据直接回写缓存 | 性能好;需保证同一主键消息在同一分区(顺序性) |
| Canal 仅作信号器 | 收到信号后从 DB 加载再回写 | 性能差,对 DB 压力大 |
本课程采用方案一,需要保证消息顺序:
# canal 的 topic 和分区配置
# 按主键 hash 分区,保证同一 ID 的消息发到同一分区
canal.mq.partitionHash = .*\\..*:$pk$
顺序性问题
分区 0: 消息1(id=1, INSERT) → 消息2(id=1, UPDATE a=2) → 消息3(id=1, UPDATE a=3)
分区 1: 消息4(id=2, INSERT)
只有同一主键的消息在同一分区,才能保证消费顺序与 Binlog 顺序一致,避免"用旧数据覆盖新数据"的问题。
flowchart LR
subgraph P0[分区 0: id=1]
M1[INSERT a=1] --> M2[UPDATE a=2] --> M3[UPDATE a=3]
end
subgraph P1[分区 1: id=2]
M4[INSERT id=2]
end
P0 --> C[消费者按序回写缓存]这张图在讲:同一主键的消息被 hash 到同一分区,Kafka 分区内有序,消费者才能拿到"INSERT→UPDATE→UPDATE"的正确顺序,不会用旧值覆盖新值。
表与 Topic 的关系
小规模:所有表共用一个 topic
大规模:不同表使用不同 topic,甚至使用不同的 Kafka 集群
避免消息积压、Kafka 集群性能瓶颈
代码实现
// CanalBinlogConsumer 监听 Canal 发送的 Binlog 消息,更新缓存
// 关键点:利用 Binlog 更新缓存是"缓存策略",不是业务逻辑
// 一般不通过 Service 更新,而是直接绕开 Service 操作 Repository 或 Cache
type CanalBinlogConsumer struct {
cache cache.UserCache
repo repository.UserRepository
}
func (c *CanalBinlogConsumer) Consume(ctx context.Context, msg kafka.Message) error {
// 步骤 1:反序列化 Binlog 消息
var event BinlogEvent
if err := json.Unmarshal(msg.Value, &event); err != nil {
return err
}
// 步骤 2:只处理 users 表
if event.Table != "users" {
return nil
}
// 步骤 3:按类型更新 / 删除缓存
switch event.Type {
case "INSERT", "UPDATE":
for _, row := range event.Data {
uid, _ := strconv.ParseInt(row["id"], 10, 64)
// 直接用 Binlog 中的数据更新缓存
user := domain.User{
ID: uid,
Nickname: row["nickname"],
Email: row["email"],
}
_ = c.cache.Set(ctx, user)
}
case "DELETE":
for _, row := range event.Data {
uid, _ := strconv.ParseInt(row["id"], 10, 64)
_ = c.cache.Del(ctx, uid)
}
}
return nil
}
设计要点:
- Canal 更新缓存是"具体缓存策略",不适合在 Repository 上定义接口
- 借助 Kafka 可以设计重试机制,解决部分失败问题
- 追求的是最终一致性,不是强一致性
9.2 案例二:借助 Canal 完成数据校验
在数据迁移场景中,使用 Canal 完成增量数据校验与修复。
数据校验流程
// Consumer 校验逻辑
// 业务方只需创建消费者并调用此方法
func (c *CanalVerifyConsumer) Consume(ctx context.Context, msg kafka.Message) error {
// 步骤 1:反序列化
var event BinlogEvent
if err := json.Unmarshal(msg.Value, &event); err != nil {
return err
}
// 步骤 2:不管增删改,只拿主键去校验两边
for _, row := range event.Data {
id, _ := strconv.ParseInt(row["id"], 10, 64)
// 步骤 3:按"以谁为准"策略修复
if c.useSourceAsTruth {
// SRC 为准:从源表读,写入/修复目标表
src, err := c.srcRepo.FindByID(ctx, id)
if err != nil {
return err
}
_ = c.dstRepo.Upsert(ctx, src)
} else {
// DST 为准:从目标表读,修复源表
dst, err := c.dstRepo.FindByID(ctx, id)
if err != nil {
return err
}
_ = c.srcRepo.Upsert(ctx, dst)
}
}
return nil
}
切换"以谁为准"的陷阱
从"源表为准"切换到"目标表为准"时,可能引起数据不一致:
1. 双写阶段,SRC 为准:业务更新 SRC.a = 2,产生 Binlog
2. Binlog 到达 Canal 时,切换为 DST 为准
3. 消费者收到 Binlog(a=2),发现是源库的 binlog,按"目标表为准"策略忽略
4. 实际上 SRC 已经是 a=2,但 DST 还是 a=1,最终未发现不一致
解决方案:切换期间被修改过的数据需要手动校验一遍。
十、可观测性体系全景(面试加分)
完整方案:ELK + Prometheus + OpenTelemetry + Grafana
- ELK:日志分析与文本搜索
- Prometheus:指标监控
- OpenTelemetry:链路追踪
- Grafana:通用仪表板与告警
flowchart LR
APP[应用] -->|日志| ELK[ELK 日志]
APP -->|指标| PROM[Prometheus]
APP -->|链路| OTEL[OpenTelemetry]
ELK --> G[Grafana 统一看板]
PROM --> G
OTEL --> G这张图在讲:可观测性三件套——日志(ELK)、指标(Prometheus)、链路(OTel),最后都可以在 Grafana 里统一看。日志分析找"发生了什么",指标监控看"趋势正不正常",链路追踪定位"慢在哪"。
面试话术思路:
- 接手项目时,可观测性极差
- 引入 ELK + Filebeat + Prometheus + OTel,系统提高可观测性
- 三年以上经验:讲自己如何推动日志规范、跨部门可观测性改造
- 性能优化:提高可观测性后发现了哪些性能问题,如何解决
- 可用性优化:发现可用性瓶颈,如何改进
晋升关键:在公司层面引入 ELK 是最好刷的 KPI,能带来"快速发现问题、快速解决问题"的收益,提高系统可用性和稳定性。
十一、自测题与动手练习
自测题(合上书能答出来,才算懂):
- ELK 四段式架构是什么?为什么生产上 Filebeat 和 Logstash 要"轻采集、重处理"分工,而不是让 Logstash 直接蹲每台机器?
- Canal 为什么能实时拿到数据库的变更?它和 MySQL 主从同步是什么关系?
- Binlog 的 Row 模式和 Statement 模式各有什么优缺点?为什么 Canal 场景几乎都用 Row 模式?
- Canal 更新缓存时,为什么要"按主键 hash 分区"?如果同主键消息分散到不同分区会发生什么?
- Canal 的 DELETE 消息里,被删除的数据在
data还是old?处理删除时该从哪取主键?
动手练习(建议真做一遍):
- 起一套 ELK + Filebeat:用 Docker Compose 起 Filebeat/Logstash/ES/Kibana,把你一个 Go 服务的 JSON 日志配置好,在 Kibana Discover 里看到自己的日志。
- 起 Canal 看 Binlog:起 Canal + MySQL(ROW 模式),改一行
users数据,用kafka-console-consumer订阅webook_binlog看消息结构。 - 写个缓存同步消费者:监听 Canal 的 Binlog,收到
users表的变更就更新本地 Map/Redis 缓存,验证"改库 → 缓存自动同步",并故意把分区改成单分区观察顺序是否仍正确。
十二、本章小结
- ELK = 采集→处理→存储→展示:Filebeat 轻量蹲机器采集,Logstash 集中做 Input→Filter→Output 管道,ES 存与搜,Kibana 展示;日志分析用 Kibana,指标监控用 Grafana。
- Logstash 三件套插件:Grok 解析、Date 标准化时间、Mutate 变换字段、Kv 提键值对;复杂逻辑应在日志源头规范掉,别堆在 Logstash。
- Canal 本质是 CDC,伪装成 MySQL 从库拿 Binlog;Binlog 三种模式里 Row 最稳(记录真实值,重放一致),Statement 有
NOW()类坑。 - Canal 更新缓存走"最终一致性":监听 Binlog 直接回写缓存,解耦业务;顺序性靠"按主键 hash 分区"保证同一主键消息同分区有序,避免旧值覆盖新值。
- 可观测性全景:ELK(日志)+ Prometheus(指标)+ OTel(链路)+ Grafana(看板),在公司层面推 ELK 是性价比极高的 KPI。
下一章(第19章)我们进入 Feed 流设计与压测——把前面学的 Kafka、异步、聚合、缓存综合起来,设计一个能扛住百万粉丝的 Feed 系统,并用 k6 把性能压出拐点。