Skip to content

02 入库链路深挖问答 ​

怎么用这篇文档:每题的「口述回答(背诵这段)」是可以原话说出口的版本,先背它;「讲解与备注」用来应对追问,理解即可不用背;「代码依据」是面试官要是让你投屏指代码时你要能马上定位的行号;「别踩的雷」每条都是这题最容易翻车的说法,面试前扫一遍。 全部数字/常量/库名都对齐当前代码,行号可以直接在 IDE 里跳转验证;凡代码里没有的,我写「代码中未找到」,不编造。 建议节奏:面试前 30 分钟只看「一、30 秒讲清链路」+「别踩的雷」+「四、背诵清单」,中间细读。


一、30 秒讲清入库链路 ​

背诵总述:我这个入库链路是「同步做轻活、异步做重活」的两段式。用户上传后,请求线程里只做四件事:卡体积、校验知识库归属和格式能力、对文本类文件跑一次可读性预检拒绝扫描件、把文件落盘并在 PostgreSQL 里插一条 documents 和一条 ingest_tasks,两条都是 pending,然后马上返回 task_id,不阻塞。真正的入库是 worker 池去轮询 ingest_tasks 表:领取任务时用一条 UPDATE ... WHERE id IN (SELECT ... LIMIT 1) RETURNING 把状态原子地改成 processing,然后走 Load→Chunk→Embed→Store:loader 按扩展名选解析器把 PDF/DOCX/Excel/CSV/HTML/Markdown/TXT 统一成 Block 列表;chunker 按配置或文档类型选固定/递归/标题三种策略切成 chunk;embedder 批量调 OpenAI 兼容接口拿向量;最后把 chunk 正文和元数据写进 Qdrant,同时把 chunk_id 数组和状态写回 PostgreSQL,并更新内存 BM25 索引。失败就回 pending 并按 2s/4s/8s 退避重试,超上限置 failed 并把文档同步成 failed。

[浏览器] ──multipart POST /api/v1/documents/upload?kb_id=xxx──▶ router.go:122
   │
   │ ═══════════ 同步段:HTTP 请求线程内(handler_doc.go:38-138)═══════════
   │ ① MaxBytesReader 卡体积(必须早于任何 body 读取)   handler_doc.go:42-43
   │ ② kb_id 必须是合法 UUID + 知识库归属校验,越权 404   handler_doc.go:46-60
   │ ③ registry.Support:格式识别 + 多媒体能力预检        handler_doc.go:74-78
   │ ④ precheckReadable:文本类 strict 解析 + 可读性判定   handler_doc.go:82-90 / 325-338
   │ ⑤ 落盘 ${FileStorageDir}/${kbID}/${docID}${ext}      handler_doc.go:96-105
   │ ⑥ INSERT documents   status=pending                 handler_doc.go:120 ┐无事务
   │ ⑦ INSERT ingest_tasks status=pending                handler_doc.go:133 ┘
   │ ⑧ 立即返回 {task_id, document_id}                    handler_doc.go:138
   │ ═══════════════════════════════════════════════════════════════════════
   ▼ 200 OK(不等解析、不等 embedding)
   ⋮  …… (这里的"队列"就是 PostgreSQL 的 ingest_tasks 表,没有 Redis/Kafka)
   │ ═══════════ 异步段:worker goroutine(worker.go:69-145)═══════════
   ▼
[worker ×N] ──无任务每 500ms 轮询一次,每次领 1 条──▶ worker.go:19-21 / 79
     pending ──(UPDATE … RETURNING,一条语句内"选+占")──▶ processing   task.go:73-96
       │
       ├─ GetDocument + os.Open(doc.FilePath)                       worker.go:107-118
       └─ pipeline.Ingest                                            worker.go:121
            ├─ DeleteByFilter{document_id} 清旧向量 + bm25.RemoveByDoc  pipeline.go:66-73
            ├─ loader.Load(tolerant:坏文件返回空文档+warning)        pipeline.go:76
            ├─ ValidateReadable(min_readable_chars=20)                pipeline.go:89
            ├─ chunker.Chunk(默认 512 token / overlap 50 / 标题层级 2) pipeline.go:94
            ├─ embedder.Embed(批量、限流、单批重试)                    pipeline.go:105
            ├─ vectorstore.Upsert(Qdrant,Wait=true 同步确认)          pipeline.go:141
            └─ bm25.AddWithDocID(内存倒排索引)                        pipeline.go:146-152
       ├─ 成功 ─▶ ingest_tasks=completed;documents=completed + chunk_ids  worker.go:132-142
       └─ 失败 ─▶ ingest_tasks=pending, retry_count++, updated_at=now+2^n s worker.go:149-156
                  超上限 ─▶ ingest_tasks=failed;documents=failed(chunk_ids={})worker.go:157-165

状态流转(两个状态机,取值完全相同:pending / processing / completed / failed,store.go:22-35)

上传                    worker 领取             成功
pending ─────────────▶ processing ────────────▶ completed
   ▲                        │  ▲                     ▲
   │  失败未超上限(退避后)    │  │                     │ 手动重试(POST /tasks/:id/retry)
   └────────────────────────┘  │                     │ 仅 failed 可重试,retry_count 清零
        进程重启 ResetProcessingTasks                  │
   ┌────────────────────────────┘                failed
   └───────────────────────────────────────────────┘

注意:documents.status 全程不会变成 processing(store.go:32 的 DocStatusProcessing 从没被写入过,全仓非测试代码零调用),这是可以主动承认的小瑕疵。


二、深挖问答 ​

Q1. 你这个上传接口,为什么能做到「上传就立刻返回」?同步那一段到底干了哪些活? ​

面试官想考:你是否清楚同步/异步边界划在哪、异步的载体是什么(PG 表当队列还是真 MQ)、有没有把重活放在请求线程里。

口述回答(背诵这段):这个接口就是「同步做轻活、异步做重活」。请求线程里我只做五件事:第一,用 http.MaxBytesReader 把 body 大小卡死,这一步必须放在任何读取之前,因为后面 FormFile 会把整个 multipart body 解析掉,放在后面就来不及了;第二,kb_id 走 query 参数并且必须是合法 UUID,防止拿着 kb_id 做路径穿越,同时校验知识库归属,越权/不存在统一返回 404;第三,用解析器注册表做格式和能力的预检,格式不支持、或者图片音频视频缺少对应的视觉/语音能力配置,直接 400 拒绝;第四,文本类文件我会用 strict 模式完整解析一遍,按 min_readable_chars 默认 20 判断它是不是扫描件或空内容,不合格当场拒绝;第五,文件落到 ${FileStorageDir}/${kbID}/${docID}${ext},然后往 documents 和 ingest_tasks 各插一条 pending 记录,返回 task_id 和 document_id。所以请求耗时只取决于落盘加一次预检解析,跟 embedding 完全无关。真正的 Load→Chunk→Embed→Store 是 worker 池轮询 ingest_tasks 表来消费的,队列就是这张 PG 表,没有额外引入 Redis 或 Kafka。

讲解与备注:这题的核心得分点是「边界清晰 + 解释为什么预检要同步」。预检之所以必须放同步段,是因为扫描件、加密 PDF、越权 kb_id 这类错误如果等到异步才暴露,用户拿到 200 之后永远不知道为什么"没入库",所以我把"能立刻判定的拒绝理由"全部前置。反过来说,我没有同步做分块和向量化,因为那是 O(文档大小) 甚至要调外部 API 的活儿。坑也要主动讲:落盘 + 两条 INSERT 之间没有事务,如果第二条 INSERT 失败,磁盘文件和 documents 记录会变成孤儿(handler_doc.go:133-136)。面试官听到「队列就是 ingest_tasks 表、没有额外 MQ」时会满意,因为这说明你知道取舍:优点是少一个中间件、天然持久化、重启不丢任务;缺点是轮询有延迟、没有背压。 代码位置:入口 internal/api/handler_doc.go:38;体积限制 :42-43;kb 校验 :46-60;格式/能力预检 :74-78;可读性预检 :82-90、:325-338;落盘 :96-105;两条 INSERT :120、:133;返回 :138。

代码依据:

go
// 大小限制必须先于任何 body 读取(FormFile 会解析整个 multipart body)
maxBytes := int64(h.cfg.UploadMaxSizeMB) * 1024 * 1024
c.Request.Body = http.MaxBytesReader(c.Writer, c.Request.Body, maxBytes)

kbID := c.Query("kb_id")
...
// 支持判断(格式识别 + 能力配置):格式不支持 / 能力未配置在此统一拒绝
info := loader.FileInfo{Filename: file.Filename, Size: file.Size}
if res := h.registry.Support(info); !res.Supported {
    Fail(c, CodeBadRequest, res.Reason)
    return
}

internal/api/handler_doc.go:41-78

追问链:

  • 追问:为什么不直接把任务丢给 goroutine,而要写张表? → 答:进程重启 goroutine 就没了,任务会丢;写表才有持久化、可查询、可手动重试,代价是轮询延迟(最长 500ms)和空转查询。
  • 追问:如果用户刚上传完就立刻 GET /documents,他会看到什么状态? → 答:看到 pending;因为 documents.status 从头到尾不会被写成 processing,只有 completed 和 failed 两种终态。
  • 追问:返回之后用户怎么知道入库结果? → 答:轮询 GET /api/v1/tasks/:id,失败任务还能 POST /api/v1/tasks/:id/retry 手动重试(仅 failed 可重试)。

