AI Engineering Academy · Lektion

Implementera centrala RAG- och agentfunktioner

Bygg pipelinen för dokumentimport, indexering i vektorlagret, hämtning med omrankning och integrationer med agentverktyg enligt de mönster ni har lärt er under kursen.

Lektion 2 av 413 steg

Implementera centrala RAG- och agentfunktioner är en gratis lektion i AI Engineering Academy på CoddyKit. Detta är lektion 2 av 4. Ni kan läsa hela lektionen gratis nedan och sedan öva praktiskt i webbläsaren med en inbyggd kodredigerare och en AI-handledare som är tillgänglig dygnet runt. Den ingår i lärvägen för AI Engineering Academy, och Era framsteg synkroniseras mellan webben och CoddyKit-appen. Kursen i AI Engineering Academy innehåller totalt 4 lektioner.

Implementationsordning och beroenden

Bygg systemet nerifrån och upp: börja med komponenterna som saknar externa beroenden och lägg sedan till komponenter som är beroende av dem. För ett RAG- och agentsystem är ordningen: (1) schema för vektorlagring, (2) inläsningspipeline, (3) retriever, (4) grundläggande Q&A-kedja, (5) streaming-endpoint, (6) agent med verktyg, (7) cachinglager, (8) spårningsinstrumentering. Testa varje komponent isolerat innan du integrerar den i pipelinen.

# Build order:
IMPL_ORDER = [
    'pgvector_schema',      # prerequisite for everything
    'document_ingestion',   # populate the vector store
    'hybrid_retriever',     # test retrieval in isolation
    'qa_chain_basic',       # integrate LLM with retrieval
    'streaming_endpoint',   # expose via API
    'function_calling',     # add agent tool calls
    'semantic_cache',       # reduce repeat API calls
    'langsmith_tracing',    # add after core works
    'injection_filter',     # harden before load testing
]

Konfigurera vektorlagringen

Skapa pgvector-tabellen med rätt schema innan du läser in dokument. Ta med kolumner för vektorn, texten i varje chunk, alla metadatafält och en tidsstämpel för updated_at som möjliggör selektiv omindexering. Skapa omedelbart ett vektorindex (HNSW eller IVFFlat) på embedding-kolumnen — att lägga till indexet efter miljontals rader är mycket långsammare än att skapa det direkt i en tom tabell.

-- PostgreSQL schema with pgvector
CREATE EXTENSION IF NOT EXISTS vector;

CREATE TABLE document_chunks (
    id          BIGSERIAL PRIMARY KEY,
    doc_id      TEXT NOT NULL,
    chunk_text  TEXT NOT NULL,
    embedding   VECTOR(1536) NOT NULL,
    source_file TEXT,
    page_number INT,
    section     TEXT,
    doc_type    TEXT,
    tenant_id   TEXT NOT NULL,  -- for data isolation
    created_at  TIMESTAMPTZ DEFAULT NOW()
);

CREATE INDEX idx_chunks_hnsw ON document_chunks
USING hnsw (embedding vector_cosine_ops)
WITH (m = 16, ef_construction = 64);

CREATE INDEX idx_chunks_tenant ON document_chunks(tenant_id);

Bygga inläsningspipelinen

Implementera inläsning som en enda asynkron funktion som tar emot en filsökväg eller URL och returnerar antalet indexerade chunks. Använd LangChains dokumentladdare för olika filtyper och den rekursiva textdelaren för tecken för chunking. Bunta ihop embedding-anrop så att du håller dig inom indatagränsen på 2048 tokens och minskar antalet API-anrop från tusentals till dussintals.

from langchain_community.document_loaders import PyPDFLoader
from langchain.text_splitter import RecursiveCharacterTextSplitter
from openai import AsyncOpenAI

async def ingest_document(file_path: str, doc_id: str, tenant_id: str) -> int:
    loader = PyPDFLoader(file_path)
    pages = loader.load()
    splitter = RecursiveCharacterTextSplitter(chunk_size=800, chunk_overlap=100)
    chunks = splitter.split_documents(pages)
    # Batch embed
    texts = [c.page_content for c in chunks]
    client = AsyncOpenAI()
    embeddings_response = await client.embeddings.create(
        model='text-embedding-3-small',
        input=texts
    )
    embeddings = [e.embedding for e in embeddings_response.data]
    await insert_chunks_to_pgvector(chunks, embeddings, doc_id, tenant_id)
    return len(chunks)

