Files
sundynix-agentix/sundynix-mcp-go/internal/rag/bleve.go
T
Blizzard 79e834e8e9 feat(rag): 大文件 RAG 做深 —— 并发入库/图谱窗口化 + Bleve落盘 + file_id治理
几十万字文件从"广而浅"到准生产级:
- 向量化串行→并发分批保序(embedAll);几十万字上千块快数倍
- 图谱整篇喂LLM(爆上下文只抽开头)→窗口化并发抽(extractGraphWindowed),
  全覆盖;窗口/封顶/并发 env 可配;图谱可单配便宜模型(GRAPH_CHAT_*,未配回退主chat)
- Bleve 内存索引(重启即丢、三路退两路)→落盘 scorch(env BLEVE_PATH,失败退内存兜底)
- 下游键改稳定 file_id:Neo4j 关系打 file_id(实体仍 kb+name 共享);
  新增 kb_delete 工具 + Engine.DeleteDoc 级联删 Milvus/Bleve/Neo4j
- 单测:窗口化/去重/env可配/落盘持久/图谱模型回退

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-30 13:37:45 +08:00

160 lines
4.6 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package rag
import (
"fmt"
"hash/fnv"
"log"
"os"
"path/filepath"
"github.com/blevesearch/bleve/v2"
"github.com/blevesearch/bleve/v2/analysis/analyzer/keyword"
"github.com/blevesearch/bleve/v2/analysis/lang/cjk"
"github.com/blevesearch/bleve/v2/mapping"
"github.com/blevesearch/bleve/v2/search/query"
)
// bleveStore 是全文(BM25)检索路。落盘索引(scorch):随 ingest 写入并持久,进程重启不丢,
// 与 Milvus(向量)/Neo4j(图谱) 同为持久存储——保证"三路混合"重启后仍是三路。
// 落盘不可用时退回内存索引(功能在、重启会丢),好过全文路全哑。
type bleveStore struct {
idx bleve.Index
}
// bleveMappingtext 字段用 cjk 分词器(bigram,能切中文,否则默认标准分词器把整段中文当一个
// token → 中文全文检索永远 0 命中);kb/doc 用 keyword(不分词,保 TermQuery 精确过滤)。
func bleveMapping() mapping.IndexMapping {
text := bleve.NewTextFieldMapping()
text.Analyzer = cjk.AnalyzerName
kw := bleve.NewTextFieldMapping()
kw.Analyzer = keyword.Name
doc := bleve.NewDocumentMapping()
doc.AddFieldMappingsAt("text", text)
doc.AddFieldMappingsAt("kb", kw)
doc.AddFieldMappingsAt("doc", kw)
im := bleve.NewIndexMapping()
im.DefaultMapping = doc
return im
}
// blevePath 返回全文索引落盘目录(env BLEVE_PATH,默认 ./.data/bleve)。
// 生产部署应指向持久卷;开发期落在工作目录下的 .data。
func blevePath() string {
if p := os.Getenv("BLEVE_PATH"); p != "" {
return p
}
return ".data/bleve"
}
// openBleve 打开(或新建)落盘全文索引;落盘失败退回内存索引兜底。
func openBleve() *bleveStore {
path := blevePath()
idx, err := bleve.Open(path)
if err == bleve.ErrorIndexPathDoesNotExist {
if mkerr := os.MkdirAll(filepath.Dir(path), 0o755); mkerr != nil {
log.Printf("[rag] bleve 目录创建失败,全文路退内存(重启会丢): %v", mkerr)
return memBleve()
}
idx, err = bleve.New(path, bleveMapping())
}
if err != nil {
log.Printf("[rag] bleve 落盘打开失败,全文路退内存(重启会丢): %v", err)
return memBleve()
}
log.Printf("[rag] bleve 全文索引落盘就绪: %s", path)
return &bleveStore{idx: idx}
}
// memBleve 内存兜底:落盘不可用时退回内存索引(功能在、重启会丢),好过全文路全哑。
func memBleve() *bleveStore {
idx, err := bleve.NewMemOnly(bleveMapping())
if err != nil {
log.Printf("[rag] bleve 内存兜底也失败,全文路降级: %v", err)
return &bleveStore{}
}
log.Printf("[rag] bleve 退回内存索引(非持久)")
return &bleveStore{idx: idx}
}
func (b *bleveStore) ready() bool { return b != nil && b.idx != nil }
// close 关闭索引(落盘版释放锁并刷盘);优雅停机时调。
func (b *bleveStore) close() {
if b.ready() {
_ = b.idx.Close()
}
}
// index 把 (kb, doc, texts) 写入全文索引(id 含 kb+doc+文本哈希,幂等)。
func (b *bleveStore) index(kb, doc string, texts []string) error {
if !b.ready() {
return nil
}
batch := b.idx.NewBatch()
for _, t := range texts {
id := fmt.Sprintf("%s:%s:%x", kb, doc, fnvHash(t))
if err := batch.Index(id, map[string]any{"text": t, "kb": kb, "doc": doc}); err != nil {
return err
}
}
return b.idx.Batch(batch)
}
// deleteDoc 删除某 (kb, doc) 的全部全文块(笔记重入库前清旧块)。
func (b *bleveStore) deleteDoc(kb, doc string) {
if !b.ready() || doc == "" {
return
}
kq := bleve.NewTermQuery(kb)
kq.SetField("kb")
dq := bleve.NewTermQuery(doc)
dq.SetField("doc")
req := bleve.NewSearchRequest(bleve.NewConjunctionQuery(kq, dq))
req.Size = 1000
res, err := b.idx.Search(req)
if err != nil {
return
}
batch := b.idx.NewBatch()
for _, h := range res.Hits {
batch.Delete(h.ID)
}
_ = b.idx.Batch(batch)
}
// search 全文检索(可按 kb 过滤),返回 BM25 排序的命中。
func (b *bleveStore) search(kb, q string, topK int) []Hit {
if !b.ready() || q == "" {
return nil
}
mq := bleve.NewMatchQuery(q)
mq.SetField("text")
var qy query.Query = mq
if kb != "" {
tq := bleve.NewTermQuery(kb)
tq.SetField("kb")
qy = bleve.NewConjunctionQuery(mq, tq)
}
req := bleve.NewSearchRequest(qy)
req.Size = topK
req.Fields = []string{"text"}
res, err := b.idx.Search(req)
if err != nil {
return nil
}
var hits []Hit
for _, h := range res.Hits {
text, _ := h.Fields["text"].(string)
if text != "" {
hits = append(hits, Hit{Text: text, Score: float32(h.Score)})
}
}
return hits
}
func fnvHash(s string) uint64 {
h := fnv.New64a()
_, _ = h.Write([]byte(s))
return h.Sum64()
}