一、问题:Canal消费追不平,ES节点CPU 85%
2024年3月,我负责的UGC平台在做大促前压测。业务链路是 MySQL → Binlog → Canal → ES,当天下午压测时,Canal 的消费延迟从 3 秒涨到了 8 分钟,而且还在涨。查监控,ES 集群 3 个数据节点 CPU 全部 85% 以上,Young GC 每分钟 12 次,单节点堆内存频繁顶满。
最恶心的是,一旦写入压力上来,Canal 客户端直接抛 EsRejectedExecutionException,写入被拒绝。看 ES 慢日志,bulk 队列每秒只能消化 4000~5000 条文档,而业务高峰期需要每秒写入 2 万+。我当时的操作和大多数人一样——加节点。
环境:ES 8.11.3,3 个 master 节点 + 3 个热数据节点(8C32G),SSD 云盘,单个日志索引已经跑到 800+ 段。数据总量约 4TB,每天新增 300GB。
二、两个无效方案:加节点、盲调刷新间隔
方案A:加节点 + 加副本
运维把热数据节点从 3 个加到 6 个,副本数从 0 改成 1。加完节点,集群状态从 green 变 yellow,因为副本在恢复,大量数据从一个节点复制到另一个节点,节点负载雪上加霜。
| 指标 | 加节点前 | 加节点后 |
|---|---|---|
| CPU(峰值) | 85% | 81% |
| P99 写入延迟 | 1.2s | 1.1s |
| bulk 被拒绝次数/分钟 | 270+ | 180+ |
| 基础设施成本 | 3×8C32G | 6×8C32G,月成本 +100% |
加节点唯一的作用是 CPU 从 85% 降到 81%,钱白花了。问题根本不在节点数量,而在单节点上的索引写入链路效率太低。
方案B:把 refresh_interval 调到 30s,bulk size 调到 50MB
这是第二板斧。把刷新间隔从默认 1s 改成 30s,bulk size 从 5MB 调到 50MB,写入吞吐确实上去了,P99 从 1.2s 降到 750ms。副作用立刻来了:查询链路要求秒级可见数据,我把 refresh 调到 30s,意味着写成功后最多要等 30s 才能查到。
上线当天就接到前端反馈:"我发了个评论,刷新页面上没显示。" 我赶紧看,评论写入成功了,但索引没刷新。于是我把 refresh 调回 15s,查询倒是能看到了,但写入性能只有 20% 提升,远不够。而且段数量在持续增加,查询变慢的趋势已经出现。
到这里我意识到:治标不治本。 得从索引设计、数据生命周期、写入链路、集群参数整体调一遍。
三、我的方案:四层优化
第1层:Mapping 重构,字段类型严格收敛
旧索引的 mapping 有 200+ 字段,一半以上的字符串字段被 mapping 动态映射成了 text + keyword。这意味着每个字段存两份倒排索引,多占 30%~50% 磁盘空间,写入时还要多走一遍分词器,CPU 全耗在这了。
我们 90% 的字段根本不需要全文搜索,只需要精确过滤。重构后:字符串统一用 keyword,只有 message 字段用 text。字段数从 200+ 砍到 45 个。
#!/usr/bin/env bash
# 创建索引模板,新索引自动套用压缩后的mapping
# ES 8.11.3,需要设置 ES_ENDPOINT / ES_USER / ES_PASSWORD 环境变量
curl -s -u "${ES_USER}:${ES_PASSWORD}" \
-X PUT "http://${ES_ENDPOINT}/_index_template/log_ingest" \
-H "Content-Type: application/json" \
-d '{
"index_patterns": ["log-*"],
"template": {
"settings": {
"number_of_shards": 6,
"number_of_replicas": 1,
"refresh_interval": "10s",
"translog": {
"durability": "async",
"sync_interval": "5s"
},
"merge": {
"scheduler": {
"max_thread_count": 2
}
}
},
"mappings": {
"properties": {
"trace_id": { "type": "keyword" },
"user_id": { "type": "keyword" },
"event_time": { "type": "date" },
"level": { "type": "keyword" },
"message": { "type": "text", "analyzer": "standard" },
"cost_ms": { "type": "integer" },
"tags": { "type": "keyword" }
}
}
}
}'
echo "index template created"
这里有三个关键点:
- refresh_interval=10s:不是默认的 1s,也不是日志场景常用的 30s。10s 是我们测试后写入性能和可见性延迟的平衡点,后续压测数据会说明。
- translog 异步化:默认 translog 每次写入都 fsync,最安全但最慢。改成
async+ 5s 同步一次,写吞吐能提 30%-50%,代价是宕机最多丢 5s 数据。我们是日志平台,这个代价能接受;交易类场景不要这么干,下面避坑部分展开说。 - merge 线程限制:merge 线程数默认会自动分配,但遇到 bulk 批量写入时 merge 会和主写入抢 CPU。限制到 2,让 CPU 优先给写入。
第2层:ILM 生命周期:按天索引 + 自动合并 + 降冷
ES 的段数量直接影响写入和查询。之前索引跑了一个月,800+ 段,每次查询要打开几百个文件,本来 10ms 的查询变成 300ms。
解决思路:用 ILM 把索引按天切分,3 天后进 warm 节点并 forcemerge 成单段,15 天后进 cold 节点并关闭副本,不再占用热节点资源。
// 提交 ILM 策略:hot -> warm -> cold
// ES 8.11.3,先用 curl 或 Kibana DevTools 执行
PUT _ilm/policy/log_lifecycle
{
"policy": {
"phases": {
"hot": {
"actions": {
"rollover": {
"max_size": "30GB",
"max_age": "1d"
}
}
},
"warm": {
"min_age": "3d",
"actions": {
"forcemerge": {
"max_num_segments": 1
},
"allocate": {
"include": {
"node_data": "warm"
}
}
}
},
"cold": {
"min_age": "15d",
"actions": {
"allocate": {
"number_of_replicas": 0,
"include": {
"node_data": "cold"
}
}
}
}
}
}
}
然后把索引模板和这个策略绑定,再按照不同节点角色打标签热节点用 node_data: hot、温节点用 node_data: warm、冷节点用 node_data: cold。
这个策略跑通后,warm 阶段的索引被强制 merge 成 1 个段。磁盘清理 30% 空间,查询速度反而暴增,因为文件句柄占用少了两个数量级。
第3层:PHP 写入端改造:缓冲聚合 + 退避重试
Canal 消费端是 PHP 写的。原来的代码每条 log 单独调一次 index(),这个写法就是灾难:每写一条产生一个 HTTP 请求,ES 要处理一次连接握手和反序列化,性能上限极低。
改造后:数据先攒在内存里凑满 5000 条,一次性通过 bulk() 接口推给 ES,配合信号量控制并发量,避免瞬时打爆节点。
待写入缓存 */
private array $buffer = [];
private int $bufferSize = 5000;
public function __construct(string $endpoint, string $username, string $password)
{
$this->client = ClientBuilder::create()
->setHosts([$endpoint])
->setBasicAuthentication($username, $password)
->setRetries(3)
->build();
}
public function append(array $doc): void
{
$index = 'log-' . gmdate('Ymd');
$this->buffer[] = ['index' => ['_index' => $index]];
$this->buffer[] = $doc;
if (count($this->buffer) >= $this->bufferSize * 2) {
$this->flush();
}
}
public function flush(): void
{
if ($this->buffer === []) {
return;
}
$body = $this->buffer;
$this->buffer = [];
$retry = 5;
for ($i = 0; $i < $retry; $i++) {
try {
$resp = $this->client->bulk(['body' => $body]);
if ($resp->getStatusCode() !== 200) {
throw new \RuntimeException('bulk failed, status=' . $resp->getStatusCode());
}
return;
} catch (ClientResponseException | ServerResponseException $e) {
// 429 或 5xx 时退避重试
$waitMs = min(2000, 100 * (2 ** $i));
usleep($waitMs * 1000);
}
}
throw new \RuntimeException('es bulk write failed after 5 retries');
}
}
// 使用示例
// $writer = new EsBulkWriter('http://es-node:9200', 'elastic', 'xxxx');
// $writer->append(['trace_id' => 'abc', 'user_id' => 'u123', 'level' => 'INFO', 'message' => 'hello', 'cost_ms' => 12, 'event_time' => gmdate('c')]);
// $writer->flush(); // 进程结束前必须调用,防丢数据
有几个细节值得说:
- bufferSize 设为 5000,不是 10000 也不是 500。5000 是压测出来的最优值:每批 5000 条 bulk 请求体约 6MB~8MB,单节点处理时间约 300ms~500ms,在 CPU 和安全水位之间平衡最好。
- 指数退避重试:bulk 失败不立刻重试,而是 100ms → 200ms → 400ms → 800ms 逐步加长,防止大量重试请求在集群恢复瞬间再次打爆线程池。
- 进程结束必须调 flush():PHP-FPM 下如果用守护进程跑消费者,建议在 pcntl_signal 捕获退出信号后执行 flush。
第4层:集群参数调优
集群层面有几个参数,不改会埋雷,改了以后写入和分片迁移稳定很多。
# elasticsearch.yml —— ES 8.11.3
# 热数据节点:只做数据,不做 master,防止 GC 影响集群选举
cluster.name: prod-log
node.name: es-${HOSTNAME}
node.roles: [data_hot]
# master 节点单独配置 node.roles: [master]
# JVM 堆内存 16G(物理内存 32G,堆外留一半给 page cache)
ES_JAVA_OPTS="-Xms16g -Xmx16g"
# 磁盘水位线默认 85% 写入就停止,太保守
# 数据盘 70% 时就禁止分配分片,85% 时禁止写入
# 我们数据盘 1TB,改成 88% / 92% / 95%
cluster.routing.allocation.disk.threshold_enabled: true
cluster.routing.allocation.disk.watermark.low: 88%
cluster.routing.allocation.disk.watermark.high: 92%
cluster.routing.allocation.disk.watermark.flood_stage: 95%
# bulk 线程池默认队列 1024,压测时会积压丢请求
thread_pool.bulk.queue_size: 4096
thread_pool.write.queue_size: 4096
# 允许在滚动重启时分片并行恢复,缩短集群恢复时间
cluster.routing.allocation.node_concurrent_incoming_recoveries: 4
cluster.routing.allocation.node_concurrent_outgoing_recoveries: 4
JVM 堆大小这里再解释一句。以前看到很多配置把 64G 机器堆内存设为 31g,但 ES 8.x 是 Lucene 9,底层严重依赖操作系统 page cache。堆设得越大,留给 page cache 的内存越少,索引文件的随机读性能反而下降。我们的 32G 物理机设 16G 堆,实测比原来 24G 堆的写入性能高 15%,原因就是文件系统缓存更充足。
四、效果数据
四层优化全部上线后,用同一套压测脚本、同样的数据量,跑了一轮对比。压测脚本基于开源 esrally 思路简化,保证公平性。
压测脚本
#!/usr/bin/env bash
# 写入压测脚本:每次发 5000 条 doc,记录耗时
# 使用方式: ES_ENDPOINT=localhost:9200 ES_USER=elastic ES_PASSWORD=xxx ./bench.sh
set -euo pipefail
ES_ENDPOINT="${ES_ENDPOINT:-http://localhost:9200}"
INDEX="log-$(date +%Y%m%d)"
NDJSON_FILE=$(mktemp)
ROWS=5000
for _ in $(seq 1 "$ROWS"); do
trace_id=$(openssl rand -hex 8)
ts=$(date -u +"%Y-%m-%dT%H:%M:%S.000Z")
printf '{"index":{"_index":"%s"}}\n' "$INDEX" >> "$NDJSON_FILE"
printf '{"trace_id":"%s","user_id":"u%d","level":"INFO","message":"bulk bench","cost_ms":%d,"event_time":"%s","tags":["bench"]}\n' \
"$trace_id" "$((RANDOM % 9000 + 1000))" "$((RANDOM % 50 + 1))" "$ts" >> "$NDJSON_FILE"
done
start=$(date +%s%N)
curl -s -u "${ES_USER:-}:${ES_PASSWORD:-}" \
-X POST "${ES_ENDPOINT}/_bulk" \
-H "Content-Type: application/x-ndjson" \
--data-binary "@${NDJSON_FILE}" \
-o /tmp/bulk_resp.json
end=$(date +%s%N)
rm -f "$NDJSON_FILE"
echo "cost_ms=$(( (end - start) / 1000000 ))"
压测结果
| 指标 | 优化前 | 优化后 | 变化 |
|---|---|---|---|
| bulk 5000条 耗时 | 2400ms | 380ms | -84.2% |
| ES 峰值 CPU | 85% | 64% | -21pp |
| Young GC 次数/分钟 | 12次 | 2次 | -83.3% |
| 单索引段数量 | 800+ | 8(rollover+forcemerge后) | -99% |
| P99 写入延迟 | 1.2s | 320ms | -73.3% |
| bulk 拒绝次数/分钟 | 270+ | 0 | -100% |
| Canal 消费积压 | 8分钟 | 2秒 | 追平 |
五、避坑指南(这5个坑我全踩过)
坑1:bulk size 不是越大越好
我一开始把 batch 从 5000 调到 20000,结果写入性能没提升,Young GC 暴增到每分钟 30 次,两个数据节点直接 FULL GC 各停顿 8 秒,从集群里掉了出去。Elasticsearch 的 bulk 是把整个 body 加载到内存再分发给分片的,5000~8000 条是安全区,超过 1 万条时堆内存压力陡增。正确姿势:定量压测,观察堆内存曲线。
坑2:refresh_interval 别盲目往大调
为了追写入性能,我把 refresh 调到 30s,结果业务方投诉"数据查不到"。后来在 10s 和 15s 之间做了灰度,10s 时写入性能损失仅 5%,查询可见性完全满足。如果你的业务要求写入秒级可见,refresh 最多给到 5s,再大就用 index buffer 和 translog 去优化,别拿 refresh 开刀。
坑3:translog 调 async 前想清楚能不能丢数据
translog 本来就是做崩溃恢复的。改成 async 后,如果机器断电,ES 会丢失掉最近 sync_interval 内的数据。我用 5s 的窗口,靠 Kafka 重放补数据;但金融账单系统里推荐保持默认的 request 模式,一次请求一次 fsync,慢一点但不会丢。
坑4:磁盘水位线默认值会坑你
ES 默认的 flood_stage 是 95%,但 high watermark 85% 时就已经拒绝写了。我们一个集群数据盘中短期冲到 84%,ES 就在集群级别禁用了新索引分配,所有索引变 yellow。排查半天,最后发现是磁盘预警的误报。把水位线调到自己有数的范围,配好 Prometheus 告警再动手。
坑5:forcemerge 会让索引短暂不可写
ILM 里 forcemerge 历史索引时,如果这个索引还在收写入,会触发写入失败。我跑完 ILM 策略后发现数据有缺口,就是因为 rollover 当天索引被误 merge 了。解决方法是:加 index.blocks.write=true 锁住索引再 forcemerge,或者强制 require 该索引在 rollover 之后才进入 warm 阶段。ILM 的 warm phase 设计就是给只读索引用的,务必保证你的 hot 阶段正确 rollover。
六、最后的建议
Elasticsearch 调优不是改一两个参数就完事的。它是一条链路:mapping 决定了索引体积和 CPU 消耗,translog 和 refresh 决定了写入路径的 I/O 模式,ILM 决定了长期运行后的文件结构,客户端写入方式决定了 ES 线程池的工作负载。
照着这个顺序排查:先解决 mapping 浪费,再控制段数量,再优化写入客户端,最后再调集群参数。别一上来就加节点。
如果你正在经历写入性能瓶颈,用文章里的模板和代码先在测试环境跑一轮压测,大概率能找到你集群里的问题。有问题欢迎交流。