Bygga den hybrida retrievern

Kombinera tät vektorsökning med BM25-nyckelordssökning och slå ihop resultaten med reciprocal rank fusion. Implementera retrievern som en klass med en enda metod, retrieve(query, tenant_id, top_k). Kör båda sökningarna samtidigt med asyncio.gather, slå ihop de rangordnade listorna med RRF, ta bort dubbletter utifrån chunk-ID och returnera de bästa top_k-resultaten med deras källmetadata.

import asyncio
from rank_bm25 import BM25Okapi

class HybridRetriever:
    def __init__(self, pool, k_rrf: int = 60):
        self.pool = pool
        self.k_rrf = k_rrf

    async def retrieve(self, query: str, tenant_id: str, top_k: int = 10) -> list:
        dense_results, sparse_results = await asyncio.gather(
            self._dense_search(query, tenant_id, top_k * 3),
            self._bm25_search(query, tenant_id, top_k * 3)
        )
        merged = self._rrf_merge(dense_results, sparse_results)
        return merged[:top_k]

    def _rrf_score(self, rank: int) -> float:
        return 1.0 / (self.k_rrf + rank + 1)

Implementera den centrala QA-kedjan

Bygg den centrala Q&A-kedjan med LangChain LCEL. Kedjan tar emot en fråga och hämtade chunks, formaterar en utökad prompt med instruktioner om att endast använda den angivna kontexten och streamar svaret. Lägg till tydliga instruktioner om att modellen ska ange källor med dokumentnamn och sidnummer samt säga ”I don't know” när svaret inte finns i den hämtade kontexten.

from langchain.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
from langchain.schema.output_parser import StrOutputParser

RAG_PROMPT = ChatPromptTemplate.from_messages([
    ('system', 'You are a precise assistant. Answer ONLY using the provided context. '
               'Cite sources as [DocName, p.N]. If the answer is not in the context, say "I don\'t have information about that."'),
    ('user', 'Context:\n{context}\n\nQuestion: {question}')
])

llm = ChatOpenAI(model='gpt-4o', temperature=0, streaming=True)
qa_chain = RAG_PROMPT | llm | StrOutputParser()

async def answer_question(question: str, chunks: list) -> str:
    context = '\n\n'.join(f'[{c["source"]}]\n{c["text"]}' for c in chunks)
    return await qa_chain.ainvoke({'context': context, 'question': question})

Lägga till Function Calling i agenten

Utöka det grundläggande Q&A-systemet med function calling för att hantera frågor som kräver realtidsdata eller beräkningar. Definiera verktyg för att: söka på webben efter aktuell information, köra SQL-frågor mot en säker skrivskyddad databas och slå upp specifika poster med ID. Agenten avgör vilka verktyg som ska anropas utifrån frågan, kör dem och väver in resultaten i det slutliga svaret.

from openai import AsyncOpenAI
import json

TOOLS = [
    {
        'type': 'function',
        'function': {
            'name': 'search_knowledge_base',
            'description': 'Search the internal document knowledge base for relevant information',
            'parameters': {
                'type': 'object',
                'properties': {
                    'query': {'type': 'string', 'description': 'Search query'},
                    'top_k': {'type': 'integer', 'default': 5}
                },
                'required': ['query']
            }
        }
    }
]

async def agent_with_tools(question: str, tenant_id: str) -> str:
    client = AsyncOpenAI()
    messages = [{'role': 'user', 'content': question}]
    while True:
        resp = await client.chat.completions.create(
            model='gpt-4o', messages=messages, tools=TOOLS)
        if resp.choices[0].finish_reason != 'tool_calls':
            return resp.choices[0].message.content
        tool_call = resp.choices[0].message.tool_calls[0]
        args = json.loads(tool_call.function.arguments)
        result = await dispatch_tool(tool_call.function.name, args, tenant_id)
        messages.append({'role': 'tool', 'tool_call_id': tool_call.id, 'content': result})

Streama agentens svar

Omslut agenten i en FastAPI StreamingResponse med server-sända händelser så att frontend-gränssnittet visar token när de anländer. För agentsvar med verktygsanrop strömmar ni en förloppsindikator medan verktyget körs (”Searching knowledge base...”) och strömmar sedan det slutliga svaret token för token. Det förhindrar att sidan verkar ha hängt sig medan verktyget körs.

