设置Refresh: wait_for仍遇es_rejected_execution_exception的原因排查
问题分析与解决方案
你遇到的es_rejected_execution_exception本质是Elasticsearch的write线程池及队列被并发请求打满,这和Refresh: wait_for的机制并不冲突,下面具体拆解原因和修复方案:
一、为什么Refresh: wait_for挡不住这个错误?
Refresh: wait_for的作用是让当前bulk请求的修改在请求完成前被刷新到可搜索状态,但它不会限制你发起请求的速率。你开了200个协程同时创建BulkIndexer并发起写入,相当于同时有200个并发的bulk流,每个还带着wait_for的阻塞逻辑,直接把ES的write线程池(队列容量200)撑爆。- 错误日志里明确显示
queued tasks = 200,刚好触及队列上限,新的请求直接被拒绝。
二、代码中的关键问题
- 每个协程独立创建BulkIndexer:
push函数每次调用都新建一个BulkIndexer,200个协程就有200个BulkIndexer同时运行,完全没有控制写入节奏。BulkIndexer的设计目的是聚合请求、自动控制批量大小和并发,你这么用等于把它当成了单次bulk请求的包装,完全浪费了它的限流能力。 - 并发数远超ES处理上限:每个BulkIndexer设置
NumWorkers:1,200个协程就会产生200个并发写入线程,而你的ES write线程池最大只有32,线程和队列直接被占满。 Refresh: wait_for加重线程占用:带wait_for的bulk请求会在ES服务器端阻塞到下一次刷新(默认1秒),这会让write线程被占用更久,队列堆积速度更快。
三、修复方案
1. 全局复用BulkIndexer,统一控制并发
把BulkIndexer的创建移到协程外部,所有协程共享同一个实例,让它帮你控制写入速率和批量大小:
// 在esClient实例级别初始化BulkIndexer func (e *esClient) InitBulkIndexer() error { var err error e.bulkIndexer, err = esutil.NewBulkIndexer(esutil.BulkIndexerConfig{ Client: e.client, Refresh: "wait_for", NumWorkers: 8, // 设置和ES write线程池匹配的数量,建议8-16,不超过32 FlushBytes: 5 * 1024 * 1024, // 累计5MB数据自动刷新 FlushInterval: 1 * time.Second, // 每1秒自动刷新一次 OnError: func(ctx context.Context, err error) { fmt.Printf("received onError %s\n", err.Error()) }, }) return err } // 协程仅负责向全局BulkIndexer添加数据 func (e *esClient) push(data []esutil.BulkIndexerItem) (*esutil.BulkIndexerStats, error) { ctx := context.Background() for _, d := range data { if err := e.bulkIndexer.Add(ctx, d); err != nil { fmt.Printf("error adding data to indexer: %s\n", err) } } return &e.bulkIndexer.Stats(), nil } // 主流程示例 func main() { esCli := &esClient{client: ...} if err := esCli.InitBulkIndexer(); err != nil { panic(err) } var wg sync.WaitGroup wg.Add(200) // 分200个协程处理数据批次 for _, batch := range dataBatches { go func(b []esutil.BulkIndexerItem) { defer wg.Done() esCli.push(b) }(batch) } wg.Wait() // 所有数据添加完成后,关闭BulkIndexer等待所有请求收尾 if err := esCli.bulkIndexer.Close(context.Background()); err != nil { fmt.Printf("error closing indexer: %s\n", err) } }
2. 临时调整ES线程池配置(应急用,非首选)
如果暂时无法修改代码,可以临时调大write线程池的队列容量:
# elasticsearch.yml thread_pool.write.queue_size: 1000
但这只是缓解手段,过多的队列堆积会导致ES内存占用飙升,甚至引发OOM,不建议长期使用。
3. 移除不必要的Refresh: wait_for(业务允许的话)
如果你的业务不需要数据写入后立即可搜索,将Refresh设为默认值false,bulk请求不会阻塞等待刷新,能大幅降低write线程的占用时间,缓解线程池压力。
总结
核心问题是并发写入速率远超ES的处理能力,Refresh: wait_for仅保证请求完成后数据可搜索,但不会帮你限流。通过复用BulkIndexer、控制并发数、调整批量参数,才能从根本上解决线程池被打满的问题。
内容的提问来源于stack exchange,提问作者Ankit Kumar
相关产品推荐
相关产品推荐

