{"slug":"rag-pipeline-embeddings","title":"RAG Pipeline: Chunking, Embeddings, Vector Search, and Retrieval","tags":["rag","embeddings","vector-db","supabase","pgvector","retrieval","ai"],"agent_summary":"Production RAG pipeline — document chunking strategy (500 tokens, 50 overlap), Voyage AI embeddings, pgvector in Supabase, hybrid search (semantic + BM25 with RRF), reranking, and augmented generation with citations.","trigger_phrases":["RAG pipeline","embeddings","vector search","pgvector","semantic search","Voyage AI","document chunking","hybrid search","retrieval augmented generation"],"runnable":false,"markdown":"\n## Overview\n\nRAG (Retrieval-Augmented Generation) grounds LLM responses in your documents. Pipeline: chunk → embed → store → search → augment → generate.\n\n## Embedding Model Selection\n\n| Model | Provider | Dimensions | Use Case |\n|-------|---------|-----------|----------|\n| `voyage-3` | Voyage AI | 1024 | General documents (best quality) |\n| `voyage-3-lite` | Voyage AI | 512 | High volume, cost-sensitive |\n| `text-embedding-3-small` | OpenAI | 1536 | Fallback, widely supported |\n| `text-embedding-3-large` | OpenAI | 3072 | Maximum quality (high cost) |\n\nVoyage AI: 200M free tokens/month. Preferred over OpenAI for cost and quality balance.\n\n## Document Chunking\n\n```python\nfrom typing import Generator\n\ndef chunk_document(\n    text: str,\n    chunk_size: int = 500,     # tokens approximate\n    overlap: int = 50,\n    metadata: dict = None\n) -> Generator[dict, None, None]:\n    \"\"\"Chunk with overlap for context continuity.\"\"\"\n    words = text.split()\n    metadata = metadata or {}\n\n    i = 0\n    chunk_index = 0\n    while i < len(words):\n        end = min(i + chunk_size, len(words))\n        chunk_text = \" \".join(words[i:end])\n\n        yield {\n            \"content\": chunk_text,\n            \"chunk_index\": chunk_index,\n            \"word_start\": i,\n            \"word_end\": end,\n            **metadata,\n        }\n\n        chunk_index += 1\n        i = end - overlap  # move back by overlap amount\n\nchunks = list(chunk_document(document_text, chunk_size=500, overlap=50))\n```\n\n## Voyage AI Embeddings\n\n```python\nimport voyageai\n\nclient = voyageai.Client(api_key=os.environ[\"VOYAGE_API_KEY\"])\n\ndef embed_texts(texts: list[str], input_type: str = \"document\") -> list[list[float]]:\n    \"\"\"\n    input_type:\n    - \"document\" for content being stored\n    - \"query\" for search queries\n    \"\"\"\n    result = client.embed(\n        texts,\n        model=\"voyage-3\",\n        input_type=input_type,\n    )\n    return result.embeddings\n\n# Embed query differently from documents\nquery_embedding = embed_texts([\"What is the refund policy?\"], input_type=\"query\")[0]\ndoc_embeddings = embed_texts([chunk[\"content\"] for chunk in chunks], input_type=\"document\")\n```\n\n## Supabase pgvector Setup\n\n```sql\n-- Enable pgvector extension\nCREATE EXTENSION IF NOT EXISTS vector;\n\n-- Documents table\nCREATE TABLE documents (\n  id BIGSERIAL PRIMARY KEY,\n  source TEXT NOT NULL,\n  chunk_index INTEGER NOT NULL,\n  content TEXT NOT NULL,\n  embedding vector(1024),  -- match Voyage-3 dimensions\n  metadata JSONB DEFAULT '{}',\n  created_at TIMESTAMPTZ DEFAULT NOW()\n);\n\n-- Vector similarity index (HNSW for fast approximate search)\nCREATE INDEX documents_embedding_idx\n  ON documents USING hnsw (embedding vector_cosine_ops)\n  WITH (m = 16, ef_construction = 64);\n\n-- Full-text search index for BM25\nCREATE INDEX documents_content_fts_idx\n  ON documents USING gin(to_tsvector('english', content));\n\n-- RLS\nALTER TABLE documents ENABLE ROW LEVEL SECURITY;\n```\n\n## Storing Embeddings\n\n```python\nimport psycopg2\nimport json\n\ndef store_chunks(chunks: list[dict], embeddings: list[list[float]]):\n    conn = psycopg2.connect(os.environ[\"DATABASE_URL\"])\n\n    with conn.cursor() as cur:\n        for chunk, embedding in zip(chunks, embeddings):\n            cur.execute(\n                \"\"\"\n                INSERT INTO documents (source, chunk_index, content, embedding, metadata)\n                VALUES (%s, %s, %s, %s::vector, %s)\n                \"\"\",\n                (\n                    chunk[\"source\"],\n                    chunk[\"chunk_index\"],\n                    chunk[\"content\"],\n                    embedding,\n                    json.dumps(chunk.get(\"metadata\", {})),\n                )\n            )\n\n    conn.commit()\n    conn.close()\n```\n\n## Semantic Search\n\n```python\ndef semantic_search(query: str, limit: int = 5) -> list[dict]:\n    query_embedding = embed_texts([query], input_type=\"query\")[0]\n\n    conn = psycopg2.connect(os.environ[\"DATABASE_URL\"])\n    with conn.cursor() as cur:\n        cur.execute(\n            \"\"\"\n            SELECT content, source, chunk_index, metadata,\n                   1 - (embedding <=> %s::vector) AS similarity\n            FROM documents\n            ORDER BY embedding <=> %s::vector\n            LIMIT %s\n            \"\"\",\n            (query_embedding, query_embedding, limit)\n        )\n        rows = cur.fetchall()\n\n    conn.close()\n    return [\n        {\"content\": r[0], \"source\": r[1], \"chunk_index\": r[2], \"metadata\": r[3], \"score\": r[4]}\n        for r in rows\n    ]\n```\n\n## Hybrid Search (Semantic + BM25 with RRF)\n\n```python\ndef hybrid_search(query: str, limit: int = 5, semantic_weight: float = 0.7) -> list[dict]:\n    \"\"\"Combine semantic and keyword search with Reciprocal Rank Fusion.\"\"\"\n    query_embedding = embed_texts([query], input_type=\"query\")[0]\n\n    conn = psycopg2.connect(os.environ[\"DATABASE_URL\"])\n    with conn.cursor() as cur:\n        cur.execute(\n            \"\"\"\n            WITH semantic AS (\n                SELECT id, content, source,\n                       ROW_NUMBER() OVER (ORDER BY embedding <=> %s::vector) AS rank\n                FROM documents\n                ORDER BY embedding <=> %s::vector\n                LIMIT 20\n            ),\n            keyword AS (\n                SELECT id, content, source,\n                       ROW_NUMBER() OVER (ORDER BY ts_rank(to_tsvector('english', content), plainto_tsquery('english', %s)) DESC) AS rank\n                FROM documents\n                WHERE to_tsvector('english', content) @@ plainto_tsquery('english', %s)\n                LIMIT 20\n            )\n            SELECT\n                COALESCE(s.id, k.id) AS id,\n                COALESCE(s.content, k.content) AS content,\n                COALESCE(s.source, k.source) AS source,\n                (COALESCE(1.0/(60 + s.rank), 0) * %s + COALESCE(1.0/(60 + k.rank), 0) * %s) AS rrf_score\n            FROM semantic s\n            FULL OUTER JOIN keyword k ON s.id = k.id\n            ORDER BY rrf_score DESC\n            LIMIT %s\n            \"\"\",\n            (query_embedding, query_embedding, query, query, semantic_weight, 1 - semantic_weight, limit)\n        )\n        rows = cur.fetchall()\n\n    conn.close()\n    return [{\"id\": r[0], \"content\": r[1], \"source\": r[2], \"score\": r[3]} for r in rows]\n```\n\n## Augmented Generation\n\n```python\nimport anthropic\n\ndef rag_query(question: str) -> str:\n    # Retrieve\n    chunks = hybrid_search(question, limit=5)\n\n    # Format context\n    context = \"\\n\\n---\\n\\n\".join(\n        f\"[Source: {c['source']}]\\n{c['content']}\"\n        for c in chunks\n    )\n\n    # Generate\n    client = anthropic.Anthropic()\n    response = client.messages.create(\n        model=\"claude-opus-4-6\",\n        max_tokens=2_000,\n        system=f\"\"\"Answer questions using ONLY the provided context. \nIf the answer isn't in the context, say so clearly.\nCite sources with [Source: filename].\n\nContext:\n{context}\"\"\",\n        messages=[{\"role\": \"user\", \"content\": question}],\n    )\n\n    return response.content[0].text\n```\n\n## Embedding Update Strategy\n\nWhen content changes, re-embed only the changed chunks:\n\n```python\nimport hashlib\n\ndef content_hash(text: str) -> str:\n    return hashlib.sha256(text.encode()).hexdigest()[:16]\n\n# Store hash with each chunk\n# On update: compare hash → only re-embed changed chunks\n# This keeps embedding costs minimal on incremental ingestion\n```\n","html":"<h2>Overview</h2>\n<p>RAG (Retrieval-Augmented Generation) grounds LLM responses in your documents. Pipeline: chunk → embed → store → search → augment → generate.</p>\n<h2>Embedding Model Selection</h2>\n<p>| Model | Provider | Dimensions | Use Case |\n|-------|---------|-----------|----------|\n| <code>voyage-3</code> | Voyage AI | 1024 | General documents (best quality) |\n| <code>voyage-3-lite</code> | Voyage AI | 512 | High volume, cost-sensitive |\n| <code>text-embedding-3-small</code> | OpenAI | 1536 | Fallback, widely supported |\n| <code>text-embedding-3-large</code> | OpenAI | 3072 | Maximum quality (high cost) |</p>\n<p>Voyage AI: 200M free tokens/month. Preferred over OpenAI for cost and quality balance.</p>\n<h2>Document Chunking</h2>\n<pre><code class=\"language-python\">from typing import Generator\n\ndef chunk_document(\n    text: str,\n    chunk_size: int = 500,     # tokens approximate\n    overlap: int = 50,\n    metadata: dict = None\n) -> Generator[dict, None, None]:\n    \"\"\"Chunk with overlap for context continuity.\"\"\"\n    words = text.split()\n    metadata = metadata or {}\n\n    i = 0\n    chunk_index = 0\n    while i &#x3C; len(words):\n        end = min(i + chunk_size, len(words))\n        chunk_text = \" \".join(words[i:end])\n\n        yield {\n            \"content\": chunk_text,\n            \"chunk_index\": chunk_index,\n            \"word_start\": i,\n            \"word_end\": end,\n            **metadata,\n        }\n\n        chunk_index += 1\n        i = end - overlap  # move back by overlap amount\n\nchunks = list(chunk_document(document_text, chunk_size=500, overlap=50))\n</code></pre>\n<h2>Voyage AI Embeddings</h2>\n<pre><code class=\"language-python\">import voyageai\n\nclient = voyageai.Client(api_key=os.environ[\"VOYAGE_API_KEY\"])\n\ndef embed_texts(texts: list[str], input_type: str = \"document\") -> list[list[float]]:\n    \"\"\"\n    input_type:\n    - \"document\" for content being stored\n    - \"query\" for search queries\n    \"\"\"\n    result = client.embed(\n        texts,\n        model=\"voyage-3\",\n        input_type=input_type,\n    )\n    return result.embeddings\n\n# Embed query differently from documents\nquery_embedding = embed_texts([\"What is the refund policy?\"], input_type=\"query\")[0]\ndoc_embeddings = embed_texts([chunk[\"content\"] for chunk in chunks], input_type=\"document\")\n</code></pre>\n<h2>Supabase pgvector Setup</h2>\n<pre><code class=\"language-sql\">-- Enable pgvector extension\nCREATE EXTENSION IF NOT EXISTS vector;\n\n-- Documents table\nCREATE TABLE documents (\n  id BIGSERIAL PRIMARY KEY,\n  source TEXT NOT NULL,\n  chunk_index INTEGER NOT NULL,\n  content TEXT NOT NULL,\n  embedding vector(1024),  -- match Voyage-3 dimensions\n  metadata JSONB DEFAULT '{}',\n  created_at TIMESTAMPTZ DEFAULT NOW()\n);\n\n-- Vector similarity index (HNSW for fast approximate search)\nCREATE INDEX documents_embedding_idx\n  ON documents USING hnsw (embedding vector_cosine_ops)\n  WITH (m = 16, ef_construction = 64);\n\n-- Full-text search index for BM25\nCREATE INDEX documents_content_fts_idx\n  ON documents USING gin(to_tsvector('english', content));\n\n-- RLS\nALTER TABLE documents ENABLE ROW LEVEL SECURITY;\n</code></pre>\n<h2>Storing Embeddings</h2>\n<pre><code class=\"language-python\">import psycopg2\nimport json\n\ndef store_chunks(chunks: list[dict], embeddings: list[list[float]]):\n    conn = psycopg2.connect(os.environ[\"DATABASE_URL\"])\n\n    with conn.cursor() as cur:\n        for chunk, embedding in zip(chunks, embeddings):\n            cur.execute(\n                \"\"\"\n                INSERT INTO documents (source, chunk_index, content, embedding, metadata)\n                VALUES (%s, %s, %s, %s::vector, %s)\n                \"\"\",\n                (\n                    chunk[\"source\"],\n                    chunk[\"chunk_index\"],\n                    chunk[\"content\"],\n                    embedding,\n                    json.dumps(chunk.get(\"metadata\", {})),\n                )\n            )\n\n    conn.commit()\n    conn.close()\n</code></pre>\n<h2>Semantic Search</h2>\n<pre><code class=\"language-python\">def semantic_search(query: str, limit: int = 5) -> list[dict]:\n    query_embedding = embed_texts([query], input_type=\"query\")[0]\n\n    conn = psycopg2.connect(os.environ[\"DATABASE_URL\"])\n    with conn.cursor() as cur:\n        cur.execute(\n            \"\"\"\n            SELECT content, source, chunk_index, metadata,\n                   1 - (embedding &#x3C;=> %s::vector) AS similarity\n            FROM documents\n            ORDER BY embedding &#x3C;=> %s::vector\n            LIMIT %s\n            \"\"\",\n            (query_embedding, query_embedding, limit)\n        )\n        rows = cur.fetchall()\n\n    conn.close()\n    return [\n        {\"content\": r[0], \"source\": r[1], \"chunk_index\": r[2], \"metadata\": r[3], \"score\": r[4]}\n        for r in rows\n    ]\n</code></pre>\n<h2>Hybrid Search (Semantic + BM25 with RRF)</h2>\n<pre><code class=\"language-python\">def hybrid_search(query: str, limit: int = 5, semantic_weight: float = 0.7) -> list[dict]:\n    \"\"\"Combine semantic and keyword search with Reciprocal Rank Fusion.\"\"\"\n    query_embedding = embed_texts([query], input_type=\"query\")[0]\n\n    conn = psycopg2.connect(os.environ[\"DATABASE_URL\"])\n    with conn.cursor() as cur:\n        cur.execute(\n            \"\"\"\n            WITH semantic AS (\n                SELECT id, content, source,\n                       ROW_NUMBER() OVER (ORDER BY embedding &#x3C;=> %s::vector) AS rank\n                FROM documents\n                ORDER BY embedding &#x3C;=> %s::vector\n                LIMIT 20\n            ),\n            keyword AS (\n                SELECT id, content, source,\n                       ROW_NUMBER() OVER (ORDER BY ts_rank(to_tsvector('english', content), plainto_tsquery('english', %s)) DESC) AS rank\n                FROM documents\n                WHERE to_tsvector('english', content) @@ plainto_tsquery('english', %s)\n                LIMIT 20\n            )\n            SELECT\n                COALESCE(s.id, k.id) AS id,\n                COALESCE(s.content, k.content) AS content,\n                COALESCE(s.source, k.source) AS source,\n                (COALESCE(1.0/(60 + s.rank), 0) * %s + COALESCE(1.0/(60 + k.rank), 0) * %s) AS rrf_score\n            FROM semantic s\n            FULL OUTER JOIN keyword k ON s.id = k.id\n            ORDER BY rrf_score DESC\n            LIMIT %s\n            \"\"\",\n            (query_embedding, query_embedding, query, query, semantic_weight, 1 - semantic_weight, limit)\n        )\n        rows = cur.fetchall()\n\n    conn.close()\n    return [{\"id\": r[0], \"content\": r[1], \"source\": r[2], \"score\": r[3]} for r in rows]\n</code></pre>\n<h2>Augmented Generation</h2>\n<pre><code class=\"language-python\">import anthropic\n\ndef rag_query(question: str) -> str:\n    # Retrieve\n    chunks = hybrid_search(question, limit=5)\n\n    # Format context\n    context = \"\\n\\n---\\n\\n\".join(\n        f\"[Source: {c['source']}]\\n{c['content']}\"\n        for c in chunks\n    )\n\n    # Generate\n    client = anthropic.Anthropic()\n    response = client.messages.create(\n        model=\"claude-opus-4-6\",\n        max_tokens=2_000,\n        system=f\"\"\"Answer questions using ONLY the provided context. \nIf the answer isn't in the context, say so clearly.\nCite sources with [Source: filename].\n\nContext:\n{context}\"\"\",\n        messages=[{\"role\": \"user\", \"content\": question}],\n    )\n\n    return response.content[0].text\n</code></pre>\n<h2>Embedding Update Strategy</h2>\n<p>When content changes, re-embed only the changed chunks:</p>\n<pre><code class=\"language-python\">import hashlib\n\ndef content_hash(text: str) -> str:\n    return hashlib.sha256(text.encode()).hexdigest()[:16]\n\n# Store hash with each chunk\n# On update: compare hash → only re-embed changed chunks\n# This keeps embedding costs minimal on incremental ingestion\n</code></pre>\n"}