Skip to main content
Version: 2.6

Embedding Pipeline

The embedding pipeline transforms raw files into searchable vectors through a 7-step process. Each step can be customized with pluggable strategies while maintaining sensible defaults.

Pipeline Overview

Opinionated Pipeline Defaults:

  • Document Loader: aktorAzureDocumentIntelligenceLoaderStrategy - Advanced document parsing with layout analysis
  • Chunking: aktorMarkdownChunkingStrategy - Semantic chunking based on markdown structure
  • Other Steps: Use standard default strategies for text cleaning, metadata enrichment, embedding generation, and post-processing

Pipeline Function

export function aktorVektorEmbeddingPipeline(input: {
file: Aktor<File>;
client: Aktor<Client>;

// Optional strategy overrides
documentLoaderStrategy?: AktorDocumentLoaderStrategy;
textCleanerStrategy?: AktorTextCleanerStrategy;
chunkingStrategy?: AktorChunkingStrategy;
metadataStrategy?: AktorMetadataEnricherStrategy;
embeddingStrategy?: AktorEmbeddingStrategy;
postProcessingStrategy?: AktorPostProcessorStrategy;
})

Strategies

1. Document Loader Strategy

Purpose: Convert raw files into structured document content

Type Definition:

type DocumentLoaderStrategy = (input: { file: File }) => Promise<DocumentContent>;

interface DocumentContent {
id: string;
title: string;
content: string;
metadata: DocumentMetadata;
}

Pipeline Default: aktorAzureDocumentIntelligenceLoaderStrategy The pipeline uses Azure Document Intelligence by default for advanced document parsing and layout understanding.

Alternative Implementation:

export const defaultDocumentLoaderStrategy: DocumentLoaderStrategy = async ({ file }) => ({
id: generateDocumentId(file.name),
title: file.name,
content: await file.text(),
metadata: {
type: getFileType(file.name),
source: file.name,
createdAt: new Date().toISOString(),
},
});

What it does:

  • Azure Document Intelligence (default): Advanced document parsing with layout analysis, table extraction, and OCR
  • Simple loader (alternative): Basic text extraction from files
  • Generates unique document IDs
  • Creates metadata (type, source, timestamp)
  • Supports extensibility for specialized file formats

2. Text Cleaner Strategy

Purpose: Clean and normalize text content for better processing

Type Definition:

type TextCleanerStrategy = (input: { document: DocumentContent }) => Promise<CleanedDocument>;

interface CleanedDocument {
id: string;
title: string;
content: string;
metadata: DocumentMetadata;
}

Default Implementation:

export const defaultTextCleanerStrategy: TextCleanerStrategy = async ({ document }) => ({
...document,
content: document.content
.replace(/\s+/g, ' ') // Normalize whitespace
.trim() // Remove leading/trailing spaces
});

What it does:

  • Removes excessive whitespace
  • Normalizes line breaks
  • Cleans up formatting artifacts
  • Prepares text for chunking

3. Chunking Strategy

Purpose: Split documents into manageable chunks for embedding

Type Definition:

type ChunkingStrategy = (input: { document: CleanedDocument }) => Promise<DocumentChunk[]>;

interface DocumentChunk {
id: string;
documentId: string;
content: string;
index: number;
metadata: DocumentMetadata;
}

Pipeline Default: aktorMarkdownChunkingStrategy The pipeline uses an advanced token-aware Markdown chunking strategy by default for optimal semantic document structure preservation.

Default Strategy Configuration:

interface ChunkingOptions {
maxTokens?: number; // Default: 800 (optimal for embedding models)
minTokens?: number; // Default: 100 (avoid tiny chunks)
overlapTokens?: number; // Default: 100 (~12.5% overlap for context)
preserveHeaders?: boolean; // Default: true
splitByHeaders?: boolean; // Default: true
}

Alternative: Fixed-size Chunking:

