79e834e8e9
几十万字文件从"广而浅"到准生产级: - 向量化串行→并发分批保序(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>
160 lines
4.6 KiB
Go
160 lines
4.6 KiB
Go
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
|
||
}
|
||
|
||
// bleveMapping:text 字段用 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()
|
||
}
|