🎯 项目背景 在现代内容平台中,搜索功能是用户获取信息的核心入口。传统的数据库 LIKE 查询在面对海量数据时性能急剧下降,而全文搜索引擎 Elasticsearch 能够提供毫秒级的搜索响应。本文将详细介绍如何使用 Go + Elasticsearch 构建一个功能完整的智能搜索系统。 ✨ 系统...
Go + Elasticsearch 实现智能搜索功能
发布时间: 2025-07-29 (a year ago)
GOES

🎯 项目背景

在现代内容平台中,搜索功能是用户获取信息的核心入口。传统的数据库 LIKE 查询在面对海量数据时性能急剧下降,而全文搜索引擎 Elasticsearch 能够提供毫秒级的搜索响应。本文将详细介绍如何使用 Go + Elasticsearch 构建一个功能完整的智能搜索系统。

✨ 系统特性

🚀 高性能搜索

  • 全文检索:支持标题、内容、摘要的多字段搜索
  • 权重排序:标题权重 3x,内容权重 2x,摘要权重 1x
  • 高亮显示:搜索关键词自动高亮标记
  • 毫秒响应:平均响应时间 < 50ms

🎯 智能推荐

  • 标签相似度:基于 Jaccard 系数计算文章相似性
  • 用户行为:点赞、收藏行为影响推荐权重
  • 时间衰减:新文章获得更高推荐权重
  • 缓存优化:Redis 缓存推荐结果,提升性能

🔄 数据同步

  • 实时同步:文章 CRUD 操作自动同步到 ES
  • 批量操作:支持批量删除和更新
  • 错误重试:网络异常自动重试机制
  • 定时全量:定时任务保证数据一致性

🏗️ 系统架构

核心组件

graph TD A[Web API] --> B[Service] B --> C[Elasticsearch] A --> D[MySQL] B --> E[Redis] C --> F[定时任务]

数据流转

  1. 写入流程:API → MySQL → ES 同步
  2. 搜索流程:API → ES 查询 → 结果处理
  3. 推荐流程:API → 相似度计算 → Redis 缓存

💡 核心实现