export const defaultChunkingStrategy: ChunkingStrategy = async ({ document }) => {
const chunkSize = 512; // tokens
const overlap = 50; // tokens
return splitIntoChunks(document, chunkSize, overlap);
};

What the Markdown Chunking Strategy does:

Multi-Level Chunking Approach:

  1. Primary: Splits by markdown headers (H1-H6)
  2. Secondary: Falls back to paragraph-based splitting for large chunks
  3. Tertiary: Uses sentence-based splitting for fine-grained control

Token-Aware Processing:

  • Uses GPT tokenizer for accurate token counting
  • Recursively splits chunks exceeding maxTokens
  • Merges small chunks below minTokens with neighbors
  • Maintains optimal chunk sizes for embedding models

Smart Features:

  • Header preservation: Maintains markdown structure and hierarchy in metadata
  • Context overlap: Adds configurable overlap between chunks for smooth transitions
  • Metadata tracking: Each chunk includes headers hierarchy and token count
  • Intelligent merging: Combines small chunks to avoid fragmentation

Output Format: Each chunk includes enhanced metadata:

{
id: '{documentId}_md_chunk_{index}',
documentId: string,
content: string,
index: number,
metadata: {
...originalMetadata,
chunkingMethod: 'markdown-headers-token-aware',
headers: string[], // Header hierarchy
tokenCount: number // Exact token count
}
}

4. Metadata Enricher Strategy

Purpose: Add additional metadata to chunks for better retrieval

Type Definition:

type MetadataEnricherStrategy = (input: { chunks: DocumentChunk[] }) => Promise<EnrichedChunk[]>;

interface EnrichedChunk extends DocumentChunk {
enrichedMetadata?: Record<string, any>;
}

Default Implementation:

export const defaultMetadataEnricherStrategy: MetadataEnricherStrategy = async ({ chunks }) =>
chunks.map(chunk => ({
...chunk,
enrichedMetadata: {
position: chunk.index / chunks.length,
isFirstChunk: chunk.index === 0,
isLastChunk: chunk.index === chunks.length - 1,
}
}));

What it does:

  • Adds positional information
  • Enriches with extracted entities
  • Includes semantic labels
  • Provides context for retrieval

5. Embedding Strategy

Purpose: Generate vector embeddings for each chunk

Type Definition:

type EmbeddingStrategy = (input: { chunks: EnrichedChunk[] }) => Promise<EmbeddedChunk[]>;

interface EmbeddedChunk extends EnrichedChunk {
embedding: number[];
embeddingModel: string;
}

Default Implementation:

export const defaultEmbeddingStrategy: EmbeddingStrategy = async ({ chunks }) => {
const embeddings = await generateEmbeddings(chunks.map(c => c.content));
return chunks.map((chunk, i) => ({
...chunk,
embedding: embeddings[i],
embeddingModel: 'text-embedding-3-small',
}));
};

What it does:

  • Calls embedding API (OpenAI by default)
  • Batch processes for efficiency
  • Tracks model version used
  • Handles API errors and retries

6. Post-Processor Strategy

Purpose: Final processing and result formatting

Type Definition:

type PostProcessorStrategy = (input: { vectors: StoredVector[] }) => Promise<ProcessingResult>;

interface ProcessingResult {
documentsProcessed: number;
chunksCreated: number;
vectorsStored: number;
errors: string[];
processingTime: number;
}

Default Implementation:

export const defaultPostProcessorStrategy: PostProcessorStrategy = async ({ vectors }) => ({
documentsProcessed: 1,
chunksCreated: vectors.length,
vectorsStored: vectors.filter(v => v.stored).length,
errors: vectors.filter(v => !v.stored).map(v => v.error || 'Unknown error'),
processingTime: Date.now() - startTime,
});

What it does:

  • Aggregates processing statistics
  • Collects and formats errors
  • Provides success metrics
  • Enables monitoring and debugging

Complete Example

Basic Usage

import { aktorVektorEmbeddingPipeline } from '@operaide/vector';

