Building RAG Applications with Vertex AI and Cloud Storage

Retrieval-Augmented Generation (RAG) makes large language models more accurate by grounding responses in your own documents. This guide builds a production RAG pipeline on GCP: document ingestion from Cloud Storage, embedding generation with Vertex AI, vector search, and answer synthesis with Gemini.

Large language models have a knowledge cutoff and no access to your organization's private documents. Ask Gemini what's in your Q3 board presentation or your company's security policy, and it will hallucinate an answer with complete confidence.

Retrieval-Augmented Generation (RAG) solves this by connecting the LLM to a retrieval system. Instead of asking the model to remember information from training, you retrieve relevant document chunks from your own knowledge base and include them in the context window of each query. The model synthesizes an answer from the retrieved content, with citations you can verify.

This architecture works. It's used in production for internal knowledge assistants, customer support bots, compliance Q&A tools, and document analysis systems. The building blocks on GCP are Vertex AI for embedding models and LLM inference, Cloud Storage for document storage, and Vertex AI Vector Search for fast approximate nearest neighbor retrieval.

RAG Architecture Overview

INGESTION (offline):
Documents in Cloud Storage
    → Extract text (PDF parsing, HTML cleaning)
    → Chunk into overlapping segments (~512 tokens)
    → Generate embeddings with text-embedding-005
    → Store in Vertex AI Vector Search Index

RETRIEVAL (online, per query):
User query
    → Generate query embedding
    → Search Vector Search for top-k chunks
    → Return chunks + metadata

GENERATION (online, per query):
Query + Retrieved chunks
    → Build prompt with context
    → Call Gemini for answer synthesis
    → Return answer with source citations

Phase 1: Document Ingestion Pipeline

Setting Up Cloud Storage