1. 智能搜索查询构建
Go 复制代码
func BuildSearchQuery(keyword, sortField, sortOrder, from, size string, tags []string, category string) *bytes.Buffer {
    var queryBody map[string]interface{}
    
    if keyword == "" {
        // 全量查询
        queryBody = map[string]interface{}{
            "query": map[string]interface{}{
                "match_all": map[string]interface{}{},
            },
        }
    } else {
        // 多字段权重搜索
        queryBody = map[string]interface{}{
            "query": map[string]interface{}{
                "multi_match": map[string]interface{}{
                    "query":  keyword,
                    "fields": []string{"title^3", "content^2", "abstract"},
                    "type":   "best_fields",
                },
            },
        }
    }

    // 构建完整查询
    query := map[string]interface{}{
        "query": queryBody["query"],
        "highlight": map[string]interface{}{
            "pre_tags":  []string{"<em class='highlight'>"},
            "post_tags": []string{"</em>"},
            "fields": map[string]interface{}{
                "title":   map[string]interface{}{},
                "content": map[string]interface{}{},
            },
        },
        "from": from,
        "size": size,
    }

    // 添加过滤条件
    if category != "" || len(tags) > 0 {
        filters := []map[string]interface{}{}
        
        if category != "" {
            filters = append(filters, map[string]interface{}{
                "term": map[string]interface{}{"category": category},
            })
        }
        
        if len(tags) > 0 {
            filters = append(filters, map[string]interface{}{
                "terms": map[string]interface{}{"tags": tags},
            })
        }

        query["query"] = map[string]interface{}{
            "bool": map[string]interface{}{
                "must":   queryBody["query"],
                "filter": filters,
            },
        }
    }

    var buf bytes.Buffer
    json.NewEncoder(&buf).Encode(query)
    return &buf
}
2. 数据同步
Go 复制代码
func SyncToES(article models.ArticleModel) error {
    if !global.Config.Elasticsearch.Enable {
        return nil
    }

    // 构造文档结构
    doc := map[string]any{
        "title":          article.Title,
        "content":        article.Content,
        "abstract":       article.Abstract,
        "tags":           article.Tags,
        "category":       article.Category,
        "created_at":     article.CreatedAt.Format("2006-01-02 15:04:05"),
        "digg_count":     article.DiggCount,
        "look_count":     article.LookCount,
        "collects_count": article.CollectsCount,
        "comment_count":  article.CommentCount,
    }

    docData, err := json.Marshal(doc)
    if err != nil {
        return err
    }

    req := esapi.IndexRequest{
        Index:      "articles",
        DocumentID: fmt.Sprintf("%d", article.ID),
        Body:       strings.NewReader(string(docData)),
        Refresh:    "wait_for", // 立即可搜索
    }

    // 带重试的执行
    // 这里最好补充熔断机制,如果服务器宕机和超载,重试对服务器是雪上加霜
    return retry.DoWithRetry(3, 1*time.Second, func() error {
        res, err := req.Do(context.Background(), global.Elasticsearch)
        if err != nil {
            return err
        }
        defer res.Body.Close()

        if res.IsError() {
            body, _ := io.ReadAll(res.Body)
            return fmt.Errorf("ES操作失败: %s | %s", res.Status(), string(body))
        }
        return nil
    })
}
3. 智能推荐算法
Go 复制代码
func GetRecommendations(articleID uint, userID uint, limit int) ([]ArticleSimilarity, error) {
    // 缓存检查
    cacheKey := fmt.Sprintf("article_recommend:%d:%d:%d", articleID, userID, limit)
    if val, err := global.Redis.Get(context.Background(), cacheKey).Result(); err == nil {
        var cachedData []ArticleSimilarity
        if err := json.Unmarshal([]byte(val), &cachedData); err == nil {
            return cachedData, nil
        }
    }

    var currentArticle models.ArticleModel
    if err := global.DB.First(&currentArticle, articleID).Error; err != nil {
        return nil, err
    }

    var allArticles []models.ArticleModel
    global.DB.Where("id != ?", articleID).Find(&allArticles)

    recommendations := make([]ArticleSimilarity, 0, len(allArticles))

    for _, article := range allArticles {
        // 基础相似度(Jaccard系数)
        score := calculateTagSimilarity(currentArticle.Tags, article.Tags)

        // 用户行为加权
        if userID > 0 {
            // 点赞行为加权
            if global.Redis.SIsMember(context.Background(), 
                fmt.Sprintf("article_digg_users:%d", article.ID), userID).Val() {
                score *= 1.3
            }

            // 收藏行为加权
            var collectExists int64
            global.DB.Model(&models.UserCollectModel{}).
                Where("user_id = ? AND article_id = ?", userID, article.ID).
                Count(&collectExists)
            if collectExists > 0 {
                score *= 1.5
            }
        }

        // 时间衰减因子
        daysOld := time.Since(article.CreatedAt).Hours() / 24
        timeFactor := math.Exp(-0.1 * daysOld)
        score *= 1 + 0.5*timeFactor

        recommendations = append(recommendations, ArticleSimilarity{
            ArticleID: article.ID,
            Score:     score,
        })
    }

    // 按分数排序
    sort.Slice(recommendations, func(i, j int) bool {
        return recommendations[i].Score > recommendations[j].Score
    })

    if limit > len(recommendations) {
        limit = len(recommendations)
    }

    result := recommendations[:limit]

    // 缓存结果
    if data, err := json.Marshal(result); err == nil {
        global.Redis.Set(context.Background(), cacheKey, data, 30*time.Minute)
    }

    return result, nil
}

