BulkProcessor写入ES出现重复数据(已关闭重试)问题求助
问题根因
你配置的BackoffPolicy.noBackoff()仅关闭了BulkProcessor组件自身的重试逻辑,没有覆盖到其他会触发重复请求的链路,这是重复写入的核心原因:
- RestHighLevelClient底层默认自带重试机制,当ES负载高、响应慢触发socket超时、连接断开等异常时,客户端会自动重新发送当前请求,这个过程BulkProcessor完全无感知,不会做去重判断。
- 你当前配置了
setConcurrentRequests(4)的异步并发提交,高负载下请求堆积时,异步回调的超时/异常通知延迟也可能导致同个批次被多次提交。 - 如果你写入时没有给文档指定唯一业务ID,依赖ES自动生成ID的逻辑,就算请求被重复发送,ES每次都会生成新的文档ID写入,直接造成重复数据。
解决方案
你的核心诉求是零重复写入,可接受少量数据丢失,按以下步骤调整即可:
- 全链路关闭重试:除了保留BulkProcessor的
noBackoff()配置,还要修改RestHighLevelClient的初始化参数,关闭客户端默认的重试策略,同时适当调短socket超时、连接超时时间,避免长时间等待触发链路重试。 - 加写入幂等兜底:所有批量写入操作统一用
Create模式代替普通Index模式,且必须为每个文档传入全局唯一的业务主键。Create模式仅当文档ID不存在时才会执行写入,就算重复请求到达ES,也会触发版本冲突直接跳过写入,从服务端彻底堵死重复写入的可能。注意绝对不要使用ES自动生成文档ID的逻辑,否则幂等配置不生效。 - 降低并发压力:把
setConcurrentRequests的值从4调整为0(即同步提交模式,上一个批次处理完成再发送下一个批次),避免多请求并发堆积压垮ES导致超时;可以根据ES集群承载能力适当调小单次批量的条数、体积阈值,减少单请求处理时长。 - 回调逻辑禁止重试:两个
afterBulk回调里不要加任何失败重试逻辑,只要捕获到写入异常、或者返回结果中存在写入失败的条目,直接记录日志后丢弃即可,符合你允许少量丢数的要求。
调整后的参考配置
@Bean public BulkProcessor bulkProcessor() { RestHighLevelClient client = restHighLevelClient(); BiConsumer<BulkRequest, ActionListener<BulkResponse>> bulkConsumer = (request, bulkListener) -> client.bulkAsync(request, RequestOptions.DEFAULT, bulkListener); return BulkProcessor.builder(bulkConsumer, new BulkProcessor.Listener() { @Override public void beforeBulk(long l, BulkRequest bulkRequest) { // 可加请求打点日志,不要做请求修改逻辑 } @Override public void afterBulk(long l, BulkRequest bulkRequest, BulkResponse bulkResponse) { // 存在写入失败条目直接记录日志丢弃,不要重试 if (bulkResponse.hasFailures()) { log.error("bulk write failed, drop data, msg:{}", bulkResponse.buildFailureMessage()); } } @Override public void afterBulk(long l, BulkRequest bulkRequest, Throwable throwable) { // 异常场景直接记录日志丢弃,不要重试 log.error("bulk write throw error, drop data", throwable); } }).setBulkActions(100) // 适当调小单次批量条数 .setBulkSize(new ByteSizeValue(1, ByteSizeUnit.MB)) // 适当调小单次批量体积 .setFlushInterval(TimeValue.timeValueSeconds(5)) .setConcurrentRequests(0) // 改为同步提交,避免并发堆积 .setBackoffPolicy(BackoffPolicy.noBackoff()) // 保持BulkProcessor层面无重试 .build(); } // 初始化RestHighLevelClient时记得关闭客户端默认重试 // 配置RequestConfig时把setSocketTimeout、setConnectTimeout设为合理短值,同时自定义HttpClient的重试策略为NO_RETRY
注意:不要依赖客户端侧的配置100%避免重复请求,服务端侧的
Create+唯一业务ID配置是必加的兜底,哪怕上游出现重复发请求的bug,ES侧也不会写入重复数据。
内容的提问来源于stack exchange,提问作者Lawrence
相关产品推荐
相关产品推荐