// Using all defaults
const result = await aktorVektorEmbeddingPipeline({
file: fileAktor,
client: dbClientAktor,
}).get();

console.log(`Processed ${result.chunksCreated} chunks, stored ${result.vectorsStored} vectors`);

Advanced Usage with Custom Strategies

import {
aktorVektorEmbeddingPipeline,
aktorAzureDocumentIntelligenceLoaderStrategy,
aktorMarkdownChunkingStrategy,
} from '@operaide/vector';

// Custom embedding strategy for a different model
const customEmbeddingStrategy = ({ chunks }) => {
// Use a different embedding model or provider
const embeddings = await myCustomEmbeddingAPI(chunks);
return chunks.map((chunk, i) => ({
...chunk,
embedding: embeddings[i],
embeddingModel: 'custom-model-v1',
}));
};

// Pipeline with custom strategies
const result = await aktorVektorEmbeddingPipeline({
file: fileAktor,
client: dbClientAktor,
documentLoaderStrategy: aktorAzureDocumentIntelligenceLoaderStrategy,
chunkingStrategy: aktorMarkdownChunkingStrategy,
embeddingStrategy: customEmbeddingStrategy,
}).get();

Processing Multiple Files

// Process multiple files in parallel
const files = [file1Aktor, file2Aktor, file3Aktor];

const results = await Promise.all(
files.map(file =>
aktorVektorEmbeddingPipeline({
file,
client: dbClientAktor,
}).get()
)
);

// Aggregate results
const totalVectors = results.reduce((sum, r) => sum + r.vectorsStored, 0);
console.log(`Processed ${files.length} files, created ${totalVectors} vectors`);

Strategy Development Guide

Creating Custom Strategies

  1. Understand the Interface: Each strategy has specific input/output types
  2. Maintain Compatibility: Ensure your output matches expected structure
  3. Handle Errors: Use try-catch and provide meaningful error messages
  4. Consider Performance: Batch operations when possible

Example: Custom Chunking Strategy

export const semanticChunkingStrategy: ChunkingStrategy = async ({ document }) => {
const chunks: DocumentChunk[] = [];

// Split by paragraphs first
const paragraphs = document.content.split(/\n\n+/);

let currentChunk = '';
let chunkIndex = 0;

for (const paragraph of paragraphs) {
// Check if adding paragraph exceeds token limit
if (tokenCount(currentChunk + paragraph) > 512) {
// Save current chunk
chunks.push({
id: `${document.id}_chunk_${chunkIndex}`,
documentId: document.id,
content: currentChunk.trim(),
index: chunkIndex,
metadata: {
...document.metadata,
chunkingStrategy: 'semantic',
},
});

currentChunk = paragraph;
chunkIndex++;
} else {
currentChunk += '\n\n' + paragraph;
}
}

// Don't forget the last chunk
if (currentChunk.trim()) {
chunks.push({
id: `${document.id}_chunk_${chunkIndex}`,
documentId: document.id,
content: currentChunk.trim(),
index: chunkIndex,
metadata: {
...document.metadata,
chunkingStrategy: 'semantic',
},
});
}

return chunks;
};

Performance Considerations

Batching

  • Embedding APIs often support batch requests
  • Database insertions are faster in batches
  • Balance batch size with memory usage

Parallelization

  • Process multiple files concurrently
  • Use connection pooling for database operations
  • Monitor API rate limits

Caching

  • Cache embeddings for unchanged content
  • Use content hashes to detect changes
  • Store processing state for resume capability

Error Handling

The pipeline provides comprehensive error tracking:

  1. File-level errors: Captured in processing result
  2. Chunk-level errors: Tracked individually
  3. Partial success: Some vectors may fail while others succeed
  4. Retry logic: Built into storage operations

Monitoring

Track these metrics for production deployments:

  • Processing time: Per file and total
  • Token usage: For cost monitoring
  • Error rates: By strategy and type
  • Storage efficiency: Deduplication statistics

Next Steps