Go语言Elasticsearch客户端单次初始化后无法索引多条记录
问题根因
你代码的核心问题是将defer res.Body.Close()放在了for循环内部:
- Go语言的
defer语句会在当前外层函数退出时才执行,而非循环迭代结束时执行。你循环100次的过程中,前99次请求的响应Body都没有被主动关闭,底层HTTP连接池的连接会一直被占用无法释放复用 - Elasticsearch客户端默认复用HTTP连接池,当连接被占满后,后续请求会一直等待空闲连接,最终触发超时
- 每次新建客户端的场景下相当于每次新建独立的连接池,不会受旧连接占用的影响,curl本身也是单次请求结束自动释放连接,所以这两个场景不会触发问题
修复方案
只需要把defer res.Body.Close()替换成每次循环迭代结束主动关闭响应Body即可,修改后代码如下:
package main import ( "context" "encoding/json" "log" "strconv" "strings" "time" "github.com/elastic/go-elasticsearch/v6" "github.com/elastic/go-elasticsearch/v6/esapi" ) func main() { cfg := elasticsearch.Config{ Addresses: []string{ "http://172.31.1.85:9200", "http://172.31.1.85:9201", }, } es, err := elasticsearch.NewClient(cfg) if err != nil { log.Fatalln(err) } for i := 0; i < 100; i++ { time.Sleep(20 * time.Second) body := `{"test": "test"}` req := esapi.IndexRequest{ Index: "test", DocumentType: "", DocumentID: strconv.Itoa(i), Body: strings.NewReader(body), Refresh: "true", } // 执行请求 res, err := req.Do(context.Background(), es) if err != nil { log.Printf("Error getting response: %s\n", err) return } if res.IsError() { log.Printf("[%s] Error indexing document ID=%d\n", res.Status(), i) res.Body.Close() continue } // 处理正常响应 var r map[string]interface{} if err := json.NewDecoder(res.Body).Decode(&r); err != nil { log.Printf("Error parsing the response body: %s", err) } else { log.Printf("[%s] %s; version=%d", res.Status(), r["result"], int(r["_version"].(float64))) } // 处理完响应后关闭Body释放连接 res.Body.Close() } }
额外优化建议
- 原代码里日志打印错误索引ID为固定值100,已经修改为循环变量
i - 原Body字符串带多余缩进和换行,建议去掉避免不必要的传输开销
- 客户端初始化失败建议直接用
log.Fatalln退出,避免后续空指针异常
内容的提问来源于stack exchange,提问作者Baruch
相关产品推荐
相关产品推荐

