ELK 与 Canal

2023-03-31T14:11:02+08:00 | 17分钟阅读 | 更新于 2026-03-31T14:11:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 讲清楚 ELK 是什么、解决什么问题,以及它在"可观测性"里占据的位置(日志维度)。
  2. 配出一条 Logstash 管道(input/filter/output),并说清 Grok、Date、Mutate、Kv 各自干嘛。
  3. 解释为什么生产上 Filebeat + Logstash 要分工,而不是让 Logstash 直接蹲在每台机器上。
  4. 讲透 Canal 的本质:它怎么靠"伪装成 MySQL 从库"拿到 Binlog,以及 binlog 三种模式的取舍。
  5. 用 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_ISO8601IPEMAIL),可以直接使用。

(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

维度KibanaGrafana
设计目标与 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 数据源

  1. 浏览器访问 http://localhost:5601
  2. 选择 Elasticsearch logs 作为数据源
  3. 默认配置会引导安装 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 等多个目的地
结合 LogstashFilebeat 采集 + 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 的作用

  1. 实时监控数据库变更:捕获 INSERT、UPDATE、DELETE 操作
  2. 数据同步:将一个数据库的变更同步到另一个数据库,保持一致性
  3. 支持实时分析:将数据及时传输到数据仓库 / 分析平台
  4. 解耦数据库系统:引入新数据库或更改结构时,不影响其他部分运作

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 的配置比较复杂,大体分为三部分:

  1. Canal Server 本体配置:Canal Server 自身运行所需的配置
  2. 数据库连接配置:每个要连接的数据库都需要一份配置,包含连接信息、用户信息(用户需具备较高权限)
  3. 转发配置: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 里统一看。日志分析找"发生了什么",指标监控看"趋势正不正常",链路追踪定位"慢在哪"。

面试话术思路

  1. 接手项目时,可观测性极差
  2. 引入 ELK + Filebeat + Prometheus + OTel,系统提高可观测性
  3. 三年以上经验:讲自己如何推动日志规范、跨部门可观测性改造
  4. 性能优化:提高可观测性后发现了哪些性能问题,如何解决
  5. 可用性优化:发现可用性瓶颈,如何改进

晋升关键:在公司层面引入 ELK 是最好刷的 KPI,能带来"快速发现问题、快速解决问题"的收益,提高系统可用性和稳定性。


十一、自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. ELK 四段式架构是什么?为什么生产上 Filebeat 和 Logstash 要"轻采集、重处理"分工,而不是让 Logstash 直接蹲每台机器?
  2. Canal 为什么能实时拿到数据库的变更?它和 MySQL 主从同步是什么关系?
  3. Binlog 的 Row 模式和 Statement 模式各有什么优缺点?为什么 Canal 场景几乎都用 Row 模式?
  4. Canal 更新缓存时,为什么要"按主键 hash 分区"?如果同主键消息分散到不同分区会发生什么?
  5. Canal 的 DELETE 消息里,被删除的数据在 data 还是 old?处理删除时该从哪取主键?

动手练习(建议真做一遍)

  1. 起一套 ELK + Filebeat:用 Docker Compose 起 Filebeat/Logstash/ES/Kibana,把你一个 Go 服务的 JSON 日志配置好,在 Kibana Discover 里看到自己的日志。
  2. 起 Canal 看 Binlog:起 Canal + MySQL(ROW 模式),改一行 users 数据,用 kafka-console-consumer 订阅 webook_binlog 看消息结构。
  3. 写个缓存同步消费者:监听 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 把性能压出拐点。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!