当前位置:首页 > 随笔 > 正文内容

Elasticsearch Java批量索引优化方案

访客 随笔 2026年7月26日 1
在进行文档写入操作时,频繁的单条请求会显著影响性能。当需要处理海量数据时,采用逐条插入的方式显然不具可行性。
传统批量操作示例:

POST /_bulk
{ "delete": { "_index": "website", "_type": "blog", "_id": "123" }} 
{ "create": { "_index": "website", "_type": "blog", "_id": "123" }}
{ "title":    "My first blog post" }
{ "index":  { "_index": "website", "_type": "blog" }}
{ "title":    "My second blog post" }
{ "update": { "_index": "website", "_type": "blog", "_id": "123", "_retry_on_conflict" : 3} }
{ "doc" : {"title" : "My updated blog post"} }
标准批量处理实现:

@Test
public void batchInsert() throws IOException {
    BulkRequestHandler bulkHandler = client.prepareBulk();

    bulkHandler.add(client.prepareIndex("social", "post", "1001")
        .setSource(jsonBuilder()
            .startObject()
            .field("author", "李四")
            .field("timestamp", new Date())
            .field("content", "关于网络事件的深度分析")
            .endObject()
        )
    );

    bulkHandler.add(client.prepareIndex("social", "post", "1002")
        .setSource(jsonBuilder()
            .startObject()
            .field("author", "王五")
            .field("timestamp", new Date())
            .field("content", "体育赛事的实时报道")
            .endObject()
        )
    );

    BulkResponse response = bulkHandler.get();
    if (response.hasFailures()) {
        // 处理失败项
    }
}
该方案存在局限性,如批量阈值、数据量限制、执行间隔等参数均需手动配置。以下为增强型处理方案:
优化批量处理实现:

@Test
public void optimizedBatch() throws Exception {
    BulkProcessor processor = BulkProcessor.builder(client, new BulkProcessor.Listener() {
        public void beforeBulk(long id, BulkRequest request) {
            System.out.println("准备提交" + request.numberOfActions() + "条记录");
        }

        public void afterBulk(long id, BulkRequest request, BulkResponse response) {
            System.out.println("成功处理" + request.numberOfActions() + "条记录");
        }

        public void afterBulk(long id, BulkRequest request, Throwable error) {
            System.out.println("处理" + request.numberOfActions() + "条记录失败");
        }
    })
    .setBulkActions(5000)
    .setBulkSize(new ByteSizeValue(500, ByteSizeUnit.MB))
    .setFlushInterval(TimeValue.timeValueSeconds(10))
    .setConcurrentRequests(2)
    .setBackoffPolicy(BackoffPolicy.exponentialBackoff(TimeValue.timeValueMillis(50), 5))
    .build();

    Map data = new HashMap<>();
    data.put("text", "异步批量测试数据");
    processor.add(new IndexRequest("news", "article", "2001").source(data));
    processor.add(new IndexRequest("news", "article", "2002").source(data));
    
    processor.flush();
    processor.awaitClose(5, TimeUnit.MINUTES);
}

相关文章

可以按小时收费的VPS

很多 VPS 提供商都支持 按小时计费(hourly billing),想短期试用 / 临时搭建节点、测试网络、短期项目等场景非常合适。下面是当前最主流且靠谱的按小时 VPS 选项,分别按不同需求场景整理: 1. Vultr(全球节点,包括日本) 按小时计费 可选机房:东京 / 大阪 / 洛杉矶 / 法兰克福 / 伦敦 … 支持 PayPal(部分情况),但更常用信用卡/PayPal+卡价格参考$...

在 iPhone 上下载国外App

地区/国家限制App Store 会根据 Apple ID 的国家或地区限制应用下载。如果你的 Apple ID 绑定的是中国大陆,就可能无法下载 OpenAI 官方的 ChatGPT 应用,因为它在大陆 App Store 不上架。解决办法:换成美国、加拿大、香港等地区的 Apple ID。或者在现有 Apple ID 上更改地区。注册一个国外 Apple ID(推荐)比如注册 美国区 Appl...

Node.js 中的异步编程:回调与 Promise

Node.js 是一个基于 JavaScript 构建的单线程、非阻塞运行环境,它通过异步编程机制来高效处理多个操作。在执行如文件读取、API 请求或数据库查询等任务时,Node.js 不会等待这些操作完成,而是使用回调函数和 Promise 来避免阻塞主线程。 回调方式实现异步 那么当异步操作完成后,Node.js 如何知道接下来要做什么呢?这就要用到 回调函数(callback)。 回调本质上...

Selenium自动化测试入门指南

Selenium自动化测试入门指南

什么是自动化测试? 自动化测试是指利用软件工具自动执行测试用例,模拟用户操作,如打开网页、点击链接、输入文本等,并验证结果是否符合预期。 其主要优点包括: 大幅减少人工成本 测试速度快 可以在非工作时间运行 支持持续集成和交付 然而,它也存在一些局限性,例如开发成本较高、不适合快速变化的项目、依赖稳定的UI界面等。 自动化测试的应用条件 适合引入自动化测试的情况包括: 手动测试耗时且需要大量...

MariaDB Galera集群故障快速恢复指南

OpenStack控制节点采用三节点MariaDB Galera集群架构。当数据库集群因故障重启时,有时会出现Galera集群无法正常启动的问题。虽然有多种方法可以恢复数据库服务,但如何实现快速启动同时确保数据完整性呢? 通过分析日志发现,MariaDB Galera集群节点宕机时会在日志中输出以下信息: [Note] WSREP: 新集群视图:全局状态: 874d8e7e-5980-11e8-8...

Android 中 EventBus 的通信机制与实现原理深度解析

EventBus 核心设计思想 EventBus 是一个基于观察者模式的事件总线框架,广泛应用于 Android 平台以实现组件解耦。它通过中心化的消息分发机制,使不同层级、不同线程的对象能够以"发布-订阅"方式通信,避免了传统接口回调或广播带来的强依赖问题。 核心角色说明 事件(Event):任意 Java 对象,作为数据载体,如网络状态变更通知、用户登录信息等。 发布者(Publi...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。