gsutil mb -l europe-west4 gs://my-project-documents
gsutil mb -l europe-west4 gs://my-project-embeddings
gsutil cp -r ./company-docs/* gs://my-project-documents/

Document Processing and Chunking

# ingestion/processor.py
import re
from dataclasses import dataclass
from typing import Generator
import fitz  # PyMuPDF for PDF parsing
from google.cloud import storage

@dataclass
class DocumentChunk:
    chunk_id: str
    document_id: str
    document_name: str
    text: str
    page_number: int
    char_start: int
    char_end: int

def extract_text_from_pdf(bucket_name: str, blob_name: str) -> list[dict]:
    """Extract text with page numbers from a GCS PDF."""
    client = storage.Client()
    bucket = client.bucket(bucket_name)
    blob = bucket.blob(blob_name)
    pdf_bytes = blob.download_as_bytes()
    pdf = fitz.open(stream=pdf_bytes, filetype="pdf")

    pages = []
    for page_num, page in enumerate(pdf):
        text = re.sub(r's+', ' ', page.get_text("text")).strip()
        if len(text) > 50:
            pages.append({"page": page_num + 1, "text": text})
    return pages

def chunk_document(
    pages: list[dict],
    document_id: str,
    document_name: str,
    chunk_size: int = 512,
    overlap: int = 50,
) -> Generator[DocumentChunk, None, None]:
    """Split document pages into overlapping chunks."""
    full_text = " ".join(p["text"] for p in pages)
    words = full_text.split()

    for idx, i in enumerate(range(0, len(words), chunk_size - overlap)):
        chunk_words = words[i:i + chunk_size]
        if len(chunk_words) < 20:
            continue
        yield DocumentChunk(
            chunk_id=f"{document_id}_chunk_{idx}",
            document_id=document_id,
            document_name=document_name,
            text=" ".join(chunk_words),
            page_number=1,
            char_start=i,
            char_end=i + len(chunk_words),
        )

Generating Embeddings with Vertex AI

# ingestion/embedder.py
import time
from typing import List
import vertexai
from vertexai.language_models import TextEmbeddingModel

vertexai.init(project="my-project", location="europe-west4")
EMBEDDING_MODEL = TextEmbeddingModel.from_pretrained("text-embedding-005")
EMBEDDING_DIMENSION = 768

def embed_texts(texts: List[str], task_type: str = "RETRIEVAL_DOCUMENT") -> List[List[float]]:
    embeddings = []
    batch_size = 250  # API limit

    for i in range(0, len(texts), batch_size):
        batch = texts[i:i + batch_size]
        response = EMBEDDING_MODEL.get_embeddings(
            batch,
            output_dimensionality=EMBEDDING_DIMENSION,
            task_type=task_type,
        )
        embeddings.extend([e.values for e in response])
        if i + batch_size < len(texts):
            time.sleep(0.5)

    return embeddings

Storing Embeddings for Vector Search

import json
from google.cloud import storage
from embedder import embed_texts
from processor import DocumentChunk

def prepare_and_upload_index_data(chunks: list[DocumentChunk], bucket_name: str, prefix: str):
    chunk_texts = [c.text for c in chunks]
    embeddings = embed_texts(chunk_texts)

    storage_client = storage.Client()
    bucket = storage_client.bucket(bucket_name)

    # Embedding records (for Vector Search index)
    embedding_lines = [
        json.dumps({"id": chunk.chunk_id, "embedding": emb})
        for chunk, emb in zip(chunks, embeddings)
    ]
    bucket.blob(f"{prefix}/embeddings.json").upload_from_string(
        "
".join(embedding_lines)
    )

    # Metadata (for retrieval context)
    metadata = {c.chunk_id: {
        "text": c.text,
        "document_id": c.document_id,
        "document_name": c.document_name,
        "page_number": c.page_number,
    } for c in chunks}
    bucket.blob(f"{prefix}/metadata.json").upload_from_string(
        json.dumps(metadata)
    )

Building the Vector Search Index

from google.cloud import aiplatform

aiplatform.init(project="my-project", location="europe-west4")

index = aiplatform.MatchingEngineIndex.create_tree_ah_index(
    display_name="company-knowledge-base",
    contents_delta_uri="gs://my-project-embeddings/index-data/",
    dimensions=768,
    approximate_neighbors_count=50,
    distance_measure_type="DOT_PRODUCT_DISTANCE",
    leaf_node_embedding_count=500,
    leaf_nodes_to_search_percent=80,
)

index_endpoint = aiplatform.MatchingEngineIndexEndpoint.create(
    display_name="knowledge-base-endpoint",
    public_endpoint_enabled=True,
)

index_endpoint.deploy_index(
    index=index,
    deployed_index_id="knowledge_base_v1",
    display_name="knowledge-base-v1",
    min_replica_count=2,
    max_replica_count=10,
)

Phase 2: Retrieval Service

# retrieval/retriever.py
import json
from google.cloud import aiplatform, storage
from vertexai.language_models import TextEmbeddingModel

class KnowledgeRetriever:
    def __init__(self, project_id: str, index_endpoint_id: str, deployed_index_id: str, metadata_bucket: str):
        self.embedding_model = TextEmbeddingModel.from_pretrained("text-embedding-005")
        self.index_endpoint = aiplatform.MatchingEngineIndexEndpoint(index_endpoint_id)
        self.deployed_index_id = deployed_index_id
        self._load_metadata(metadata_bucket)

    def _load_metadata(self, bucket_name: str):
        client = storage.Client()
        bucket = client.bucket(bucket_name)
        blob = bucket.blob("index-data/metadata.json")
        self.metadata = json.loads(blob.download_as_text())

    def retrieve(self, query: str, top_k: int = 5) -> list[dict]:
        query_embedding = self.embedding_model.get_embeddings(
            [query], task_type="RETRIEVAL_QUERY"
        )[0].values

        results = self.index_endpoint.find_neighbors(
            deployed_index_id=self.deployed_index_id,
            queries=[query_embedding],
            num_neighbors=top_k,
        )

        chunks = []
        for neighbor in results[0]:
            if neighbor.id in self.metadata:
                chunk = self.metadata[neighbor.id].copy()
                chunk["similarity_score"] = neighbor.distance
                chunks.append(chunk)

        return sorted(chunks, key=lambda x: x["similarity_score"], reverse=True)

Phase 3: Generation with Gemini

import vertexai
from vertexai.generative_models import GenerativeModel, GenerationConfig

vertexai.init(project="my-project", location="europe-west4")
model = GenerativeModel("gemini-1.5-pro")

SYSTEM_PROMPT = """You are a helpful assistant that answers questions based on provided document excerpts.
Rules:
1. Answer ONLY based on the provided context. Do not use outside knowledge.
2. If the context doesn't contain enough information, say so.
3. Always cite the source document and page number for each claim.
"""

def generate_answer(query: str, retrieved_chunks: list[dict]) -> dict:
    context_parts = [
        f"[Source {i}: {c['document_name']}, page {c['page_number']}]
{c['text']}"
        for i, c in enumerate(retrieved_chunks, 1)
    ]
    context = "

".join(context_parts)

    prompt = f"Context documents:
{context}

Question: {query}

Answer (cite sources):"

    response = model.generate_content(
        prompt,
        generation_config=GenerationConfig(temperature=0.1, max_output_tokens=1000),
        system_instruction=SYSTEM_PROMPT,
    )

    cited_sources = {i for i, _ in enumerate(retrieved_chunks, 1) if f"[Source {i}]" in response.text}

    return {
        "answer": response.text,
        "sources": [
            {"index": i, "document": c["document_name"], "page": c["page_number"]}
            for i, c in enumerate(retrieved_chunks, 1)
            if i in cited_sources
        ],
    }

Phase 4: Cloud Run API Service

from flask import Flask, request, jsonify
from retriever import KnowledgeRetriever
from generator import generate_answer
import os

app = Flask(__name__)

retriever = KnowledgeRetriever(
    project_id=os.environ["PROJECT_ID"],
    index_endpoint_id=os.environ["INDEX_ENDPOINT_ID"],
    deployed_index_id=os.environ["DEPLOYED_INDEX_ID"],
    metadata_bucket=os.environ["METADATA_BUCKET"],
)

@app.route("/query", methods=["POST"])
def query():
    body = request.get_json()
    query_text = body.get("query", "").strip()
    if not query_text:
        return jsonify({"error": "Query is required"}), 400

    chunks = retriever.retrieve(query_text, top_k=body.get("top_k", 5))
    if not chunks:
        return jsonify({"answer": "No relevant documents found.", "sources": []}), 200

    return jsonify(generate_answer(query_text, chunks)), 200

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=int(os.environ.get("PORT", 8080)))

Keeping the Index Fresh

# Incremental index update (no full rebuild)
index = aiplatform.MatchingEngineIndex("INDEX_RESOURCE_NAME")
index.update_embeddings(
    contents_delta_uri="gs://my-project-embeddings/updates/batch-2025-01/",
)

Trigger updates automatically: when a new document is uploaded to Cloud Storage, an Eventarc trigger fires a Cloud Run job that processes the document and updates the index.

For the LLM serving infrastructure, see our Vertex AI Gemini enterprise guide. For the MLOps pipeline that keeps models and indexes updated, see our Vertex AI MLOps guide.