from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import asyncio

app = FastAPI()

async def event_stream(question: str, tenant_id: str):
    yield 'data: {"type": "start"}\n\n'
    chunks = await retriever.retrieve(question, tenant_id)
    yield 'data: {"type": "retrieving", "count": ' + str(len(chunks)) + '}\n\n'
    async for token in qa_chain.astream({'context': format_context(chunks), 'question': question}):
        yield f'data: {{"type": "token", "content": {repr(token)}}}\n\n'
    yield 'data: {"type": "done"}\n\n'

@app.post('/query/stream')
async def stream_query(question: str, tenant_id: str):
    return StreamingResponse(event_stream(question, tenant_id), media_type='text/event-stream')

Koppla in den semantiska cachen

Lägg till den semantiska cachen som ett steg före informationshämtningen i frågepipelinen. Innan ni anropar hämtaren och LLM:en bäddar ni in användarens fråga och kontrollerar cachen. Om cachen innehåller en träff med likhet över tröskelvärdet returnerar ni det cachade svaret direkt med flaggan cached: true. Vid cachemiss fortsätter frågan genom hela pipelinen, och det resulterande svaret lagras i cachen för framtida liknande frågor.

async def query_pipeline(question: str, tenant_id: str) -> dict:
    # 1. Check semantic cache
    cache_hit = await semantic_cache.lookup(question, tenant_id, threshold=0.92)
    if cache_hit:
        return {'answer': cache_hit.answer, 'cached': True, 'sources': cache_hit.sources}

    # 2. Retrieve
    chunks = await retriever.retrieve(question, tenant_id, top_k=5)
    chunks = await reranker.rerank(question, chunks, top_n=3)

    # 3. Generate
    answer = await answer_question(question, chunks)
    sources = [c['source'] for c in chunks]

    # 4. Cache result
    await semantic_cache.store(question, tenant_id, answer, sources)

    return {'answer': answer, 'cached': False, 'sources': sources}

Lägg till spårning med LangSmith

Instrumentera pipelinen med LangSmith genom att ange två miljövariabler. Varje anrop till en LangChain-kedja spåras automatiskt med antal token, fördröjning, indata, utdata och eventuella fel. För anpassad kod som inte använder LangChain (informationshämtning, omrankning) omsluter ni anropen med dekoratorer av typen @traceable så att de inkluderas i spårningen. Det ger fullständig insyn i varje steg i pipelinen från början till slut.

import os
from langsmith import traceable

os.environ['LANGCHAIN_TRACING_V2'] = 'true'
os.environ['LANGCHAIN_API_KEY'] = os.environ['LANGSMITH_API_KEY']
os.environ['LANGCHAIN_PROJECT'] = 'document-qa-production'

# Wrap non-LangChain steps with @traceable
@traceable(name='hybrid_retrieval')
async def traced_retrieval(question: str, tenant_id: str, top_k: int) -> list:
    return await retriever.retrieve(question, tenant_id, top_k)

@traceable(name='cohere_reranking')
async def traced_reranking(question: str, chunks: list) -> list:
    return await reranker.rerank(question, chunks)

# LangChain LCEL chains are automatically traced — no extra code needed

Kör integrationstester

Skriv integrationstester som kör hela frågepipelinen från indata i form av en fråga till det slutliga svaret. Använd en liten test-vektorlagerlösning med kända dokument så att ni kan skriva deterministiska påståenden om vilka textavsnitt som ska hämtas och vad svaret ska innehålla. Kör integrationstester mot en staging-miljö som motsvarar produktionsinfrastrukturen men använder ett separat vektorlager och en separat LLM-nyckel.

import pytest

@pytest.mark.asyncio
async def test_full_pipeline_returns_grounded_answer():
    # Setup: ingest known document
    await ingest_document('tests/fixtures/policy.pdf', 'policy_v1', 'test_tenant')

    # Query with a question that has a known answer in the document
    result = await query_pipeline(
        question='What is the cancellation policy?',
        tenant_id='test_tenant'
    )

    assert result['answer'] is not None
    assert len(result['answer']) > 50
    assert '24 hours' in result['answer']  # known fact in document
    assert 'policy.pdf' in str(result['sources'])
    assert result['cached'] is False  # fresh query