// Jaccard 相似度计算
func calculateTagSimilarity(currentTags, targetTags []string) float64 {
    tagSet := make(map[string]bool)
    for _, tag := range currentTags {
        tagSet[tag] = true
    }

    commons := 0
    for _, tag := range targetTags {
        if tagSet[tag] {
            commons++
        }
    }

    // Jaccard 系数 = 交集 / 并集
    return float64(commons) / float64(len(currentTags)+len(targetTags)-commons)
}
4. 搜索结果解析
Go 复制代码
type SearchResult struct {
    Total int
    Hits  []SearchHit
}

type SearchHit struct {
    ID        string                 `json:"id"`
    Source    map[string]interface{} `json:"_source"`
    Highlight map[string][]string    `json:"highlight"`
}

func ParseSearchResponse(body io.Reader) (*SearchResult, error) {
    var response struct {
        Hits struct {
            Total struct {
                Value int `json:"value"`
            } `json:"total"`
            Hits []struct {
                ID        string                 `json:"_id"`
                Source    map[string]interface{} `json:"_source"`
                Highlight map[string][]string    `json:"highlight"`
            } `json:"hits"`
        } `json:"hits"`
    }

    if err := json.NewDecoder(body).Decode(&response); err != nil {
        return nil, err
    }

    result := &SearchResult{
        Total: response.Hits.Total.Value,
        Hits:  make([]SearchHit, len(response.Hits.Hits)),
    }

    for i, hit := range response.Hits.Hits {
        result.Hits[i] = SearchHit{
            ID:        hit.ID,
            Source:    hit.Source,
            Highlight: hit.Highlight,
        }
    }

    return result, nil
}

📊 性能优化

1. 搜索性能

  • 索引优化:合理设置分片和副本数量
  • 查询缓存:ES 内置查询结果缓存
  • 分页优化:使用 from/size 进行高效分页

2. 推荐性能

  • Redis 缓存:推荐结果缓存 30 分钟
  • 异步计算:后台定时更新推荐数据
  • 批量处理:批量计算相似度减少数据库查询

3. 同步性能

  • 批量操作:支持批量索引和删除
  • 异步同步:写入 MySQL 后异步同步 ES
  • 错误重试:网络异常自动重试机制

🛡️ 可靠性保障

1. 数据一致性

Go 复制代码
// 定时全量同步任务
func FullSyncToES() {
    var articles []models.ArticleModel
    global.DB.Find(&articles)

    for _, article := range articles {
        if err := SyncToES(article); err != nil {
            global.Log.Error("全量同步失败", 
                zap.Uint("id", article.ID), zap.Error(err))
        }
    }
}

2. 错误处理

Go 复制代码
// 重试机制
func retry.DoWithRetry(maxRetries int, delay time.Duration, fn func() error) error {
    for i := 0; i < maxRetries; i++ {
        if err := fn(); err == nil {
            return nil
        }
        if i < maxRetries-1 {
            time.Sleep(delay)
        }
    }
    return fmt.Errorf("重试 %d 次后仍然失败", maxRetries)
}

3. 监控告警

  • 搜索延迟监控:记录搜索响应时间
  • 同步失败告警:ES 同步失败自动告警
  • 索引健康检查:定期检查 ES 集群状态

🚀 使用示例

基础搜索

Go 复制代码
// 关键词搜索
query := BuildSearchQuery("Go语言", "created_at", "desc", "0", "10", nil, "")
// 执行搜索...

高级搜索

Go 复制代码
// 带标签和分类的搜索
tags := []string{"后端", "微服务"}
query := BuildSearchQuery("分布式", "digg_count", "desc", "0", "20", tags, "技术")

智能推荐

Go 复制代码
// 获取相关文章推荐
recommendations, err := GetRecommendations(articleID, userID, 5)

技术栈:Go + Elasticsearch + Redis + MySQL
适用场景:内容平台、电商搜索、知识库
核心算法:TF-IDF、Jaccard 相似度、时间衰减