You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Golang中实现Elasticsearch索引仅保留10条文档并自动替换

嘿,这个需求我之前也帮人处理过,有两种靠谱的实现方式,一种是让Elasticsearch自己来维护(用索引生命周期管理ILM),另一种是在Go代码里手动控制文档数量。我给你详细拆解下:

方案一:让Elasticsearch自动维护(ILM方式)

这种方式不用在代码里写太多逻辑,交给ES自己管就行,适合不想折腾代码的场景。核心思路是用ILM的rollover功能,当索引满10条文档时,自动切换到新索引,然后删掉旧索引,同时用别名来统一访问最新的索引。

步骤如下:

  1. 创建生命周期策略:定义一个策略,指定索引满10条时触发rollover,然后立即删除旧索引
curl -X PUT "http://localhost:9200/_ilm/policy/keep-10-docs-policy" -H "Content-Type: application/json" -d'
{
  "policy": {
    "phases": {
      "hot": {
        "actions": {
          "rollover": {
            "max_docs": 10
          }
        }
      },
      "delete": {
        "min_age": "0ms",
        "actions": {
          "delete": {}
        }
      }
    }
  }
}
'
  1. 创建索引模板:把上面的策略绑定到索引模板上,同时设置别名
curl -X PUT "http://localhost:9200/_index_template/keep-10-docs-template" -H "Content-Type: application/json" -d'
{
  "index_patterns": ["my-index-*"],
  "template": {
    "settings": {
      "index.lifecycle.name": "keep-10-docs-policy",
      "index.lifecycle.rollover_alias": "my-index"
    }
  },
  "priority": 100
}
'
  1. 初始化第一个索引:创建第一个带别名的索引,后续写入都通过别名操作
curl -X PUT "http://localhost:9200/my-index-000001" -H "Content-Type: application/json" -d'
{
  "aliases": {
    "my-index": {
      "is_write_index": true
    }
  }
}
'

之后你的Go代码只需要往my-index这个别名写入文档就好,ES会自动在索引满10条时切换到新索引,旧索引会被删除,始终保留最新的10条文档。

方案二:Go代码手动控制(更灵活)

如果不想改ES配置,或者需要在删除旧文档前做一些自定义操作,那手动控制是更好的选择。核心逻辑很简单:每次写入新文档前,先检查当前索引的文档数,要是已经有10条了,就删掉最早的那条,再写入新的。

实现步骤(附代码):

首先安装官方的Go ES客户端:

go get github.com/elastic/go-elasticsearch/v8

然后是完整的示例代码:

package main

import (
	"context"
	"encoding/json"
	"fmt"
	"log"
	"time"

	"github.com/elastic/go-elasticsearch/v8"
	"github.com/elastic/go-elasticsearch/v8/esapi"
)

func main() {
	// 配置ES客户端,替换成你的ES地址和认证信息
	cfg := elasticsearch.Config{
		Addresses: []string{"http://localhost:9200"},
		// Username: "elastic",
		// Password: "your-password",
	}

	es, err := elasticsearch.NewClient(cfg)
	if err != nil {
		log.Fatalf("创建ES客户端失败: %s", err)
	}

	// 要操作的索引名
	indexName := "your-target-index"

	// 准备要写入的新文档,记得带一个创建时间字段用来排序
	newDoc := map[string]interface{}{
		"content":    "这是第11条测试文档",
		"created_at": time.Now().UTC().Format(time.RFC3339),
	}

	// 1. 查询当前索引的文档数量
	countResp, err := es.Count(
		es.Count.WithContext(context.Background()),
		es.Count.WithIndex(indexName),
	)
	if err != nil {
		log.Fatalf("统计文档数失败: %s", err)
	}
	defer countResp.Body.Close()

	var countResult map[string]interface{}
	if err := json.NewDecoder(countResp.Body).Decode(&countResult); err != nil {
		log.Fatalf("解析统计响应失败: %s", err)
	}
	currentCount := int(countResult["count"].(float64))
	fmt.Printf("当前索引文档数: %d\n", currentCount)

	// 2. 如果文档数达到10,删除最早的文档
	if currentCount >= 10 {
		// 搜索最早的文档,按created_at升序取第一条,只返回文档ID
		searchResp, err := es.Search(
			es.Search.WithContext(context.Background()),
			es.Search.WithIndex(indexName),
			es.Search.WithBody(json.RawMessage(`{
				"size": 1,
				"sort": [{"created_at": "asc"}],
				"_source": false,
				"fields": ["_id"]
			}`)),
		)
		if err != nil {
			log.Fatalf("搜索最早文档失败: %s", err)
		}
		defer searchResp.Body.Close()

		var searchResult map[string]interface{}
		if err := json.NewDecoder(searchResp.Body).Decode(&searchResult); err != nil {
			log.Fatalf("解析搜索响应失败: %s", err)
		}

		hits := searchResult["hits"].(map[string]interface{})["hits"].([]interface{})
		if len(hits) > 0 {
			oldestDocID := hits[0].(map[string]interface{})["_id"].(string)
			fmt.Printf("即将删除最早的文档ID: %s\n", oldestDocID)

			// 执行删除操作
			deleteResp, err := es.Delete(
				indexName,
				oldestDocID,
				es.Delete.WithContext(context.Background()),
			)
			if err != nil {
				log.Fatalf("删除最早文档失败: %s", err)
			}
			defer deleteResp.Body.Close()

			if deleteResp.IsError() {
				log.Fatalf("删除请求失败: %s", deleteResp.String())
			}
			fmt.Println("成功删除最早的文档")
		}
	}

	// 3. 写入新文档
	docBytes, err := json.Marshal(newDoc)
	if err != nil {
		log.Fatalf("序列化文档失败: %s", err)
	}

	indexResp, err := es.Index(
		indexName,
		json.RawMessage(docBytes),
		es.Index.WithContext(context.Background()),
		es.Index.WithRefresh("wait_for"), // 可选,确保写入后立即可见
	)
	if err != nil {
		log.Fatalf("写入新文档失败: %s", err)
	}
	defer indexResp.Body.Close()

	if indexResp.IsError() {
		log.Fatalf("写入请求失败: %s", indexResp.String())
	}
	fmt.Println("新文档写入成功")
}

注意事项:

  • 一定要用可靠的排序字段:比如示例里的created_at,别用ES自动生成的_id,因为它不代表文档的创建顺序。
  • 并发场景要注意竞态问题:如果多个请求同时写入,可能会出现同时检测到文档数是9,然后都写入导致变成11条的情况。这种情况可以加分布式锁,或者改用ILM方案更省心。

两种方案各有优劣:ILM省心省力,适合不需要自定义逻辑的场景;手动控制更灵活,适合需要在删除前做额外处理的情况,你可以根据自己的需求选。

内容的提问来源于stack exchange,提问作者kish

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 11:20:22