Mät den grundläggande kvaliteten på informationshämtningen

Innan ni optimerar något bör ni fastställa en baslinje för informationshämtning. Använd ert utvärderingstestset för att mäta träffrekvensen (visas rätt dokument bland de fem främsta resultaten?) och MRR (hur högt rankas det?). Kör denna baslinje efter att det initiala vektorlagret har konfigurerats, men innan ni lägger till omrankning eller hybridsökning. Baslinjen visar vilka faktiska förbättringar varje optimering ger, så att ni har underlag för vilka tekniker som är värda att behålla.

async def measure_retrieval_baseline(test_cases: list) -> dict:
    hits = 0
    reciprocal_ranks = []
    for case in test_cases:
        results = await retriever.retrieve(case['question'], case['tenant_id'], top_k=5)
        result_docs = [r['doc_id'] for r in results]
        if case['relevant_doc'] in result_docs:
            hits += 1
            rank = result_docs.index(case['relevant_doc']) + 1
            reciprocal_ranks.append(1.0 / rank)
        else:
            reciprocal_ranks.append(0.0)
    return {
        'hit_rate_at_5': hits / len(test_cases),
        'mrr': sum(reciprocal_ranks) / len(reciprocal_ranks)
    }

Snabbkontroll

Testa er förståelse av hur man bygger grundläggande RAG- och agentfunktioner.

Lektionssammanfattning

I den här lektionen har ni lärt er att implementering i bottom-up-ordning säkerställer att varje komponent testas isolerat före integreringen, att hybridinformationshämtning med samtidig dense- och BM25-sökning som kombineras genom RRF ger bäst kvalitet på informationshämtningen, samt att LangSmith-spårning med @traceable-dekoratorer ger fullständig insyn i hela pipelinen. Härnäst gör vi systemet robust med mönster för säkerhet, cachning och tillförlitlighet.

Gratis att börja

Lär dig Python med en AI-lärare – gratis

Skriv och kör riktig kod i webbläsaren, få omedelbar hjälp av en AI-lärare dygnet runt och fortsätt där du slutade – på webben eller i appen.

Kurser
30
Lektioner
120

Vanliga frågor

Är lektionen ”Implementera centrala RAG- och agentfunktioner” gratis?

Ja – hela texten till ”Implementera centrala RAG- och agentfunktioner” kan läsas gratis här på webben. Om Ni vill öva interaktivt med en inbyggd kodredigerare och en AI-handledare som är tillgänglig dygnet runt och låsa upp resten av kursen i AI Engineering Academy, kan Ni uppgradera till CoddyKit PRO. Kursen i AI Engineering Academy innehåller totalt 4 lektioner.

Vad lär jag mig i ”Implementera centrala RAG- och agentfunktioner”?

Bygg pipelinen för dokumentimport, indexering i vektorlagret, hämtning med omrankning och integrationer med agentverktyg enligt de mönster ni har lärt er under kursen. Ni övar på AI Engineering Academy med praktisk kod som körs direkt i webbläsaren, medan en AI-handledare som är tillgänglig dygnet runt svarar på Era frågor under lektionen.

Behöver jag någon erfarenhet för att börja lära mig AI Engineering Academy?

Du behöver inga förkunskaper. Utbildningen i AI Engineering Academy på CoddyKit är upplagd för allt från nybörjare till avancerade elever, så att du kan börja här eller från början och gå fram i din egen takt. Detta är lektion 2 av 4.

Hur lång tid tar lektionen ”Implementera centrala RAG- och agentfunktioner”?

De flesta CoddyKit-lektioner tar cirka 5–10 minuter. Varje lektion är kort och interaktiv, så att du gör stadiga framsteg och kan fortsätta precis där du slutade – på webben eller i appen.

Kan jag skriva och köra kod i den här AI Engineering Academy-lektionen?

Ja. Varje AI Engineering Academy-lektion innehåller en inbyggd kodredigerare, så att du kan skriva och köra riktig kod direkt i webbläsaren och få omedelbar AI-feedback – utan lokal installation.

Alla lektioner i den här kursen

  1. Utforma produktionsarkitekturen
  2. Implementera centrala RAG- och agentfunktioner
  3. Härda systemet: säkerhet, cachning och tillförlitlighet
  4. Utvärdering, driftsättning och efteranalys
← Tillbaka till AI Engineering Academy