Chuyển đến nội dung chính

Lesson 5: RAG Pipeline — Vector Store, Chunking, Hybrid Search & Re-ranking

End-to-end RAG pipeline, document ingestion (PDF/HTML/DOCX/code), chunking strategies (semantic, recursive, sentence-window), embedding models, vector store (Qdrant/Pgvector), hybrid search (BM25 + semantic), re-ranking (Cohere/cross-encoder), citation generation.

🏗️ Architecture — Lesson 5 Lesson 5: RAG Pipeline — Vector Store, Chunking, Hybrid Search & Re-ranking

Enterprise AI Chatbot Platform Architecture — From Prototype to Production

Part 2: Core Chatbot Engine

xdev.asia

1. RAG in AI Chatbot — Why is it needed?

LLM only knows the training data. Chatbot enterprise needs to respond based on private data of the organization: documentation, policy, product information, FAQs. RAG (Retrieval-Augmented Generation) solves this problem.


┌──────────────────────── RAG PIPELINE ────────────────────────────┐
│                                                                    │
│  INGESTION (Offline)                  RETRIEVAL (Online)           │
│  ┌──────────┐                         ┌──────────┐                │
│  │ Documents│──▶ Parse ──▶ Chunk ──▶  │  Query   │                │
│  │ PDF/HTML │     │          │     │   │          │                │
│  │ DOCX/MD  │     ▼          ▼     │   └────┬─────┘                │
│  └──────────┘   Clean     Embed   │        │                      │
│                   │          │     │   ┌────▼─────┐                │
│                   ▼          ▼     │   │  Hybrid  │                │
│               ┌──────┐  ┌──────┐  │   │  Search  │                │
│               │ BM25 │  │Vector│  │   │BM25+Embd │                │
│               │Index │  │Store │  │   └────┬─────┘                │
│               └──────┘  └──────┘  │        │                      │
│                                   │   ┌────▼─────┐                │
│                                   │   │ Re-rank  │                │
│                                   │   └────┬─────┘                │
│                                   │        │                      │
│                                   │   ┌────▼─────┐ ┌──────────┐  │
│                                   │   │  Format  │─▶│ Citation │  │
│                                   │   │ Context  │ │ Generate │  │
│                                   │   └──────────┘ └──────────┘  │
└──────────────────────────────────────────────────────────────────┘

2. Document Ingestion Pipeline


interface DocumentParser {
  parse(input: Buffer, mimeType: string): Promise<ParsedDocument>;
}

interface ParsedDocument {
  title: string;
  content: string;         // Plain text content
  sections: Section[];     // Structured sections with headings
  metadata: {
    source: string;
    mimeType: string;
    pageCount?: number;
    lastModified?: Date;
    language?: string;
  };
}

class IngestionPipeline {
  private parsers: Map<string, DocumentParser> = new Map([
    ['application/pdf', new PDFParser()],
    ['text/html', new HTMLParser()],
    ['text/markdown', new MarkdownParser()],
    ['application/vnd.openxmlformats-officedocument.wordprocessingml.document', new DOCXParser()],
    ['text/plain', new PlainTextParser()],
    ['application/json', new JSONParser()],
  ]);

  async ingest(
    tenantId: string,
    knowledgeBaseId: string,
    file: UploadedFile,
  ): Promise<IngestionResult> {
    // 1. Parse document
    const parser = this.parsers.get(file.mimeType);
    if (!parser) throw new UnsupportedFormatError(file.mimeType);
    const parsed = await parser.parse(file.buffer, file.mimeType);

    // 2. Clean text
    const cleaned = this.cleanText(parsed.content);

    // 3. Chunk
    const chunks = await this.chunker.chunk(cleaned, {
      strategy: 'recursive',
      chunkSize: 512,
      chunkOverlap: 50,
    });

    // 4. Generate embeddings
    const embeddings = await this.embeddingService.embedBatch(
      chunks.map(c => c.content),
    );

    // 5. Store in vector database
    const points = chunks.map((chunk, i) => ({
      id: crypto.randomUUID(),
      vector: embeddings[i],
      payload: {
        tenantId,
        knowledgeBaseId,
        documentId: file.id,
        content: chunk.content,
        chunkIndex: i,
        metadata: {
          title: parsed.title,
          section: chunk.section,
          source: file.name,
          pageNumber: chunk.pageNumber,
        },
      },
    }));

    await this.vectorStore.upsert('knowledge_chunks', points);

    // 6. Index for BM25 full-text search
    await this.bm25Index.index(points.map(p => ({
      id: p.id,
      content: p.payload.content,
      metadata: p.payload.metadata,
    })));

    return {
      documentId: file.id,
      chunksCreated: chunks.length,
      totalTokens: chunks.reduce((sum, c) => sum + c.tokenCount, 0),
    };
  }

