🎯 项目背景
在现代内容平台中,搜索功能是用户获取信息的核心入口。传统的数据库 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[定时任务]
数据流转
- 写入流程:API → MySQL → ES 同步
- 搜索流程:API → ES 查询 → 结果处理
- 推荐流程: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(¤tArticle, 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 相似度、时间衰减