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

Go中间件使用Elasticsearch Go Client索引文档时遇EOF JSON格式错误

问题原因分析

  1. 为何文档已成功索引仍报"Invalid JSON format: EOF"?
    HTTP请求体的r.Body是一次性的io.ReadCloser,你的中间件通过json.NewDecoder(r.Body).Decode(&doc)已经将请求体的数据流读取完毕。当中间件调用next.ServeHTTP(w, r)把请求传给主处理器时,主处理器尝试再次读取请求体,此时数据流已到末尾,会返回EOF错误导致JSON解析失败。而Elasticsearch的索引操作基于已解码到内存的doc对象,不受请求体是否被读完的影响,因此文档能成功写入ES。

解决方案

要让中间件和主处理器都能读取请求体,需先将请求体内容缓存到内存,再构造可重复读取的Body替换原请求体:

修改后的中间件代码:

package middlewares

import (
    "bytes"
    "encoding/json"
    "io"
    "net/http"

    "github.com/elastic/go-elasticsearch/v8"
    "github.com/elastic/go-elasticsearch/v8/typedapi/types"
    "github.com/google/uuid"
)

// IndexDocumentMiddleware creates a middleware to index documents into Elasticsearch
func IndexDocumentMiddleware(es *elasticsearch.TypedClient) func(http.Handler) http.Handler {
    return func(next http.Handler) http.Handler {
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
            ctx := r.Context()
            
            // 读取请求体所有字节并缓存
            bodyBytes, err := io.ReadAll(r.Body)
            if err != nil {
                http.Error(w, "Error reading request body", http.StatusBadRequest)
                return
            }
            // 关闭原Body避免资源泄漏
            _ = r.Body.Close()
            
            // 用缓存的字节解码JSON
            var doc map[string]interface{}
            if err := json.Unmarshal(bodyBytes, &doc); err != nil {
                http.Error(w, "Error parsing request body", http.StatusBadRequest)
                return
            }

            var indexName string
            if typeName, ok := doc["type"].(string); ok {
                indexName = typeName
            } else {
                http.Error(w, "Error: 'type' is not a string or is missing", http.StatusBadRequest)
                return
            }

            existsRes, err := es.Indices.Exists(indexName).Do(ctx)
            if err != nil {
                http.Error(w, "Error existsRes: "+err.Error(), http.StatusInternalServerError)
                return
            }

            if !existsRes {
                _, err := es.Indices.Create(indexName).Mappings(types.NewTypeMapping()).Do(ctx)
                if err != nil {
                    http.Error(w, "Error creating index: "+err.Error(), http.StatusInternalServerError)
                    return
                }
            }

            docID := uuid.New().String()

            _, err = es.Index(indexName).
                Id(docID).
                Document(doc).Do(ctx)
            if err != nil {
                http.Error(w, "Error indexing document: "+err.Error(), http.StatusInternalServerError)
                return
            }
            
            // 重新构造可重复读取的Body,替换原请求体
            r.Body = io.NopCloser(bytes.NewBuffer(bodyBytes))
            // 重置Content-Length,保证后续处理器正确识别请求体大小
            r.ContentLength = int64(len(bodyBytes))

            next.ServeHTTP(w, r)
        })
    }
}

关键修改说明

  • 用io.ReadAll读取请求体到字节数组,缓存全部内容
  • 改用json.Unmarshal解码JSON,避免直接消耗请求体数据流
  • 用io.NopCloser(bytes.NewBuffer(bodyBytes))重新构造请求体,让主处理器可重复读取
  • 重置r.ContentLength为缓存字节的长度,避免后续处理器误解请求体大小

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 11:53:19