  private cleanText(text: string): string {
    return text
      .replace(/\s+/g, ' ')           // Normalize whitespace
      .replace(/\n{3,}/g, '\n\n')     // Max 2 newlines
      .replace(/[^\S\n]+/g, ' ')      // Collapse spaces (keep newlines)
      .trim();
  }
}

3. Chunking Strategies

StrategyHow it worksBest for
Fixed-sizeSplit every N tokensSimple, general purpose
RecursiveSplit by headers → paragraphs → sentencesStructured documents
SemanticsSplit when embedding similarity dropsUnstructured, narrative text
Sentence-windowChunk = sentence, context = surrounding sentencesPrecise retrieval + context
Parent-childSmall chunks for search, return parent for contextLong documents

class RecursiveChunker {
  private readonly SEPARATORS = [
    '\n## ',   // H2 headers
    '\n### ',  // H3 headers
    '\n\n',    // Paragraphs
    '\n',      // Lines
    '. ',      // Sentences
    ' ',       // Words
  ];

  async chunk(
    text: string,
    config: { chunkSize: number; chunkOverlap: number },
  ): Promise<Chunk[]> {
    return this.recursiveSplit(text, this.SEPARATORS, config);
  }

  private recursiveSplit(
    text: string,
    separators: string[],
    config: { chunkSize: number; chunkOverlap: number },
  ): Chunk[] {
    if (text.length <= config.chunkSize) {
      return [{ content: text, tokenCount: this.estimateTokens(text) }];
    }

    const separator = separators[0];
    const parts = text.split(separator);

    const chunks: Chunk[] = [];
    let currentChunk = '';

    for (const part of parts) {
      const candidate = currentChunk
        ? currentChunk + separator + part
        : part;

      if (this.estimateTokens(candidate) > config.chunkSize) {
        if (currentChunk) {
          chunks.push({
            content: currentChunk.trim(),
            tokenCount: this.estimateTokens(currentChunk),
          });
        }

        if (this.estimateTokens(part) > config.chunkSize && separators.length > 1) {
          // Recursively split with next separator
          chunks.push(...this.recursiveSplit(part, separators.slice(1), config));
        } else {
          currentChunk = part;
        }
      } else {
        currentChunk = candidate;
      }
    }

    if (currentChunk) {
      chunks.push({
        content: currentChunk.trim(),
        tokenCount: this.estimateTokens(currentChunk),
      });
    }

    return this.addOverlap(chunks, config.chunkOverlap);
  }

  private estimateTokens(text: string): number {
    return Math.ceil(text.length / 4); // Rough estimate
  }
}

class HybridSearchEngine {
  constructor(
    private vectorStore: VectorStore,   // Qdrant
    private bm25Index: BM25Index,       // Elasticsearch/Meilisearch
    private embeddingService: EmbeddingService,
  ) {}

  async search(
    query: string,
    options: SearchOptions,
  ): Promise<SearchResult[]> {
    // Run both searches in parallel
    const queryEmbedding = await this.embeddingService.embed(query);

    const [semanticResults, bm25Results] = await Promise.all([
      this.vectorStore.search({
        collection: 'knowledge_chunks',
        vector: queryEmbedding,
        filter: {
          tenantId: options.tenantId,
          knowledgeBaseId: options.knowledgeBaseId,
        },
        limit: options.topK * 2,
      }),
      this.bm25Index.search({
        query,
        filter: {
          tenantId: options.tenantId,
          knowledgeBaseId: options.knowledgeBaseId,
        },
        limit: options.topK * 2,
      }),
    ]);

    // Reciprocal Rank Fusion (RRF) to combine results
    return this.reciprocalRankFusion(semanticResults, bm25Results, options.topK);
  }