别踩的雷:

  • ❌ 说「上传后就异步返回,其他啥也没做」——同步段有落盘和两条 DB 写入,说漏了会显得不清楚边界。
  • ❌ 说「事务保证两条记录同时成功」——没有事务,grep Begin( 在 store/pipeline/task 零命中,说错了会被当场追问。
  • ❌ 说用了 Redis/Kafka/NSQ——代码里没有,队列就是 ingest_tasks 表。

Q2. 文件落在哪、什么时候落、事务边界在哪?失败了会留下什么脏数据? ​

面试官想考:有没有"资源与记录一致性"的意识,会不会出现磁盘文件孤儿、DB 记录孤儿。

口述回答(背诵这段):落盘位置是 ${server.file_storage_dir}/${kbID}/${docID}${ext},默认目录是 ./data/uploads,docID 是 UUID v4,存的时候用 docID 而不是原文件名,一是防重名、二是防路径穿越。落盘发生在写数据库之前,用的是 Gin 的 SaveUploadedFile。事务边界这块我要老实说:整条链路没有任何数据库事务,CreateDocument 和 CreateTask 是两次独立的 Exec,所以理论上存在两个孤儿窗口——第一条 INSERT 成功、第二条失败,会留下「文件 + documents 记录但没有任务」,这个文档永远不会被 worker 消费;还有落盘成功但第一条 INSERT 失败,会留下纯磁盘孤儿文件。我的取舍是:这两个窗口都是本地短事务、失败率极低,我选择了「用日志和文件覆盖来兜底」而不是引入分布式事务;如果要做严格版本,最小改法是把两条 INSERT 包进一个 pgx 事务(一次 Begin/Commit),再拿 documents.task_id 做外键约束,然后在启动时做一次扫描把没有任务的 pending 文档补成任务。

讲解与备注:这题真正想验证的是「你知不知道自己的系统在哪一刻会不一致」。回答得分点是三点:位置可预测、命名用 UUID、主动承认无事务 + 给出最小改进方案。注意区分两层"不一致":文件与 DB 记录的不一致(这里),以及后面 Q13 的 PG 与 Qdrant 不一致(外部存储层),别混着讲。另外 documents.file_path 存的是本地绝对路径(handler_doc.go:101-102),这意味着多实例部署必须共享存储卷,否则 A 实例落盘的文件 B 实例读不到——worker 是按 doc.FilePath 直接 os.Open(worker.go:113),没有任何对象存储抽象。这点面试官很可能会追。 代码位置:internal/api/handler_doc.go:92-105(docID/落盘)、:120-136(两条 INSERT)、internal/store/document.go:13-27、internal/store/task.go:10-20。

代码依据:

go
docID := uuid.New().String()
ext := filepath.Ext(file.Filename)

// 保存文件
dir := filepath.Join(h.cfg.FileStorageDir, kbID)
if err := os.MkdirAll(dir, 0o755); err != nil {
    Fail(c, CodeInternal, "创建存储目录失败")
    return
}
filePath := filepath.Join(dir, docID+ext)
if err := c.SaveUploadedFile(file, filePath); err != nil {
    Fail(c, CodeInternal, "保存文件失败")
    return
}

internal/api/handler_doc.go:92-105

追问链:

  • 追问:多实例部署会怎样? → 答:file_path 是本地绝对路径,worker 直接 os.Open,所以必须共享存储卷(NFS/云盘);代码里没有对象存储抽象,这是已知限制。
  • 追问:删文档时文件会删吗? → 答:会,deleteDocument 里 os.Remove(doc.FilePath),但错误被 _ = 忽略了(handler_doc.go:319)。
  • 追问:怎么补这些孤儿? → 答:加一个启动巡检,扫 pending 且无任务的文档补建任务,或者直接用一个事务把两条 INSERT 合并,再给 task_id 加外键。

别踩的雷:

  • ❌ 说「事务保证一致」——没有事务,会被反问「你怎么保证的」。
  • ❌ 说「文件存到对象存储/S3」——代码中是本地磁盘 os.MkdirAll + SaveUploadedFile。
  • ❌ 说「原文件名做存储名」——存的是 docID+ext,原文件名只进 documents.filename 字段。

Q3. 你这个多格式加载器是怎么做的?PDF、DOCX、Excel、HTML 分别用什么库解析,为什么? ​

面试官想考:是不是真读过代码,还是抄了个 loader 就完事;能不能说清选型理由和统一抽象。

口述回答(背诵这段):加载器是「注册表 + Parser 接口」的结构:每个格式实现 Parse / SupportedExts / SupportedMIMEs 三个方法,Register 的时候把扩展名和 MIME 都登记进两张 map,Resolve 先按扩展名精确命中、没有再按 MIME 兜底,都没有就抛 ErrUnsupportedFormat。具体实现:PDF 用 pdfcpu,先 api.PageCount 拿页数,再对每一页调 api.ExtractContent 抽内容流,一页产出一个带 page 元数据的段落 Block;DOCX 用 fumiama/go-docx,遍历 Body 里的 Paragraph,看 Style.Val 是不是 Heading1 到 Heading6 来决定层级;Excel 用 excelize/v2,每个 sheet 一个一级标题 Block,每行用 tab 拼成一个表格 Block;CSV 用标准库 encoding/csv,开了 LazyQuotes 和 FieldsPerRecord=-1 容错,首行当表头;HTML 用 golang.org/x/net/html 递归遍历,h1-h6/p/li/pre 分别映射成标题、段落、列表项、代码块,并且跳过 script/style/nav/footer/head 这些噪音标签;Markdown 用 goldmark 走 AST,Heading、Paragraph、ListItem、FencedCodeBlock 四类;TXT 就是 bufio.Scanner 按空行分段。选型原则我当时想的是:能用官方/主流库就不用外部命令行,避免部署时还要装二进制;但视频必须用 ffmpeg,因为要抽帧和拆音轨,那是另一条线。

讲解与备注:这题最容易被追问"为什么不用 X"。诚实的说法是:代码里没有写选型理由(我确实没在注释里留下选型记录),所以要讲成"我当时基于这几点权衡",而不是"因为代码注释说"。可以讲的取舍:① PDF 抽文本用 pdfcpu 而不是 poppler/mutool,是因为 pdfcpu 是纯 Go 库,不引入 CGO 和外部二进制,go build 出来就是单文件;代价是它专注于 PDF 结构与内容流,没有版面分析,所以 PDF 我拿到的是一页一坨文本,没有做行/列还原。② DOCX 用 go-docx 而不是先转 PDF 或调 LibreOffice,也是为了纯 Go;代价是它只认标准 Heading 样式,用自定义样式的标题我识别不出来。③ Excel 用 excelize 是 Go 生态里最成熟的选择,代价是 GetRows 是全量读进内存。④ HTML 用 x/net/html 而不是 goquery,是因为我只需要遍历节点、不需要 CSS 选择器。⑤ Markdown 用 goldmark 是因为它有完整 AST、还兼容 CommonMark。得分句是那句"能用纯 Go 库就不引入外部二进制,只有视频例外(ffmpeg)",它体现了一致的工程判断,而不是每个格式随便挑一个。 代码位置:注册表 internal/loader/registry.go:21-49;注册清单 internal/loader/loader.go:21-35;各 parser 见下表。

格式库/实现关键行
PDFgithub.com/pdfcpu/pdfcpu(api.PageCount / api.ExtractContent)internal/loader/parser_pdf.go:44,61
DOCXgithub.com/fumiama/go-docxinternal/loader/parser_docx.go:38
Excelgithub.com/xuri/excelize/v2internal/loader/parser_excel.go:29
CSV标准库 encoding/csv(LazyQuotes、FieldsPerRecord=-1)internal/loader/parser_csv.go:25-27
HTMLgolang.org/x/net/htmlinternal/loader/parser_html.go:26
Markdowngithub.com/yuin/goldmark(AST)internal/loader/parser_markdown.go:39-45
TXT标准库 bufio.Scannerinternal/loader/parser_txt.go:28-52
图片/音频抽象接口 VisionProvider / SpeechProvider(OpenAI 兼容 / dashscope)internal/loader/capability.go:13-37
视频ffmpeg / ffprobe(抽帧、拆音轨)+ VLM + ASRinternal/multimedia/frame_extractor.go:45、audio_extractor.go:41

代码依据:

go
func (r *defaultRegistry) Resolve(info FileInfo) (Parser, error) {
    ext := strings.ToLower(filepath.Ext(info.Filename))
    if ext != "" {
        if p, ok := r.extMap[ext]; ok {
            return p, nil
        }
    }
    if info.MIMEType != "" {
        mime := strings.ToLower(info.MIMEType)
        if p, ok := r.mimeMap[mime]; ok {
            return p, nil
        }
    }
    return nil, &ErrUnsupportedFormat{Filename: info.Filename, MIMEType: info.MIMEType}
}

internal/loader/registry.go:30-49

追问链:

  • 追问:PDF 里表格怎么办? → 答:处理不了。pdfcpu 给的是内容流文本,我按页整块存,没有做版面/表格还原,这是明确的能力边界。
  • 追问:加密 PDF 怎么处理? → 答:代码里没有显式加密检测,我把校验模式放宽到 ValidationRelaxed,加密文件会在取页数或抽内容时报错,tolerant 模式下变成"空文档 + warning",最后被可读性判定拒掉,用户看到的是"文档无可读文本"而不是"文件已加密"——这是措辞上的缺陷。
  • 追问:加一个格式要改几处? → 答:写一个实现三方法的 parser,Register 一下就完了,pipeline 和 API 都不用改;app.go 里用真实 provider 覆盖同名扩展名也是靠这个 map 后写覆盖先写。

别踩的雷:

  • ❌ 说 PDF 用 unidoc/pdftotext/go-fitz——都不是,是 pdfcpu。
  • ❌ 说"支持扫描件 OCR"——没有 OCR,图片走的是视觉模型描述,PDF 扫描件是被拒绝的。
  • ❌ 说 MIME 优先——是扩展名优先、MIME 兜底,而且入库链路里 worker 只传了 Filename,MIME 兜底实际只有 API 预检路径能用。

Q4. 解析完统一成什么结构?后面分块和检索怎么用这个结构? ​

面试官想考:抽象设计能力——是不是每加一个格式就要改下游。

口述回答(背诵这段):我把所有格式统一成一个 Document,里面有 Blocks []Block 和文档级 Metadata,Block 是四个字段:Type、Content、Level、Metadata map[string]any。Type 目前有七种:段落、标题、列表项、代码块、表格、图片描述、音频分段;Level 只对标题有效,1 到 6。文档级元数据里有文件名、格式、标题、字节大小、PDF 页数和一个 Extra 扩展 map,多媒体会往里塞时长、编码、宽高、来源文件名这些。这个结构的好处是下游只认 Block:分块器根据 Type 决定怎么渲染文本——标题加 #、列表项加 - 、代码块加三反引号;根据 Metadata["page"] 判断是不是 PDF 从而按页分块;根据 BlockType 是不是图片描述或音频分段判断走多媒体按块切分;检索侧要用到的 start_ms/end_ms/source_type/page_number 全部是从 Block.Metadata 一路透传到 Qdrant payload 的。另外还有 Warnings []string,用来承载非阻断告警,比如视频没配语音能力就跳过音轨,这个 warning 会一路传到任务的 warning_message 字段,让用户看到"入库成功了但音轨没转写"。

讲解与备注:这题的关键词是「Block 是稳定契约」。面试官想确认你的抽象不是"为了分层而分层"。可以强调两点取舍:① 元数据用 map[string]any 而不是强类型 struct,好处是各格式自由扩展(PDF 塞 page、媒体塞时间戳),坏处是丢了类型安全,取的时候要类型断言,我在分块器里 b.Metadata["page"].(int) 就吃过默认值的亏(断言失败会变成 0,我在 chunker.go:185-188 兜底成第 1 页)。② 我用一个 BlockType 枚举而不是"每个格式一个 struct",是为了让 chunker 只依赖一个类型,代价是格式特有信息只能塞 Metadata。得分句:"加一种格式只需要实现 Parser 三方法,chunker / pipeline / 检索一行都不用改"。 代码位置:internal/loader/types.go:11-70(BlockType/Block/Document/DocumentMeta/LoadResult)、internal/chunker/strategy_heading.go:95-107(按 Type 渲染)、internal/chunker/chunker.go:49-56,165-172(按 Metadata 选分块路径)、internal/pipeline/pipeline.go:120-138(元数据透传成 payload)。

代码依据:

go
const (
    BlockParagraph        BlockType = iota // 普通段落
    BlockHeading                           // 标题
    BlockListItem                          // 列表项
    BlockCode                              // 代码块
    BlockTable                             // 表格
    BlockImageDescription                  // 图片/视频帧视觉描述(多媒体)
    BlockAudioSegment                      // 音频转写分段(多媒体,含起止时间戳 metadata)
)

type Block struct {
    Type     BlockType
    Content  string
    Level    int            // 仅标题有效,1-6
    Metadata map[string]any // 扩展字段
}

internal/loader/types.go:11-27

追问链:

  • 追问:map[string]any 有什么问题? → 答:类型不安全,取值要断言,断言失败静默变零值,我在 PDF 页码那里就加了兜底;更好的做法是定义强类型 PageMeta/MediaMeta 再嵌进 map。
  • 追问:Warnings 有什么用? → 答:承载非阻断降级,比如视频音轨没转写;它一路传到 ingest_tasks.warning_message(worker.go:134-136),任务依然算 completed。
  • 追问:Table 类型有没有做特殊处理? → 答:渲染阶段走了默认分支按原样输出,没有表头补全、也没有跨 chunk 保护,这是不足。

别踩的雷:

  • ❌ 说"每个格式返回自己的 struct,下游做类型转换"——统一成一个 Document/Block。
  • ❌ 说"PDF 表格会识别成 BlockTable"——不会,PDF 页内所有内容是一个 BlockParagraph,只有 CSV/Excel 才产 BlockTable。
  • ❌ 说 Warning 会导致任务失败——不会,是非阻断的。

Q5. 用户传上来一个扫描件 PDF,你怎么识别并拒绝?min_readable_chars 默认是多少? ​

面试官想考:有没有防脏数据的意识,以及有没有真的想过"怎么判定一页 PDF 是图不是字"。

口述回答(背诵这段):min_readable_chars 代码里的默认是 20,配置项在 loader.min_readable_chars。判定不是简单数字符,而是先算「可读文本量」:一个汉字算 1,一个纯字母单词(长度 2 到 20)算 1,标点和空白不算;关键是我维护了一张 52 个 PDF 内容流操作符的白名单,像 q/Q/cm/re/BT/ET/Tf/TJ/Tj/Do/Im 这些,出现它们就不计入可读量。因为扫描件 PDF 抽出来往往是 q 595.44 0 0 841.68 cm 1 g /Im10 Do Q 这种图像绘制指令,如果不排除,靠数字符会被误判成"有内容"。判定条件是双阈值:全文可读量要 ≥ 20,并且至少有一个 block 的可读量 ≥ 20,第二条是为了防"大量 5 到 15 字的指令碎片累加过线",我专门写了一个测试锁定这个行为。防线有两层:上传时用 strict 模式同步预检一次,不合格就直接 400,用户马上收到明确错误;入库时 pipeline 还会用 tolerant 结果再兜一次,因为文件可能在两次读取之间被替换。多媒体是豁免的,因为音频转写可能就是"大家好"这种短句,按单块 20 字判会被误杀。

讲解与备注:这题的杀手锏是那张 52 项操作符白名单和双阈值,说出这两点面试官会觉得你真的处理过 PDF 脏数据。技术原理讲清:PDF 扫描件本身是"一张图片 + 图像绘制指令",ExtractContent 拿到的就是指令流,所以要么上 OCR,要么就当无文本拒掉——我选了拒掉,因为项目没有 OCR 能力,硬塞进去只会污染知识库。取舍讲法:"宁缺毋滥:拒掉一个扫描件,用户重传一份文字版 PDF 就好;放进去一个乱码文档,会长期污染检索结果且很难发现。" 阈值 20 的依据代码注释给了:正常文档动辄数百,扫描件的指令碎片通常小于 15(config.go:56-59)。另外要主动说这套启发式的弱点:纯英文短文档(比如只有 hello world,2 个词)会被误拒;如果要做严,可以按"可读字符数 / 总字符数"的密度来判断,而不是绝对量。 代码位置:阈值默认 internal/config/config.go:521-523;YAML 配置 configs/config.yaml:37;算法 internal/loader/validate.go:26-53;白名单 :9-18;双条件 :79-94;多媒体豁免 :74-77;上传预检 internal/api/handler_doc.go:82-90,325-338;入库兜底 internal/pipeline/pipeline.go:82-91。

代码依据:

go
total := 0
maxBlock := 0
for _, b := range doc.Blocks {
    n := ReadableCharCount(b.Content)
    total += n
    if n > maxBlock {
        maxBlock = n
    }
}
if total < minChars || maxBlock < minChars {
    return &ErrNoReadableContent{
        Format:   doc.Metadata.Format,
        Readable: maxBlock, // 展示代表性 block 的可读量
        MinChars: minChars,
    }
}

internal/loader/validate.go:79-94

追问链:

  • 追问:为什么要有"单 block 也要达标"这条? → 答:防碎片累加。测试 TestValidateReadableFragments 就是四个 PDF 指令碎片,总量可能过线但没有一个 block 达标,必须拒。
  • 追问:如果用户传一份只有 hello world 的 txt? → 答:会被拒(可读量 2 < 20),这是启发式的代价;要更准可以引入 OCR 或按密度判定。
  • 追问:加密 PDF 走哪条路? → 答:代码里没有加密检测,它会在解析阶段失败变成空文档,最后也是被这里拒掉,报的是"无可读文本",措辞不够准确。

别踩的雷:

  • ❌ 说"用 OCR 识别扫描件"——项目没有 OCR,扫描件是拒绝的。
  • ❌ 说阈值是 10 或 50——是 20(代码默认与 YAML 一致)。
  • ❌ 说"只看总量"——是双条件,单 block 也要达标,这是最容易被追问的点。

Q6. 你写了三种分块策略,它们的算法差别是什么?一个文档进来你走哪一种? ​

面试官想考:分块是 RAG 效果的核心,这里能看出你是抄参数还是理解语义边界。

口述回答(背诵这段):三种策略:第一种固定大小,先算全文字数够不够 chunk_size,不够就直接一整块;够的话按字符窗口切,先估一个结束位置、再收缩到 token 限额内、再往后扩到刚好不超限,然后回退到最近的空白或者中文标点,块之间按 token 数回退做 overlap。第二种递归字符,我的分隔符列表是十个:双换行、单换行、中文的句号感叹号问号、英文的点号感叹号问号、空格、空串兜底;先把文本按第一个分隔符切,超限的段落用剩下的分隔符递归切,再用贪心合并把小段拼到接近上限,最后给每块前缀拼上一块的尾部做 overlap。第三种 Markdown 标题,按 heading_level(默认 2)把 block 流切成节,标题层级用栈维护,产出的 chunk 里正文本身含标题原文,另外把 主标题 > 二级标题 这样的路径写进 heading_context 元数据、最近一级标题写进 heading,再 slug 化成 anchor 给前端跳转用;如果某一节超过 chunk_size,降级用递归策略切,子块继承同一个标题上下文。选哪一种是代码决定优先级、配置只能影响第三顺位:只要文档里有图片描述或音频分段的 block,就走多媒体按块切分,一块一个 chunk;只要有 block 带 page 元数据(也就是 PDF),就按页分块,同一页合并成一个 chunk、超长页才降级递归;这两个都不满足,才看配置里的 strategy,heading 走标题策略,fixed/recursive 走对应策略,未知值回退递归。

讲解与备注:这题有两个高分点。第一是优先级明确的自动路由:媒体 > PDF 分页 > 配置策略,说明你不是"配了啥就用啥",而是认为"切块边界应该由文档语义决定"。第二是能讲出每种策略的适用场景:固定大小适合无结构的长文本、最省事但对语义不敏感;递归字符适合中英文混排的通用文档,靠标点保住句子完整;标题策略适合结构化文档(Markdown、技术文档),能天然实现"按章节检索 + 引用跳转"。要主动说出代价:递归策略我用了 strings.Split 切分但没有把分隔符拼回去,所以超过 chunk_size 的文档会丢掉句末标点(strategy_recursive.go:52-64),这是一个真实缺陷,修法是切开后把分隔符附加到前一段末尾;另外非标题策略下的 heading_context 是错的,我那个函数叫 extractFirstHeadingContext,但实现是遍历完所有 block 后取标题栈的最终态,所以拿到的是文档最后一个标题路径,还被赋给了每个 chunk(chunker.go:245-259),准确说法应该是按 chunk 位置回溯当时的栈。主动讲这两个 bug 比藏着强,说明你知道边界在哪。 代码位置:入口与路由 internal/chunker/chunker.go:41-105;固定 internal/chunker/strategy_fixed.go:14-121;递归 internal/chunker/strategy_recursive.go:13-161;标题 internal/chunker/strategy_heading.go:28-107;媒体/分页 internal/chunker/chunker.go:136-222;配置解析 internal/chunker/types.go:15-34。

代码依据:

go
// 多媒体 Document:按块切分(每 Image/Audio Block 一个 chunk,时间戳 1:1)
if isMediaDoc(doc.Blocks) {
    return c.chunkMediaBlocks(doc)
}
// PDF Document:按页分块(block 带 page metadata)
if isPagedDoc(doc.Blocks) {
    return c.chunkPagedBlocks(doc, config)
}

var rawChunks []rawChunk
if config.Strategy == StrategyHeading {
    sections := c.headingStrategy.SplitByBlocks(doc.Blocks, config, c.tokenizer)
    ...
} else {
    strategy, ok := c.strategies[config.Strategy]
    if !ok {
        strategy = c.strategies[StrategyRecursive] // 未知策略静默回退递归
    }

internal/chunker/chunker.go:48-74

追问链:

  • 追问:PDF 为什么不按标题切? → 答:pdfcpu 只给内容流文本,没有标题语义,PDF 的 Block 全是段落加页码,所以我只能按页切。
  • 追问:标题策略里 breadcrumb 会拼进 chunk 正文吗? → 答:不拼。正文里含标题原文(渲染成 ## 标题),A > B 这种路径只进 heading_context 元数据,放在 payload 里给检索和展示用。
  • 追问:一个视频的转写很长,会切吗? → 答:不会,媒体走的是"一块一个 chunk",完全不看 chunk_size,这是为了保住时间戳 1:1;代价是长转写会变成超大 chunk。

别踩的雷:

  • ❌ 说"配置了 fixed 就是固定大小"——媒体和 PDF 会无视配置,优先级是媒体 > PDF 分页 > 配置。
  • ❌ 说"标题路径拼进了正文"——正文只有标题原文,路径在 metadata。
  • ❌ 说"分隔符切开后保留了标点"——递归策略是丢的,主动承认更好。

Q7. chunk_size、chunk_overlap 真实默认值是多少?token 是怎么算的,你们用 tiktoken 吗? ​

面试官想考:知不知道自己的常量、以及是否理解"token 估算与模型真实 tokenizer 的差异"。

口述回答(背诵这段):代码里的默认值:chunk_size 512、chunk_overlap 50、标题目标层级 heading_level 2,配置在 chunker 段;生产 YAML 里实际配的是 chunk_size: 500、overlap: 50、heading_level: 2,策略是 recursive。有个细节要说清楚:overlap 的填充判据写的是 < 0 才填默认值,所以如果 YAML 里把这一行删掉,得到的是 0、也就是完全不重叠,跟"默认 50"这个注释是有出入的,生产配置里显式写了 50 所以没问题,但这是个隐患。token 这块我要老实讲:我没有用真正的 tokenizer,项目里是一个启发式的估算器,规则是——一个汉字算 2 个 token、一个标点或符号算 1、一段连续的非空白字符(比如一个英文单词)算 1、空白只做分隔不计分。所以 chunk_size=500 对纯中文大约是 250 个汉字,因为每个汉字按 2 计。为什么汉字要算 2?因为主流 tokenizer 里一个汉字通常要占 1 到 2 个 BPE token,我用 2 是取保守上界,宁可切小也不要超模型上下文。这个设计的代价是:chunk_size 的语义和模型上下文窗口的 token 不是同一套计量,如果需要精确对齐(比如要卡到 8192 上下文),应该换成真实的 tokenizer 库或者直接调 embedding 服务的分词接口;我是靠"保守估计 + chunk 只有几百 token"来兜住风险的。

讲解与备注:这题的正确答案是大方承认不是真 tokenizer,然后讲清估算规则和保守策略。面试官最反感的是"我们用了 tiktoken"这种一查就穿的说法。加分项是主动指出 ChunkOverlap < 0 这个判据问题,它说明你真读了自己的默认值代码。另外可以补充:chunker 的 Tokenizer 是接口(tokenizer.go:9-11),可以注入,测试里就用了一个"每字符 1 token"的 mock 来精确断言切块边界,所以换成真 tokenizer 只需要实现一个 Count 方法,改动面很小——这是我留的扩展点。 代码位置:默认值 internal/chunker/types.go:37-48(512/50/2)、internal/config/config.go:399-407(同一套判据);YAML configs/config.yaml:86,88,90,92(recursive / 500 / 50 / 2);估算器 internal/chunker/tokenizer.go:20-57;接口 :9-11;测试断言 internal/chunker/chunker_test.go:27-48,253-275。

代码依据:

go
func (t *DefaultTokenizer) Count(text string) int {
    ...
    for _, r := range text {
        if unicode.Is(unicode.Han, r) {
            // 中文字符计 2 token
            if wordBuf.Len() > 0 { count++; wordBuf.Reset() }
            count += 2
        } else if unicode.IsPunct(r) || unicode.IsSymbol(r) {
            if wordBuf.Len() > 0 { count++; wordBuf.Reset() }
            count++
        } else if unicode.IsSpace(r) {
            if wordBuf.Len() > 0 { count++; wordBuf.Reset() }
        } else {
            wordBuf.WriteRune(r)
        }
    }

internal/chunker/tokenizer.go:20-50

追问链:

  • 追问:那 500 token 大概是多少字? → 答:纯中文约 250 字;中英混排取决于词的比例,英文一个词算 1 token,所以同样的 chunk_size 在英文文档里能装更多字符。
  • 追问:为什么不直接上 tiktoken? → 答:Go 生态里对齐 OpenAI/国产模型 BPE 的库不统一,而且不同 embedding 模型 tokenizer 不同;我留了 Tokenizer 接口,需要精确对齐时实现一个 Count 就行。
  • 追问:overlap 有什么用、会不会重复检索? → 答:防止一句话被切断导致语义丢失;代价是相邻 chunk 有重复内容,召回时可能返回多条高度相似的片段,靠后续 RRF 融合和 topK 截断控制。

别踩的雷:

  • ❌ 说"用了 tiktoken / HuggingFace tokenizer"——没有。
  • ❌ 说 chunk_size 默认 500——代码默认是 512,500 是 YAML 配的值,两个都要说清。
  • ❌ 说 overlap 默认一定是 50——判据是 < 0,省略配置就是 0。

Q8. chunk_id 是怎么生成的?我把同一份文档传两遍会怎样,你怎么保证幂等? ​

面试官想考:幂等性意识,这是 RAG 系统最容易被追问的工程点。

口述回答(背诵这段):chunk_id 是 UUID v4,uuid.New().String(),不是内容哈希——这点我不绕:它随机、无内容指纹,所以不存在天然的幂等键。它同时承担四个角色:Qdrant 的 point id、payload 里的 chunk_id、PG 里 documents.chunk_ids 数组的元素、还有内存 BM25 索引的 docID。幂等我是靠"document_id 维度的删后写"来实现的:pipeline.Ingest 一开始就按 document_id 调 Qdrant 的 DeleteByFilter 把这篇文档的旧向量全删掉,同时 bm25.RemoveByDoc 清掉内存索引里这篇文档的所有 chunk,然后再重新写一遍,所以同一篇文档重试、或者用户手动重试失败任务,都是覆盖式幂等,不会产生重复向量。但要区分两种"重复":同一个 document_id 重跑是覆盖;同一份文件第二次上传是新的 document_id,那就真的是两份独立内容,向量库里会有两份、chunk_id 也完全不同,我在上传时没有做内容去重、也没有 (kb_id, filename) 的唯一约束。还有一个坑要说:那个前置删除失败时我只打了 warn 就继续入库(pipeline.go:68),这时旧向量残留加上新向量,检索会出现重复结果。

讲解与备注:诚实 + 有方案是这题的得分姿势。可以给出改进路线:① 把 chunk_id 换成 sha256(document_id + chunk_index + content) 的确定性 ID,天然幂等、便于增量更新;② 上传时算文件内容 hash,同 KB 内 hash 命中就拒绝或提示"已存在";③ 加 (kb_id, content_hash) 唯一索引;④ 前置删除失败应该中断任务而不是继续(否则会写出重复数据)。另外要能解释"为什么当初用 UUID":Qdrant point id 和 PG 数组都需要全局唯一、可跨文档引用,UUID 是最省事的选择,GET /api/v1/chunks/:id 还会用 uuid.Parse 校验入参(handler_chunk.go:28-31),说明这个设计是贯通的;代价就是放弃了内容幂等。 代码位置:生成与 payload internal/pipeline/pipeline.go:116-138;前置删除 :66-73;Qdrant DeleteByFilter internal/vectorstore/qdrant.go:168-184;BM25 RemoveByDoc internal/retriever/bm25.go:104-112;上传造 ID internal/api/handler_doc.go:92;chunk 查询校验 internal/api/handler_chunk.go:28-31。

代码依据:

go
// 重试补偿:清理同一 document_id 的旧 chunk(向量 + BM25),防止重试产生孤儿向量
if req.DocumentID != "" {
    if err := p.vectorstore.DeleteByFilter(ctx, map[string]any{"document_id": req.DocumentID}); err != nil {
        slog.Warn("清理旧 chunk 向量失败(继续入库)", "doc", req.DocumentID, "err", err)
    }
    if p.bm25Index != nil {
        p.bm25Index.RemoveByDoc(req.DocumentID)
    }
}
...
for i, c := range chunks {
    chunkID := uuid.New().String()
    chunkIDs[i] = chunkID

internal/pipeline/pipeline.go:65-73,117-119

追问链:

  • 追问:那你的幂等键到底是什么? → 答:严格说是 document_id,不是 chunk 内容;同 document_id 是覆盖写,同内容不同 document_id 就不幂等。
  • 追问:怎么改成内容幂等? → 答:chunk_id 改成 sha256(document_id + index + content),Qdrant 用同名 point id 做 upsert 覆盖,PG 加 (kb_id, content_hash) 唯一索引,上传时先查 hash。
  • 追问:为什么前置删除失败只告警? → 答:我当时想的是"入库优先、别因为清理失败把用户文档卡死",但这确实是错的,应该让它失败重试,否则会产生重复向量。

别踩的雷:

  • ❌ 说"chunk_id 是内容的 md5/sha256"——是 UUID v4,全仓 sha256 只用于 API Key。
  • ❌ 说"重复上传不会重复入库"——会,新 document_id 就是新文档。
  • ❌ 说"删除失败重试了"——只 slog.Warn 后继续,没有重试。

Q9. Worker 池是怎么起的、并发多少、任务怎么领?为什么这么领不会重复? ​

面试官想考:并发安全 + SQL 细节,这题几乎是必问。

口述回答(背诵这段):worker 数量代码默认是 2,生产 YAML 配的是 5,每个 worker 一个 goroutine,启动时会先做一次悬挂任务重置,然后进入循环。领取方式是数据库轮询,不是 channel、也不是 Redis:每个 worker 每轮调 ClaimPendingTasks(ctx, 1),里面的 SQL 是一条 UPDATE ingest_tasks SET status='processing', updated_at=now() WHERE id IN (SELECT id FROM ingest_tasks WHERE status='pending' AND updated_at <= NOW() ORDER BY created_at LIMIT $1) RETURNING ...。为什么这样不会重复领?因为"选"和"占"在同一条 SQL 里完成,PostgreSQL 的行级锁保证同一行不会被两个事务同时更新到 processing,第二个 worker 的 UPDATE 会等第一个提交后重新求值,发现状态已经不是 pending 就不再返回它——用 RETURNING 拿到行的 worker 才是真正的主人。每次只领 1 条是为了让任务在 worker 之间均匀分摊,不会出现一个 worker 抱一堆、其他空转。没任务的时候每 500ms 轮询一次,领取报错退避 1s。我要承认两点:第一,我没有写 FOR UPDATE SKIP LOCKED,所以并发下可能白跑几次(拿到 0 行)而不是立刻跳过去领下一批;第二,5 个 worker 空转时大概每秒 10 次空查询,量不大但没有做基于 LISTEN/NOTIFY 或 channel 的唤醒优化。单测里我是用 pgxmock 对这个 SQL 做断言的,也在 worker_test.go 里用内存 fakeStore 覆盖了成功、重试、超限、并发 20 个任务这些路径。

讲解与备注:这题要讲出"原子性来自 UPDATE,而不是 SELECT"这个认知——很多人会说"我用 SELECT ... FOR UPDATE 选出来的",但真正安全的做法确实是让写语句自己带 where 条件去抢。可以补两点:① ORDER BY created_at 是 FIFO,但 created_at 相同(同一毫秒)时顺序不确定;② 领取条件是 updated_at <= NOW(),这个条件同时被用来做退避(见 Q10),一箭双雕但要付出语义混淆的代价。面试官听到"单条 UPDATE...RETURNING 原子领取 + 为什么没有 SKIP LOCKED"会满意,因为这说明你知道方案的天花板在哪。测试细节也要能说:TestClaimPendingTasks 用 pgxmock 期望一条 UPDATE ingest_tasks 且返回状态是 processing(store_test.go:116-146);worker 侧 TestProcess_Concurrent 用 20 个 goroutine 跑 process 验证无竞争(worker_test.go:260-285)。 代码位置:常量 internal/task/worker.go:18-22;启动 :44-57;循环 :69-103;领取 SQL internal/store/task.go:73-96;配置默认 internal/config/config.go:534-536、YAML configs/config.yaml:244;测试 internal/store/store_test.go:116-146、internal/task/worker_test.go:260-285。

代码依据:

go
rows, err := s.pool.Query(ctx,
    `UPDATE ingest_tasks SET status = 'processing', updated_at = now()
     WHERE id IN (
         SELECT id FROM ingest_tasks WHERE status = 'pending' AND updated_at <= NOW() ORDER BY created_at LIMIT $1
     )
     RETURNING id, kb_id, document_id, status, retry_count, error_message, warning_message, created_at, updated_at`,
    limit)

internal/store/task.go:74-80

追问链:

  • 追问:为什么不用 channel 做队列? → 答:channel 是进程内的,重启就丢;用表做队列天然持久化、可查询、可手动重试,代价是轮询延迟和空转查询。
  • 追问:为什么没有 SKIP LOCKED? → 答:单条 UPDATE 已经保证了不重复领,SKIP LOCKED 解决的是"减少锁等待、提高吞吐",是性能优化不是正确性必需;如果 worker 数很多、任务很密,我会补上。
  • 追问:领取批量是 1,会不会太慢? → 答:claimBatchSize=1 是为了任务公平分摊;吞吐不够时先加 worker,或者把批量调大并配合 SKIP LOCKED。

别踩的雷:

  • ❌ 说"worker 默认 5 个"——代码默认 2,YAML 配的是 5,两个数都要说。
  • ❌ 说"用 SELECT ... FOR UPDATE 先选再更新"——代码是单条 UPDATE ... WHERE id IN (SELECT ...) RETURNING,没有 FOR UPDATE,也没有 SKIP LOCKED。
  • ❌ 说"用 Redis 做分布式锁防重复领取"——没有,靠 PG 行锁。

Q10. 失败怎么重试?指数退避具体怎么实现的?退避期间任务怎么被跳过? ​

面试官想考:是否真的写过重试,还是只会说"我们会重试三次"。

口述回答(背诵这段):失败时我调 fail,逻辑是:如果 retry_count 还没到 task_max_retries(代码默认 3,生产 YAML 也是 3),就把任务状态改回 pending、retry_count 加一、把错误写进 error_message,然后算一个退避时间:time.Second << retry_count,因为我是先自增再移位,所以第一次失败退避 2 秒,然后是 4 秒、8 秒。退避怎么落地?我把这个"未来时间"写进 updated_at 字段,而领取 SQL 的条件是 status='pending' AND updated_at <= NOW(),所以在退避窗口内这条任务虽然状态是 pending,但不会被领走,时间到了自然被捡起来——不需要额外的调度器或者延迟队列,纯数据库实现。失败超过上限时,任务置 failed 并写错误信息,同时把文档状态也同步成 failed,这样文档列表能反映真实结果,不会一直挂在 pending 上。另外还有一层:embedding 客户端自己也有重试,只对网络错误、429 和 5xx 重试,退避是 1 秒、2 秒、4 秒,最多 3 次;两层叠加的话,最坏情况单个任务会打十几次外部调用。手动重试接口是 POST /tasks/:id/retry,只允许 failed 状态,会把 retry_count 清零、error_message 清空,等于重开一轮。

讲解与备注:得分点是"退避寄生在 updated_at 上"这个实现细节,以及能说清它的代价:updated_at 同时承担了"最后更新时间"和"下次可领取时间"两个语义,任何基于它做监控、排序、归档的逻辑都会被这个未来时间污染;更规范的做法是加一个 next_attempt_at 列(或 available_at),updated_at 保持真实语义。另外可以讲清楚"为什么用指数退避":外部依赖(embedding API、Qdrant)的故障通常有恢复期,固定间隔重试会把压力持续打上去,指数退避给下游恢复窗口。还要主动说清重试的分类问题:我现在是"所有错误一视同仁地重试",包括文件不存在、格式不支持这种永久性错误,这些重试 3 次纯属浪费;更细的做法是按错误类型区分,永久错误直接 failed,临时错误才退避。TaskMaxRetries=3 时总共有 4 次执行机会(初始 1 次 + 3 次重试)。 代码位置:internal/task/worker.go:148-169(fail 全逻辑)、:132-142(成功路径)、internal/store/task.go:59-69(UpdateTask 写 updated_at)、:77(领取时按 updated_at 过滤)、internal/api/handler_task.go:56-86(手动重试);embedding 重试 internal/embedding/embedder.go:71-96,143-147;配置 internal/config/config.go:537-539、configs/config.yaml:246。

代码依据:

go
func (w *defaultWorkerPool) fail(ctx context.Context, t store.Task, err error) {
    if t.RetryCount < w.cfg.TaskMaxRetries {
        t.Status = store.TaskStatusPending
        t.RetryCount++
        t.ErrorMessage = err.Error()
        // 指数退避:1s/2s/4s/8s...,避免持续错误密集打满外部服务
        backoff := time.Second << uint(t.RetryCount)
        t.UpdatedAt = time.Now().Add(backoff)
        slog.Warn("入库任务失败,退避后重试", "task", t.ID, "retry", t.RetryCount, "backoff", backoff, "err", err)
    } else {
        t.Status = store.TaskStatusFailed
        t.ErrorMessage = err.Error()
        ...
        if err := w.store.UpdateDocumentStatus(ctx, t.DocumentID, store.DocStatusFailed, nil); err != nil {

internal/task/worker.go:148-163

追问链:

  • 追问:第一次失败到底退避几秒? → 答:2 秒,因为我是先 RetryCount++ 再做 1 << RetryCount,序列是 2/4/8,不是 1/2/4。
  • 追问:文件不存在这种错误也会重试 3 次吗? → 答:会,这是我现在的缺陷——没有区分永久错误和临时错误,永久错误应该直接 failed。
  • 追问:用户手动重试有没有次数限制? → 答:没有,接口把 retry_count 清零,可以无限次重试,这也是要补的(加频率限制或总次数上限)。

别踩的雷:

  • ❌ 说"退避 1s/2s/4s"——是 2s/4s/8s(先自增后移位),代码注释写的是 1s/2s/4s/8s 的示意,实际首退是 2s。
  • ❌ 说"用 Redis 延迟队列实现退避"——是写进 updated_at 靠 SQL 过滤。
  • ❌ 说"重试只针对可恢复错误"——现在是一视同仁,主动承认更安全。

Q11. 服务重启后,那些正在处理中的任务会怎样?多实例部署安全吗? ​

面试官想考:恢复语义 + 多实例并发,这是分布式系统的常识题。

口述回答(背诵这段):worker 池启动的时候第一件事就是调 ResetProcessingTasks,执行一条 UPDATE ingest_tasks SET status='pending', updated_at=now() WHERE status='processing',把所有停留在 processing 的任务全部打回 pending,然后才开始起 goroutine 领任务。这样单实例重启场景是安全的:重启前正在跑的任务会被重新捡起来执行,因为 pipeline 是"删后写"的覆盖式幂等,重跑不会产生重复向量。但我要老实承认它的代价:这个重置是无条件全表的,不看任务年龄、不看属于哪个实例,所以多实例部署时不安全——A 实例滚动重启会把 B 实例正在跑的任务抢回 pending,两份 worker 并行入库同一篇文档,两次 DeleteByFilter 和两次 Upsert 交错,可能出现"A 删掉了 B 刚写入的向量"或者同一文档残留两套向量。根因是表里没有 owner、lease、heartbeat 这些字段:正确做法是领取时写 worker_id 和 lease_expires_at,恢复时只重置"租约已经过期"的任务,处理过程中定期续租。另外还有一个隐患是任务重跑的触发点不止重启:如果任务跑成功了但 PG 状态更新失败(我只打了 warn),任务会停在 processing,等下次重启又被重置一遍,等于是用重跑兜底一致性——能收敛,但代价是重复计算。

讲解与备注:这题答好的标志是"先说单实例正确,再主动指出多实例不安全,最后给出 lease 方案"。可以补充一个真实场景说明为什么现在这样写也能用:这个项目实际部署是单实例(cmd/server 一份进程 + 一个 PG + 一个 Qdrant),所以暂停在"启动兜底"这一步是合理的工程取舍;一旦要水平扩容,ResetProcessingTasks 必须先改成"只重置租约过期的任务",否则一定会串。另外可以提一句文件存储也是多实例的障碍:documents.file_path 是本地绝对路径,多实例必须挂共享卷(见 Q2)。 代码位置:internal/task/worker.go:44-47(Start 里调用重置)、internal/store/task.go:99-103(重置 SQL)、:73-96(领取)、internal/app/app.go:122-123(装配时 Start);相关测试 internal/task/worker_test.go:244-257(断言 Start 调用了重置)、internal/store/store_test.go:149-165。

代码依据:

go
// Start 启动 worker 池:先重置 processing 悬挂任务,再启动 WorkerCount 个 worker
func (w *defaultWorkerPool) Start(ctx context.Context) {
    if err := w.store.ResetProcessingTasks(ctx); err != nil {
        slog.Warn("重置悬挂任务失败", "err", err)
    }
    ...
}
// ResetProcessingTasks 把 processing 任务重置为 pending(启动恢复,防悬挂)
func (s *pgStore) ResetProcessingTasks(ctx context.Context) error {
    _, err := s.pool.Exec(ctx,
        `UPDATE ingest_tasks SET status = 'pending', updated_at = now() WHERE status = 'processing'`)
    return err
}

internal/task/worker.go:43-47 + internal/store/task.go:98-103

追问链:

  • 追问:多实例怎么做才对? → 答:领取时记 worker_id,并在任务上写租约过期时间;恢复只重置"租约过期且心跳停止"的任务;处理中定时续租,续租失败就主动放弃并回 pending。
  • 追问:如果任务跑到一半进程被杀,Qdrant 会不会残留半篇向量? → 答:会,但我单次 Upsert 是整批一起提交且 Wait=true,要么没写要么全写;真正的风险是"旧向量已删、新向量没写"的窗口,靠重跑恢复。
  • 追问:有没有 graceful shutdown? → 答:有,Shutdown 会 cancel 上下文并等所有 worker 退出,正在跑的任务因为用了 WithoutCancel 不会被中断,会跑完再退——代价是可能等很久(见 Q12)。

别踩的雷:

  • ❌ 说"多实例也安全,靠数据库行锁"——领取是安全的,但重置不安全,这是两件事。
  • ❌ 说"重启时只重置超时任务"——代码是无条件全表重置。
  • ❌ 说"重启会丢任务"——不会丢,任务在 PG 表里,只是会被重跑一遍。

Q12. 任务状态机有哪几个状态?为什么会出现"卡在 processing"?context.WithoutCancel 是干嘛的? ​

面试官想考:状态机枚举是否背得住、悬挂(stuck)的成因与防护、以及那个不常见的 API 你到底理解没有。

口述回答(背诵这段):任务和文档各自有一套状态机,取值一样,都是四个:pending、processing、completed、failed,定义在 store 包里。流转是:上传写 pending;worker 领取时原子改成 processing;成功写 completed;失败未超限回 pending 并退避;超上限写 failed;启动时把 processing 全部重置回 pending;手动重试只允许 failed 回到 pending。会有两个细节:文档状态永远不会变成 processing——我定义了 DocStatusProcessing 这个常量但从来没写过,所以文档列表里入库中的文档一直显示 pending,这算个小瑕疵;没有 cancelled / dead 这种状态,用户想取消入库是做不到的,只能等它失败或者删文档。至于"卡在 processing",成因有三个:进程在处理中被 kill(没有机会写回状态)、PG 状态更新失败但 Qdrant 已经写完、以及外部 API 长时间不返回导致任务一直占着。防护就是启动时的 ResetProcessingTasks。context.WithoutCancel 是我有意加的:worker 收到的是 wctx,Shutdown 时会 cancel,但我在处理单个任务时用 context.WithoutCancel(ctx) 剥掉取消信号,目的是——服务下线时不要让任务在半路被中断,否则状态可能来不及落库,留下悬挂任务;代价是任务没有取消能力、也没有超时,一个卡住的 embedding 请求会让 Shutdown 一直等下去,唯一兜底是 embedding 客户端 30 秒的 HTTP 超时。更好的做法是给每个任务加一个带超时的 context,让任务有硬上限,同时保证状态写入用一个不被取消的 context。

讲解与备注:这题信息密度高,要按"枚举 → 流转 → 悬挂成因 → WithoutCancel 的作用与代价 → 改进"的顺序讲。context.WithoutCancel 是 Go 1.21 才有的 API,能说清"它保留 value、剥掉取消信号"是加分项。注意"文档状态不写 processing"和"任务状态写 processing"是两件事,别答混:卡在 processing 指的是 ingest_tasks,不是 documents。另外可补:handler_task.go:73-76 明确只允许 failed 重试,processing 的任务既不能重试也不能取消,这也说明缺少人工干预手段(没有 force-reset 接口)。 代码位置:枚举 internal/store/store.go:22-35;领取置 processing internal/store/task.go:75;成功/失败 internal/task/worker.go:132-142,148-165;WithoutCancel internal/task/worker.go:98-101;Shutdown :60-66;重置 internal/store/task.go:99-103;手动重试 internal/api/handler_task.go:73-80;embedding 30s 超时 internal/embedding/embedder.go:37。

代码依据:

go
for _, t := range tasks {
    // 用不随 Shutdown 取消的 ctx 处理任务:保证失败时状态能落库(防残留 processing)
    w.process(context.WithoutCancel(ctx), t)
}
...
// Shutdown 停止 worker,等待当前任务完成
func (w *defaultWorkerPool) Shutdown() {
    if w.cancel != nil { w.cancel() }
    w.wg.Wait()   // 任务不可中断,只能等它跑完
    slog.Info("入库 worker 已停止")
}

internal/task/worker.go:98-101,60-66

追问链:

  • 追问:那 Shamdown 会不会一直卡住? → 答:有可能,任务里没有 deadline,唯一上限是 embedding 的 30 秒 HTTP 超时;如果 parser 卡在超大文件上,就要等到解析完。
  • 追问:为什么不让任务随取消信号中断? → 答:中断会让任务停留在 processing 且没有错误信息,重启才会被重置;我优先保证"状态能落库"。
  • 追问:用户能取消一个正在排队/执行的任务吗? → 答:不能,状态机里没有 cancelled,接口也只提供 failed→pending 的重试,没有取消或强制重置。

别踩的雷:

  • ❌ 说文档状态也会经过 processing——常量存在但从未写入。
  • ❌ 说"Shutdown 会中断正在跑的任务"——不会,WithoutCancel + wg.Wait() 是等它跑完。
  • ❌ 说状态机里有 cancelled/retrying/dead——只有四个。

Q13. PostgreSQL 和 Qdrant 是双写,你怎么保证一致性?先写谁?中间挂了会留下什么? ​

面试官想考:跨存储一致性,这是最能区分"做过"和"想过"的题目。

口述回答(背诵这段):我先说结论:没有分布式事务,也没有 outbox,本质是"顺序写 + 覆盖式补偿"。顺序是:先按 document_id 删掉这篇文档的旧向量和内存 BM25 里的旧条目,然后解析、分块、算向量,先写 Qdrant(一次批量 Upsert,Wait=true 保证服务端确认落盘),再更新内存 BM25,最后回 PG 更新任务状态和文档状态(回填 chunk_ids)。这样排的原因是:chunk_id 是 Qdrant 写入时才生成的,所以 PG 的 chunk_ids 必须等 Qdrant 成功后才能回填;而 BM25 放最后是为了保证"BM25 里出现的 id 在 Qdrant 里一定存在"。中间挂掉会留下三种脏数据:一、Qdrant 写成功但 PG 更新失败,我只打了 warn,后果是 Qdrant 有向量、但 documents.chunk_ids 是空的——这篇文档后续被删除时,向量删不掉,因为删除是按 chunk_ids 逐个删的;二、前置删除失败但我继续入库,旧向量残留、新向量叠加,检索会出现重复结果;三、向量更新成功但 PG 任务状态没写,任务停在 processing,要等下次重启被重置后重跑。补偿手段就是"重跑"——因为入库开头会按 document_id 清旧向量,重跑是覆盖式的,所以最终能收敛,只是要多花一次 embedding 的钱。如果要做得更严,我的方案是:① 加 outbox 表,把"待删/待写"意图先落库再执行;② 写 Qdrant 成功后立刻把 chunk_ids 落库,失败就让任务失败重试;③ 删除文档时改成按 document_id 过滤删除(DeleteByFilter 已经在 pipeline 里实现好了,只是删除路径没用它)。

讲解与备注:这题的黄金句式是"我用的是最终一致,不是强一致;收敛靠 document_id 维度的覆盖式重写"。面试官会满意你主动说出"chunk_ids 为空导致删除漏向量"这个具体的脏数据场景,因为它说明你真推演过失败路径。技术原理可以补:PG 和 Qdrant 之间没有两阶段提交,跨系统 2PC 在工程上基本不可用,所以主流做法就是"幂等重写 + 补偿 + 定期对账"。对账方案我可以给:定时用 PG 的 documents 表去 Qdrant 按 payload 的 kb_id/document_id 做聚合比对,数量或哈希不一致就重跑该文档。另外 BM25 要单独强调——它是纯内存索引,Qdrant 写入成功后同步更新,进程重启后它是空的,Rebuild 方法在全仓非测试代码里没有任何调用点,所以重启后混合检索会静默退化成纯向量检索,这是必须主动承认的问题。 代码位置:pipeline 顺序 internal/pipeline/pipeline.go:66-152;Qdrant 写入与等待 internal/vectorstore/qdrant.go:85-90、DeleteByFilter :168-184;PG 回填 internal/task/worker.go:137-142;删除路径只按 chunk_ids internal/api/handler_doc.go:306-321;BM25 更新 internal/pipeline/pipeline.go:146-152、重启为空 internal/retriever/retriever.go:104,120-123。

代码依据:

go
// Store(payload 携带知识库/文档/chunk 维度)
...
if err := p.vectorstore.Upsert(ctx, records); err != nil {
    return nil, nil, fmt.Errorf("向量入库失败: %w", err)
}

// 更新 BM25 索引(携带知识库 + 文档维度,RemoveByDoc 补偿用)
if p.bm25Index != nil {
    for _, rec := range records {
        content, _ := rec.Payload["content"].(string)
        kbID, _ := rec.Payload["kb_id"].(string)
        p.bm25Index.AddWithDocID(rec.ID, content, kbID, req.DocumentID)
    }
}

internal/pipeline/pipeline.go:141-152

追问链:

  • 追问:为什么必须先写 Qdrant 再写 PG? → 答:chunk_id 是写 Qdrant 时生成的,PG 的 chunk_ids 依赖它;反过来先写 PG 就只能先写空数组,还得再更新一次。
  • 追问:向量写成功、PG 更新失败,用户看到什么? → 答:文档状态还是 pending(或 failed 之前的旧态),但检索其实已经能搜到内容;用户如果删文档,向量删不掉,会残留。
  • 追问:BM25 和向量库不一致怎么办? → 答:BM25 是内存索引、只用于融合排序,即使不一致也不会出现"检索到不存在的文档",因为 RRF 融合时以向量库返回的文档为准;但它重启后为空,混合检索会退化成纯向量。

别踩的雷:

  • ❌ 说"用事务保证 PG 和 Qdrant 一致"——跨系统做不到,会被立刻追问。
  • ❌ 说"删除文档会按 document_id 清向量"——删除路径只按 chunk_ids 逐个删,DeleteByFilter 只在入库补偿里用了。
  • ❌ 说"BM25 有持久化/重启会重建"——没有,Rebuild 零调用点。

Q14. 边界情况:大文件、超长文档、Excel 多 sheet、PDF 分页,你是怎么处理的? ​

面试官想考:有没有想过极端输入,以及知道自己的限制在哪。

口述回答(背诵这段):分几个说。大文件:上传层用 http.MaxBytesReader 卡体积,代码默认上限是 50MB,生产配置里我把它放大到了 1024MB;loader 层没有任何大小保护,PDF、DOCX、Markdown、图片、音频都是 io.ReadAll 一次性读进内存,所以真正的大文件风险在这里,1GB 的 PDF 会直接把内存打满——正确做法是流式解析或者把上限调回几十 MB。超长文档:分块器会把长文本切成多个 chunk,递归策略会按段落和标点切、再贪心合并到接近 chunk_size,单块不会超限;但多媒体不走这条路,一个长视频的转写按 block 一块一个 chunk,可能很长。Excel 多 sheet:每个 sheet 产出一个一级标题 Block,后面每行一个表格 Block,所以 sheet 名天然成了分块的语义边界,标题策略下每个 sheet 会被切成独立的节。PDF:因为有 page 元数据,分块器识别出这是分页文档后按页分块,同一页的所有 Block 合并成一个 chunk 并把页码写进元数据;单页超过 chunk_size 时降级用递归策略再切,子 chunk 继承同一页码,所以引用的时候能精确跳到第几页。音频视频还有时间戳:每个 chunk 带 start_ms/end_ms,检索结果能定位到视频的具体时间点。

讲解与备注:这题的加分项是"PDF 按页、Excel 按 sheet 是我有意设计的语义边界",而不是"随便切的"。同时要主动给出量化边界:① 上传上限代码 50MB / 生产 1024MB(这个放大是我配的,属于风险点);② loader 无体积保护、无页数上限(PageCount 只写进元数据没用来校验);③ 无 OCR、无表格结构还原;④ 明文限制"1GB 文件会 OOM"。面试官听到你主动说"我把上限放大到 1024MB 但 loader 没有流式解析,这是隐患"会给你加分,因为这是真实的取舍。改进方案:上传上限按场景压到 50-100MB、PDF 改成逐页流式处理、给解析加超时和内存上限、大文件走异步预校验并把源文件放对象存储。 代码位置:体积上限 internal/api/handler_doc.go:42-43、默认值 internal/config/config.go:531-533、生产值 configs/config.yaml:242;io.ReadAll 位置 internal/loader/parser_pdf.go:29、parser_docx.go:27、parser_markdown.go:28、parser_image.go:50、parser_audio.go:48;PDF 逐页 parser_pdf.go:58-88;Excel 多 sheet parser_excel.go:43-48;按页分块 internal/chunker/chunker.go:176-222;媒体按块 :136-162。

代码依据:

go
// chunkPagedBlocks PDF 按页分块:同页 block 合并为一个 chunk,PageNumber=page;
// 单页超长时用 recursive 切分,子 chunk 沿用同一 PageNumber。
for _, b := range doc.Blocks {
    page, _ := b.Metadata["page"].(int)
    if page <= 0 { page = 1 }
    idx, ok := indexByPage[page]
    ...
}
...
text := strings.TrimSpace(blocksToText(g.blocks))
parts := []string{text}
if c.tokenizer.Count(text) > config.ChunkSize {
    parts = fallback.Split(text, config, c.tokenizer)
}

internal/chunker/chunker.go:174-208

追问链:

  • 追问:100MB 的 PDF 会怎样? → 答:会被读进内存,可能 OOM;loader 层没有大小和页数上限,这是要补的。
  • 追问:Excel 10 万行会怎样? → 答:excelize.GetRows 全量读进内存,然后每行一个 Block,Block 数量会非常大,分块和 embedding 都会很慢;我应该加行数上限和分批。
  • 追问:视频很长的转写会不会超模型上下文? → 答:媒体走按块分块,不切分,所以理论上单 chunk 可以很长;不过它受 ASR 分段粒度约束,通常是几秒到几十秒一段。

别踩的雷:

  • ❌ 说"上传上限是 50MB"——代码默认 50MB,但生产 YAML 是 1024MB,两个都要提。
  • ❌ 说"loader 有大小/页数限制"——没有,代码中未找到。
  • ❌ 说"大文件是流式处理的"——除了视频落临时文件,其余都是 io.ReadAll。

Q15. 入库过程的可观测性怎么样?线上一个任务很慢,你怎么定位? ​

面试官想考:有没有线上意识,是不是只写了功能没写日志。

口述回答(背诵这段):可观测性目前是结构化日志 + 状态字段这个级别,没有接指标、没有 trace。具体有三层:第一,worker 启动时会打 入库 worker 已启动 并带上并发数;第二,每个任务成功时打一条 入库任务完成,字段里有 task id、doc id、chunk 数量,还有一个耗时毫秒数,是我在 process 里用 time.Now() 前后相减算的,这条日志就是定位慢任务的主要依据;第三,所有失败路径都是 slog.Warn/Error,包括领取失败、更新任务状态失败、更新文档状态失败、重试退避、超上限失败。定位慢任务的流程就是:先看这条完成日志的耗时分布,如果只是个别慢,通常是那篇文档特别大或者 embedding 限流排队;如果是整体变慢,就看是不是 embedding 服务的 QPS 限制(我配了限流器)、或者 worker 数量不够导致队列堆积——队列堆积可以直接查 ingest_tasks 里 pending 的数量和最早的 created_at。不够的地方我也很清楚:① 耗时只打了总数,没有分阶段(load/chunk/embed/store 各占多少),所以定位不到是哪一步慢;② 任务一旦卡死在 processing,没有任何告警,只能等人发现;③ 没有 metrics(队列深度、成功率、P95 耗时),也没有把这些日志接到 Prometheus。改进方向就是给每阶段加打点、暴露队列深度和耗时的 Prometheus 指标、加一个"processing 超过 N 分钟"的巡检告警。

讲解与备注:这题不要吹,说清"有什么、缺什么、怎么补"就很好。可以强调一个细节:耗时是用 time.Since(start).Milliseconds() 算的,单位是毫秒(worker.go:143-144),字段名是中文 耗时ms。也可以顺带说 warning_message 的用途:视频没配语音能力时会写进去,任务还是 completed,这样用户能区分"失败了"和"成功了但降级了"。加分点是提出"processing 超时告警",因为它直接对上了 Q12 的悬挂问题。 代码位置:完成日志 internal/task/worker.go:143-144;启动日志 :56;各类 warn :68,84,138,141,156,160,163,167;状态字段 internal/store/schema.go:29-39,99-101。

代码依据:

go
start := time.Now()
chunkIDs, warnings, err := w.pipeline.Ingest(ctx, pipeline.IngestRequest{...})
if err != nil {
    w.fail(ctx, t, err)
    return
}
...
if len(warnings) > 0 {
    t.WarningMessage = strings.Join(warnings, "; ")
}
slog.Info("入库任务完成", "task", t.ID, "doc", t.DocumentID, "chunks", len(chunkIDs),
    "耗时ms", time.Since(start).Milliseconds())

internal/task/worker.go:120-144

追问链:

  • 追问:怎么知道是 embedding 慢还是解析慢? → 答:现在知道不了,只有总耗时;要加阶段打点,在 Load、Chunk、Embed、Store 四处各记一次。
  • 追问:怎么知道队列积压? → 答:查 ingest_tasks 里 status='pending' 的数量和最早 created_at;没有指标面板,这也是要补的。
  • 追问:失败任务有人知道吗? → 答:只有日志 error 和任务状态 failed,用户在文档列表能看到失败状态;没有告警推送。

别踩的雷:

  • ❌ 说"接了 Prometheus/Grafana/OpenTelemetry"——代码里没有,只有 log/slog。
  • ❌ 说"有分阶段耗时统计"——只有一个总耗时。
  • ❌ 说"任务卡住会告警"——没有任何巡检或告警。

Q16. embedding 这一侧你是怎么调的?批量、限流、重试、维度校验都是怎么做的? ​

面试官想考:外部依赖的工程化处理——这是"能上线"和"能跑通"的分界线。

口述回答(背诵这段):入库时我是把整篇文档的所有 chunk 文本一次性交给 embedder,由 embedder 内部按 batch_size 切片串行发请求。我配的 batch_size 是 10,代码默认是 100,qps 是 10,max_retries 是 3。限流用的是 golang.org/x/time/rate 的限流器,每一批请求前先 Wait,这样能把对外部 API 的并发压住,不会因为上传量大把对方打挂或者触发 429。重试只针对可重试错误:网络错误、HTTP 429、HTTP 5xx 会被包成 RetryableError,退避是 1 秒、2 秒、4 秒;HTTP 4xx(比如鉴权失败、参数错误)直接返回不重试,因为重试没有意义。向量回来之后我会校验返回条数和 chunk 条数是否相等,不等就报错,避免出现"少一个向量但 chunk_id 已经对齐错位"这种更隐蔽的问题;返回结果按响应里的 index 字段回填下标,不是靠数组顺序猜。HTTP 客户端超时是 30 秒,写在代码里没有做成配置项,这是个可以做的小改进。另外向量维度这块:embedder 配了 1024 维,Qdrant collection 也是 1024 维,pipeline 里没有做二次校验,是靠 Qdrant 自己拒绝不一致的向量。

讲解与备注:得分点是"限流 + 只重试可重试错误 + 按 index 回填 + 数量校验"这四个工程细节,它们说明你把外部依赖当成不可靠资源来处理。可以补一句取舍:"批量是为了减少往返、限流是为了保护下游,这两个参数是矛盾的两端,我用 batch_size=10 + qps=10 取的是稳而不是快。" 也可以主动承认:超时 30s 硬编码、Qdrant 的 Upsert 没有重试(只有 embedding 有重试),都是可以改进的点。 代码位置:接口 internal/embedding/embedder.go:17-19;构造(超时/限流器):28-40;批量与限流 :42-69;单批重试 :71-96;可重试判定 :132-147;按 index 回填 :158-163;pipeline 侧数量校验 internal/pipeline/pipeline.go:110-112;配置默认 internal/config/config.go:381-392、YAML configs/config.yaml:14,16,18,20。

代码依据:

go
for i := 0; i < len(texts); i += e.config.BatchSize {
    end := i + e.config.BatchSize
    if end > len(texts) { end = len(texts) }
    batch := texts[i:end]

    if err := e.limiter.Wait(ctx); err != nil {
        return nil, fmt.Errorf("限流等待失败: %w", err)
    }
    vectors, err := e.embedBatch(ctx, batch)
    ...
}
vectors := make([][]float32, len(texts))
for _, d := range embResp.Data {
    if d.Index < len(vectors) { vectors[d.Index] = d.Embedding }
}

internal/embedding/embedder.go:49-66,158-163

追问链:

  • 追问:为什么要限流而不是直接并发? → 答:入库是后台任务,突发并发会把下游打爆并触发 429,反而更慢;限流让吞吐稳定可控。
  • 追问:批内有一条失败会怎样? → 答:整批失败并按批重试,不会只重试失败那一条;粒度粗但实现简单,重试次数上限 3,之后整个任务进任务级重试。
  • 追问:向量维度不匹配会怎样? → 答:pipeline 不校验,Qdrant 写入会报错,任务失败重试;应该在启动时校验 embedder 配置维度和 collection 维度是否一致。

别踩的雷:

  • ❌ 说"每个 chunk 一个请求"——是按 batch_size 批量发的。
  • ❌ 说"所有错误都重试"——只有网络错误 / 429 / 5xx 重试,4xx 不重试。
  • ❌ 说 batch_size 默认 100 但没说实际配的是 10——两个都要说。

Q17. 图片、音频、视频这条线在入库链路里是怎么走的?能力没配置会怎样? ​

面试官想考:多媒体这块是不是"简历写了但没实现",以及依赖倒置的设计能力。

口述回答(背诵这段):多媒体走的是同一套 Loader 抽象,只是能力靠接口注入。我在 loader 包里定义了三个能力接口:VisionProvider(图片/视频帧出文本描述)、SpeechProvider(音频出带时间戳的分段)、还有抽帧和探测的接口;实现放在 multimedia 包,loader 不反向依赖它,就是为了可插拔替换。装配的时候,如果配置里 vision 或 speech 的 api_key 是空的,NewVisionProvider 会返回 nil,注册表里对应的 parser 拿着 nil;上传预检时 parser 实现了一个 CheckCapabilities 方法,registry 一旦发现它返回错误就直接 400 拒绝,错误信息是"multimedia.vision 未配置,无法处理该类型文件"。这样设计的目的是防止产生"空文档"脏数据——如果允许上传,最后会得到一个没有任何内容的文档,用户不知道为什么搜不到。图片的处理是:读字节、尽量拿宽高(拿不到就不阻断)、调视觉模型出描述,产出一个图片描述 Block,把宽高和来源文件名写进元数据;音频是把转写按时间戳分成多个音频分段 Block;视频最重——先落一个临时文件,因为 ffmpeg 需要文件路径输入,然后 ffprobe 探测有没有音轨、时长、编码,再按固定间隔或者场景检测抽帧,每帧调视觉模型出描述并带上起止时间戳,音轨单独拆出来走 ASR;如果只配了视觉没配语音,视觉照常做,音轨跳过并记一条 warning,任务还是成功,这就是"非阻断降级"。

讲解与备注:这题的亮点是"能力缺失在预检阶段就明确 400 拒绝,而不是让它变成空文档",以及"视觉和音频两条轨道独立降级"。面试官可能会追问"为什么抽帧要落临时文件"——因为 ffmpeg 需要 seek,stdin 不方便,而且要防命令注入所以我用数组传参而不是拼 shell 字符串(exec.CommandContext)。还要能说出场景检测那条线:frame_strategy=scene 的时候需要额外配视觉 embedding 做帧间相似度,没配的话抽帧策略构造返回 nil,CheckCapabilities 会报"scene 需要配置 vision_embedding",同样在上传阶段拒绝。我要主动承认的边界:图片/音频/视频都要调外部多模态服务,成本和延迟都远高于文本;媒体 chunk 不做长度切分,长视频转写会是大 chunk。 代码位置:能力接口 internal/loader/capability.go:13-118;能力缺失错误 :64-76;图片 internal/loader/parser_image.go:34-93;音频 parser_audio.go:32-95;视频 parser_video.go:37-183;provider 构造(api_key 为空返回 nil)internal/multimedia/provider.go:17-51;装配覆盖注册 internal/app/app.go:228-252;预检拒绝 internal/api/handler_doc.go:74-78 + internal/loader/support.go:21-32。

代码依据:

go
// MediaCapabilityChecker 能力缺失检查。
// 多媒体 Parser 实现该接口:上传预检阶段只调用 CheckCapabilities 判定能力是否可用
type MediaCapabilityChecker interface {
    CheckCapabilities() error // 能力缺失返回 *ErrMediaCapabilityMissing
}

func (r *defaultRegistry) Support(info FileInfo) SupportResult {
    parser, err := r.Resolve(info)
    if err != nil {
        return SupportResult{Supported: false, Reason: err.Error()}
    }
    if checker, ok := parser.(MediaCapabilityChecker); ok {
        if err := checker.CheckCapabilities(); err != nil {
            return SupportResult{Supported: false, Reason: err.Error()}
        }
    }
    return SupportResult{Supported: true}
}

internal/loader/capability.go:51-56 + internal/loader/support.go:21-32

追问链:

  • 追问:为什么把接口定义在 loader 包里? → 答:依赖倒置——parser 只依赖抽象,multimedia 包实现它,loader 不 import multimedia,所以换成别的视觉模型只需要换实现。
  • 追问:视频没配语音能力会怎样? → 答:视觉帧照常出 chunk,音轨跳过并在 ingest_tasks.warning_message 记一条告警,任务状态是 completed。
  • 追问:音轨转写的时间戳准确吗? → 答:取决于 ASR 提供方;代码注释里写了 dashscope 的 qwen ASR 不返回时间戳,那种情况下音频 chunk 的时间戳是 0,跳转定位会失效。

别踩的雷:

  • ❌ 说"图片走 OCR"——图片走的是视觉模型描述,OCR 代码里没有。
  • ❌ 说"没配能力也能上传,只是不处理"——上传阶段就 400 拒绝了。
  • ❌ 说"视频抽帧用的是外部服务"——是本地 ffmpeg 抽帧,只有理解和转写调外部模型。

三、诚实承认的不足与改进 ​

不足(代码事实 + 位置)面试怎么讲得体改进方案
没有幂等键:chunk_id 是 UUID v4(pipeline.go:118),同一文件二次上传会产生两份内容(handler_doc.go:92 新 docID)"我早期按 document_id 维度做覆盖式重写,重试是幂等的,但内容级幂等我确实没做——同一文件重传会真重复,这是我知道的缺陷。"chunk_id 改 sha256(document_id+index+content);上传算文件 hash,KB 内命中就提示已存在;PG 加 (kb_id, content_hash) 唯一索引
没有死信队列 / 人工干预通道:失败任务停在 failed(worker.go:157-165),只有 POST /tasks/:id/retry 手动重试(handler_task.go:56-86),无 DLQ、无告警"失败超过 3 次就置 failed 留在表里,靠用户在文档列表看到再点重试;没有自动 DLQ 和告警,任务多了会很难发现。"加失败任务告警(超阈值推送/计数指标);建 failed 视图或归档表;提供批量重试接口
没有背压/准入控制:上传只管插 pending(handler_doc.go:133),消费端固定 5 个 worker、每次 1 条(worker.go:19-21)"入口没有队列深度限制,也不按知识库限流,极端情况下任务表会堆积、检索端可能先看到半成品。"加 pending 数量阈值,超阈值返回 429;按 KB 做并发配额;任务表加 TTL/归档
多实例领取与恢复不安全:ResetProcessingTasks 无条件全表重置 processing(task.go:99-103),表里没有 owner/lease/heartbeat"单实例重启是安全的;多实例我承认有问题——滚动重启会把别人正在跑的任务抢回 pending,必须加租约才能水平扩容。"加 worker_id + lease_expires_at,处理中续租,恢复只重置过期租约;领取加 FOR UPDATE SKIP LOCKED
跨存储无事务、无 outbox:Qdrant 写成功但 PG 更新失败只打 warn(worker.go:138,141),可能造成 chunk_ids 为空导致删向量漏删(handler_doc.go:307-317)"我用最终一致 + 覆盖式重写来收敛,代价是要重跑一次;具体脏数据场景是 chunk_ids 为空时删除会漏向量,这个我说得出来。"写 Qdrant 成功后立刻落 chunk_ids,失败即任务失败;删除文档改用 DeleteByFilter(document_id);定时对账任务
BM25 是纯内存、重启不重建:Rebuild 在生产代码零调用(bm25.go:217-231),重启后 DocCount()==0,混合检索退化为纯向量(retriever.go:104,120-123)"BM25 只服务融合排序,不影响数据正确性,但重启后混合检索会静默降级,这是我最想先修的一个点。"启动时从 Qdrant 拉全部 chunk 调 Rebuild 预热;或给 BM25 做落盘快照
没有取消/超时能力:状态机只有四个状态(store.go:22-35),用 context.WithoutCancel 处理任务(worker.go:100),任务无 deadline,Shutdown 可能长时间阻塞"为了让状态能落库,我主动放弃了任务级取消和超时,代价是下线可能慢、卡死的任务只能等外部 30 秒超时;状态机里也没有 cancelled。"给每个任务加带超时的 context(如 10 分钟);增加 cancelled 状态 + 取消接口;外部调用统一带 deadline
重试不区分错误类型:所有错误一视同仁重试 3 次(worker.go:149),文件不存在、格式不支持这种永久错误也重试"现在是无差别重试,永久错误重试纯属浪费 2s/4s/8s 三次;应该按错误类型分流。"定义可重试错误标记(复用 embedding 的 RetryableError 思路),永久错误直接 failed;同时把退避单独放 next_attempt_at 列
updated_at 语义复用:退避时间写进 updated_at(worker.go:155),领取时用它过滤(task.go:77)"我用 updated_at 兼作 next_attempt_at,省了一张调度表,但破坏了字段语义,监控和归档都会被未来时间污染。"新增 next_attempt_at 列,updated_at 只表示最后更新时间
loader 无体积/页数保护:上传上限生产配成 1024MB(configs/config.yaml:242),而 loader 全部 io.ReadAll(parser_pdf.go:29 等),PageCount 只写元数据不校验(parser_pdf.go:95)"上传被放大到 1GB 但解析是全量读内存,这是明显的隐患,应该回到 50-100MB 或改成流式解析。"上限调回并做成按格式可配;PDF 逐页流式处理;加解析超时与内存上限
分块质量缺陷(可主动坦白):递归策略 strings.Split 丢弃标点(strategy_recursive.go:52-64);非标题策略的 heading_context 取的是文档最后一个标题路径(chunker.go:245-259);ChunkOverlap 判据 < 0 导致省略配置即 0(types.go:41-43)"切分后我没把分隔符拼回去,长文档会丢句末标点;breadcrumb 那个函数名和实现不符,返回的是最后一个标题栈;overlap 的默认值判据也有 bug——这三个我都是有具体位置能改的。"切分后把分隔符附加到前段末尾;按 chunk 起始位置回溯标题栈;判据改成 <= 0

四、背诵清单(10 条一句话要点) ​

  1. 边界:上传同步只做体积校验、kb 校验、格式/能力预检、可读性预检(min_readable_chars 默认 20)、落盘、两条 INSERT(documents + ingest_tasks 都是 pending),然后立刻返回 task_id;重活全在 worker。handler_doc.go:38-138
  2. 队列:异步队列就是 PG 的 ingest_tasks 表,worker 每 500ms 轮询、每次领 1 条,没有 Redis/Kafka/channel。worker.go:18-22
  3. 领取 SQL:UPDATE ... WHERE id IN (SELECT ... status='pending' AND updated_at <= NOW() ORDER BY created_at LIMIT 1) RETURNING,选与占在同一条语句里所以不会重复领;没有 FOR UPDATE SKIP LOCKED。task.go:73-96
  4. 并发:worker 数代码默认 2、生产 YAML 5;领取出错退避 1s。config.go:534-536、configs/config.yaml:244
  5. 重试:默认最多 3 次,退避写在 updated_at 里,序列是 2s / 4s / 8s(先自增后移位),领取时靠 updated_at <= NOW() 跳过退避期;超上限置 failed 并同步文档 failed。worker.go:148-165
  6. 恢复:启动时 ResetProcessingTasks 把所有 processing 无条件打回 pending —— 单实例安全,多实例会互相抢(无 lease/heartbeat)。worker.go:44-47、task.go:99-103
  7. 分块:三种策略固定/递归/标题,默认 chunk_size 512(YAML 500)、chunk_overlap 50(判据 < 0 才生效)、heading_level 2;路由优先级是多媒体按块 > PDF 按页 > 配置策略。types.go:37-48、chunker.go:41-56
  8. token:不是真 tokenizer,是启发式——汉字算 2 token、标点 1、连续非空白字符(英文词)1,所以 500 token ≈ 250 个汉字;Tokenizer 是接口,可替换。tokenizer.go:20-57
  9. chunk_id:uuid.New().String()(UUID v4,不是内容哈希),同时是 Qdrant point id / BM25 docID / PG chunk_ids 元素;幂等靠 document_id 维度的"先删后写"(DeleteByFilter + RemoveByDoc),同一文件重传会产生真重复。pipeline.go:66-73,116-119
  10. 双写:顺序是 清旧向量 → 写 Qdrant(Wait=true)→ 更新内存 BM25 → 回写 PG;无事务、无 outbox,PG 更新失败只告警,靠重跑收敛;BM25 无持久化,重启后混合检索退化为纯向量。pipeline.go:141-152、bm25.go:217-231

附加保险句(被问到"你觉得自己最大的问题"时用):"最想先修三件:BM25 重启预热、多实例的任务租约、以及把 chunk_id 改成内容哈希做真幂等;这三件事分别对应检索质量、可扩容性和数据正确性。"

持续学习,持续构建。