Elasticsearch6.8异步delete_by_query生成无限运行任务问题咨询
Elasticsearch 6.8 重复delete_by_query问题解决方案
问题1:无法避免删除查询重叠时的大批量数据正确删除方案
- 新增任务提交前去重逻辑
提交删除任务前,先调用GET _tasks?detailed=true&actions=*/delete/byquery接口查询当前集群正在运行的delete_by_query任务,解析任务描述中的查询参数,匹配ownerId和对应的时间范围,如果已有相同参数的任务在运行,直接跳过本次提交,从根源避免重复任务生成。 - 优化delete_by_query执行参数
- 移除
Refresh("true")配置:批量删除场景下无需每执行一个子批次就刷新索引,待整个删除任务完成后手动触发一次索引refresh即可,可大幅降低集群IO开销。 - 增加
scroll_size配置:默认值为1000,可根据集群性能调整到2000~5000,减少索引扫描的请求次数。 - 增加
requests_per_second限流配置:可设置为1000~5000不等,避免删除任务占满集群带宽影响业务。
- 移除
- 优先使用时间分片索引方案
如果数据是按时间时序写入,建议按天/周维度拆分索引,超出保留期的老旧数据直接删除整个索引,性能比delete_by_query高数个量级,完全不会产生版本冲突问题。 - 补充任务取消能力
提交异步任务时会返回对应taskId,将taskId持久化存储后,遇到异常可调用POST _tasks/{task_id}/_cancel接口直接终止对应删除任务,无需重建索引。
问题2:多个重复delete_by_query任务导致集群异常的原因
Elasticsearch 6.8的delete_by_query本质是先通过scroll接口扫描匹配文档,再逐批发送删除请求,每个删除请求都会触发版本校验流程:
- 多个相同参数的delete_by_query任务并行运行时,会先扫描到同一批未删除的文档,先后发送删除请求,第一个请求删除成功后文档版本号会上涨,后续请求删除同一文档时就会触发版本冲突。
- 你开启了
ProceedOnVersionConflict()配置,任务遇到冲突不会终止,会继续扫描下一批文档,此时大部分待删除文档已经被其他任务标记为删除,扫描到的有效待删除文档占比极低,但任务仍然会持续消耗CPU、内存和节点间带宽做全量扫描和冲突校验。 - 当重复提交的任务数量足够多时,集群所有节点的资源都被这些无意义的扫描和冲突校验占满,正常业务请求无法拿到资源就会出现超时、集群挂起的现象,并非进入生成新任务的死循环,而是你提交的大量重复任务都卡在低效运行状态,长时间无法结束。
优化后代码逻辑参考
// 先校验是否有相同参数的删除任务正在运行 tasks, err := esClient.TasksList().Detailed(true).Actions("*delete/byquery").Do(context.Background()) if err != nil { return err } hasRunningTask := false for _, task := range tasks.Tasks { // 匹配任务描述中的查询参数 desc := task.Description if strings.Contains(desc, fmt.Sprintf("ownerId:%s", ownerId)) && strings.Contains(desc, fmt.Sprintf("timestamp<=%d", to)) && strings.Contains(desc, fmt.Sprintf("timestamp>=%d", from)) { hasRunningTask = true break } } if hasRunningTask { // 已有相同任务运行,跳过提交 return nil } // 提交优化参数后的删除任务 query := elastic.NewBoolQuery().Must( elastic.NewTermQuery("ownerId", ownerId), elastic.NewRangeQuery("timestamp").Lte(to).Gte(from), ) taskId, err := elastic.NewDeleteByQueryService(esClient). Index("my-index"). Query(query). ProceedOnVersionConflict(). ScrollSize(3000). RequestsPerSecond(1000). Refresh("false"). DoAsync(context.Background()) if err != nil { return err } // 可选:将taskId持久化到存储,用于后续任务状态查询或取消
内容的提问来源于stack exchange,提问作者SubChord
相关产品推荐
相关产品推荐