  private reciprocalRankFusion(
    semanticResults: SearchResult[],
    bm25Results: SearchResult[],
    topK: number,
    k: number = 60, // RRF constant
  ): SearchResult[] {
    const scoreMap = new Map<string, { result: SearchResult; score: number }>();

    // Score semantic results
    semanticResults.forEach((result, rank) => {
      const rrf = 1 / (k + rank + 1);
      const existing = scoreMap.get(result.id);
      if (existing) {
        existing.score += rrf;
      } else {
        scoreMap.set(result.id, { result, score: rrf });
      }
    });

    // Score BM25 results
    bm25Results.forEach((result, rank) => {
      const rrf = 1 / (k + rank + 1);
      const existing = scoreMap.get(result.id);
      if (existing) {
        existing.score += rrf;
      } else {
        scoreMap.set(result.id, { result, score: rrf });
      }
    });

    // Sort by combined RRF score
    return Array.from(scoreMap.values())
      .sort((a, b) => b.score - a.score)
      .slice(0, topK)
      .map(entry => ({ ...entry.result, score: entry.score }));
  }
}

5. Re-ranking — Cross-Encoder


interface ReRanker {
  rerank(query: string, documents: SearchResult[], topK: number): Promise<SearchResult[]>;
}

class CohereReRanker implements ReRanker {
  constructor(private apiKey: string) {}

  async rerank(
    query: string,
    documents: SearchResult[],
    topK: number,
  ): Promise<SearchResult[]> {
    const response = await fetch('https://api.cohere.ai/v1/rerank', {
      method: 'POST',
      headers: {
        'Authorization': `Bearer ${this.apiKey}`,
        'Content-Type': 'application/json',
      },
      body: JSON.stringify({
        model: 'rerank-english-v3.0',
        query,
        documents: documents.map(d => d.content),
        top_n: topK,
        return_documents: false,
      }),
    });

    const data = await response.json();

    return data.results.map((r: { index: number; relevance_score: number }) => ({
      ...documents[r.index],
      score: r.relevance_score,
    }));
  }
}

6. Citation Generation


class CitationGenerator {
  formatContextWithCitations(results: SearchResult[]): {
    context: string;
    citations: Citation[];
  } {
    const citations: Citation[] = [];

    const context = results
      .map((result, index) => {
        const citationId = `[${index + 1}]`;
        citations.push({
          id: citationId,
          documentId: result.metadata.documentId,
          chunkId: result.id,
          content: result.content.substring(0, 200),
          source: result.metadata.source,
          pageNumber: result.metadata.pageNumber,
          score: result.score,
        });

        return `${citationId} ${result.content}`;
      })
      .join('\n\n');

    return { context, citations };
  }

  // System prompt instruction for citation
  getCitationInstruction(): string {
    return `When answering, cite sources using [1], [2], etc.
Only cite information that directly comes from the provided context.
If the context doesn't contain the answer, say you don't have enough information.
Do NOT make up information or citations.`;
  }
}

7. Full RAG Pipeline Service


class RAGPipeline {
  constructor(
    private searchEngine: HybridSearchEngine,
    private reRanker: ReRanker,
    private citationGenerator: CitationGenerator,
    private queryTransformer: QueryTransformer,
  ) {}

  async search(
    query: string,
    options: RAGOptions,
  ): Promise<RAGResult> {
    // 1. Query transformation (expand, decompose, HyDE)
    const transformedQuery = await this.queryTransformer.transform(query);

    // 2. Hybrid search (BM25 + Semantic)
    const searchResults = await this.searchEngine.search(transformedQuery, {
      tenantId: options.tenantId,
      knowledgeBaseId: options.knowledgeBaseId,
      topK: 20, // Retrieve more for re-ranking
    });

    // 3. Re-rank
    const reranked = await this.reRanker.rerank(query, searchResults, options.topK ?? 5);

    // 4. Format with citations
    const { context, citations } = this.citationGenerator.formatContextWithCitations(reranked);

    return { context, citations, totalResults: searchResults.length };
  }
}

Summary of Lesson 5

StageComponentKey Decision
IngestionParser + Chunker + EmbedderChunk size 512 tokens, 50 overlap
StorageVector Store + BM25 IndexQdrant for semantics, Meilisearch for keywords
RetrievalHybrid SearchRRF fusion (k=60) for combining scores
Re-rankingCross-encoderCohere rerank-v3 or local cross-encoder
CitationCitation generatorInline citations [1][2] with source links

Next article: Prompt Engineering Engine — template system, chain-of-thought, dynamic prompt assembly, persona management, and prompt A/B testing.