Skip to main content
EVOKORE// BROWSE
>

./browse/prompts

8 NODES
๐Ÿค–system promptโ€ข7 months ago

spark-optimization

Optimize Apache Spark jobs with partitioning, caching, shuffle

data
โญ1
# Apache Spark Optimization Production patterns for optimizing Apache Spark jobs including partitioning strategies, memory management, shuffle optimization, and performance tuning. ## When to Use This Skill - Optimizing slow Spark jobs - Tuning memory and executor configuration - Implementing efficient partitioning strategies - Debugging Spark performance issues - Scaling Spark pipelines for large datasets - Reducing shuffle and data skew ## Core Concepts ### 1. Spark Execution Model ``` Driver Program โ†“ Job (triggered by action) โ†“ Stages (separated by shuffles) โ†“ Tasks (one per partition) ``` ### 2. Key Performance Factors | Factor | Impact | Solution | | ----------------- | --------------------- | ----------------------------- | | **Shuffle** | Network I/O, disk I/O | Minimize wide transformations | | **Data Skew** | Uneven task duration | Salting, broadcast joins | | **Serialization** | CPU overhead | Use Kryo, columnar formats | | **Memory** | GC pressure, spills | Tune executor memory | | **Partitions** | Parallelism | Right-size partitions | ## Quick Start ```python from pyspark.sql import SparkSession from pyspark.sql import functions as F # Create optimized Spark session spark = (SparkSession.builder .appName("OptimizedJob") .config("spark.sql.adaptive.enabled", "true") .config("spark.sql.adaptive.coalescePartitions.enabled", "true") .config("spark.sql.adaptive.skewJoin.enabled", "true") .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") .config("spark.sql.shuffle.partitions", "200") .getOrCreate()) # Read with optimized settings df = (spark.read .format("parquet") .option("mergeSchema", "false") .load("s3://bucket/data/")) # Efficient transformations result = (df .filter(F.col("date") >= "2024-01-01") .select("id", "amount", "category") .groupBy("category") .agg(F.sum("amount").alias("total"))) result.write.mode("overwrite").parquet("s3://bucket/output/") ``` ## Patterns ### Pattern 1: Optimal Partitioning ```python # Calculate optimal partition count def calculate_partitions(data_size_gb: float, partition_size_mb: int = 128) -> int: """ Optimal partition size: 128MB - 256MB Too few: Under-utilization, memory pressure Too many: Task scheduling overhead """ return max(int(data_size_gb * 1024 / partition_size_mb), 1) # Repartition for even distribution df_repartitioned = df.repartition(200, "partition_key") # Coalesce to reduce partitions (no shuffle) df_coalesced = df.coalesce(100) # Partition pruning with predicate pushdown df = (spark.read.parquet("s3://bucket/data/") .filter(F.col("date") == "2024-01-01")) # Spark pushes this down # Write with partitioning for future queries (df.write .partitionBy("year", "month", "day") .mode("overwrite") .parquet("s3://bucket/partitioned_output/")) ``` ### Pattern 2: Join Optimization ```python from pyspark.sql import functions as F from pyspark.sql.types import * # 1. Broadcast Join - Small table joins # Best when: One side < 10MB (configurable) small_df = spark.read.parquet("s3://bucket/small_table/") # < 10MB large_df = spark.read.parquet("s3://bucket/large_table/") # TBs # Explicit broadcast hint result = large_df.join( F.broadcast(small_df), on="key", how="left" ) # 2. Sort-Merge Join - Default for large tables # Requires shuffle, but handles any size result = large_df1.join(large_df2, on="key", how="inner") # 3. Bucket Join - Pre-sorted, no shuffle at join time # Write bucketed tables (df.write .bucketBy(200, "customer_id") .sortBy("customer_id") .mode("overwrite") .saveAsTable("bucketed_orders")) # Join bucketed tables (no shuffle!) orders = spark.table("bucketed_orders") customers = spark.table("bucketed_customers") # Same bucket count result = orders.join(customers, on="customer_id") # 4. Skew Join Handling # Enable AQE skew join optimization spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true") spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5") spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256MB") # Manual salting for severe skew def salt_join(df_skewed, df_other, key_col, num_salts=10): """Add salt to distribute skewed keys""" # Add salt to skewed side df_salted = df_skewed.withColumn( "salt", (F.rand() * num_salts).cast("int") ).withColumn( "salted_key", F.concat(F.col(key_col), F.lit("_"), F.col("salt")) ) # Explode other side with all salts df_exploded = df_other.crossJoin( spark.range(num_salts).withColumnRenamed("id", "salt") ).withColumn( "salted_key", F.concat(F.col(key_col), F.lit("_"), F.col("salt")) ) # Join on salted key return df_salted.join(df_exploded, on="salted_key", how="inner") ``` ### Pattern 3: Caching and Persistence ```python from pyspark import StorageLevel # Cache when reusing DataFrame multiple times df = spark.read.parquet("s3://bucket/data/") df_filtered = df.filter(F.col("status") == "active") # Cache in memory (MEMORY_AND_DISK is default) df_filtered.cache() # Or with specific storage level df_filtered.persist(StorageLevel.MEMORY_AND_DISK_SER) # Force materialization df_filtered.count() # Use in multiple actions agg1 = df_filtered.groupBy("category").count() agg2 = df_filtered.groupBy("region").sum("amount") # Unpersist when done df_filtered.unpersist() # Storage levels explained: # MEMORY_ONLY - Fast, but may not fit # MEMORY_AND_DISK - Spills to disk if needed (recommended) # MEMORY_ONLY_SER - Serialized, less memory, more CPU # DISK_ONLY - When memory is tight # OFF_HEAP - Tungsten off-heap memory # Checkpoint for complex lineage spark.sparkContext.setCheckpointDir("s3://bucket/checkpoints/") df_complex = (df .join(other_df, "key") .groupBy("category") .agg(F.sum("amount"))) df_complex.checkpoint() # Breaks lineage, materializes ``` ### Pattern 4: Memory Tuning ```python # Executor memory configuration # spark-submit --executor-memory 8g --executor-cores 4 # Memory breakdown (8GB executor): # - spark.memory.fraction = 0.6 (60% = 4.8GB for execution + storage) # - spark.memory.storageFraction = 0.5 (50% of 4.8GB = 2.4GB for cache) # - Remaining 2.4GB for execution (shuffles, joins, sorts) # - 40% = 3.2GB for user data structures and internal metadata spark = (SparkSession.builder .config("spark.executor.memory", "8g") .config("spark.executor.memoryOverhead", "2g") # For non-JVM memory .config("spark.memory.fraction", "0.6") .config("spark.memory.storageFraction", "0.5") .config("spark.sql.shuffle.partitions", "200") # For memory-intensive operations .config("spark.sql.autoBroadcastJoinThreshold", "50MB") # Prevent OOM on large shuffles .config("spark.sql.files.maxPartitionBytes", "128MB") .getOrCreate()) # Monitor memory usage def print_memory_usage(spark): """Print current memory usage""" sc = spark.sparkContext for executor in sc._jsc.sc().getExecutorMemoryStatus().keySet().toArray(): mem_status = sc._jsc.sc().getExecutorMemoryStatus().get(executor) total = mem_status._1() / (1024**3) free = mem_status._2() / (1024**3) print(f"{executor}: {total:.2f}GB total, {free:.2f}GB free") ``` ### Pattern 5: Shuffle Optimization ```python # Reduce shuffle data size spark.conf.set("spark.sql.shuffle.partitions", "auto") # With AQE spark.conf.set("spark.shuffle.compress", "true") spark.conf.set("spark.shuffle.spill.compress", "true") # Pre-aggregate before shuffle df_optimized = (df # Local aggregation first (combiner) .groupBy("key", "partition_col") .agg(F.sum("value").alias("partial_sum")) # Then global aggregation .groupBy("key") .agg(F.sum("partial_sum").alias("total"))) # Avoid shuffle with map-side operations # BAD: Shuffle for each distinct distinct_count = df.select("category").distinct().count() # GOOD: Approximate distinct (no shuffle) approx_count = df.select(F.approx_count_distinct("category")).collect()[0][0] # Use coalesce instead of repartition when reducing partitions df_reduced = df.coalesce(10) # No shuffle # Optimize shuffle with compression spark.conf.set("spark.io.compression.codec", "lz4") # Fast compression ``` ### Pattern 6: Data Format Optimization ```python # Parquet optimizations (df.write .option("compression", "snappy") # Fast compression .option("parquet.block.size", 128 * 1024 * 1024) # 128MB row groups .parquet("s3://bucket/output/")) # Column pruning - only read needed columns df = (spark.read.parquet("s3://bucket/data/") .select("id", "amount", "date")) # Spark only reads these columns # Predicate pushdown - filter at storage level df = (spark.read.parquet("s3://bucket/partitioned/year=2024/") .filter(F.col("status") == "active")) # Pushed to Parquet reader # Delta Lake optimizations (df.write .format("delta") .option("optimizeWrite", "true") # Bin-packing .option("autoCompact", "true") # Compact small files .mode("overwrite") .save("s3://bucket/delta_table/")) # Z-ordering for multi-dimensional queries spark.sql(""" OPTIMIZE delta.`s3://bucket/delta_table/` ZORDER BY (customer_id, date) """) ``` ### Pattern 7: Monitoring and Debugging ```python # Enable detailed metrics spark.conf.set("spark.sql.codegen.wholeStage", "true") spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true") # Explain query plan df.explain(mode="extended") # Modes: simple, extended, codegen, cost, formatted # Get physical plan statistics df.explain(mode="cost") # Monitor task metrics def analyze_stage_metrics(spark): """Analyze recent stage metrics""" status_tracker = spark.sparkContext.statusTracker() for stage_id in status_tracker.getActiveStageIds(): stage_info = status_tracker.getStageInfo(stage_id) print(f"Stage {stage_id}:") print(f" Tasks: {stage_info.numTasks}") print(f" Completed: {stage_info.numCompletedTasks}") print(f" Failed: {stage_info.numFailedTasks}") # Identify data skew def check_partition_skew(df): """Check for partition skew""" partition_counts = (df .withColumn("partition_id", F.spark_partition_id()) .groupBy("partition_id") .count() .orderBy(F.desc("count"))) partition_counts.show(20) stats = partition_counts.select( F.min("count").alias("min"), F.max("count").alias("max"), F.avg("count").alias("avg"), F.stddev("count").alias("stddev") ).collect()[0] skew_ratio = stats["max"] / stats["avg"] print(f"Skew ratio: {skew_ratio:.2f}x (>2x indicates skew)") ``` ## Configuration Cheat Sheet ```python # Production configuration template spark_configs = { # Adaptive Query Execution (AQE) "spark.sql.adaptive.enabled": "true", "spark.sql.adaptive.coalescePartitions.enabled": "true", "spark.sql.adaptive.skewJoin.enabled": "true", # Memory "spark.executor.memory": "8g", "spark.executor.memoryOverhead": "2g", "spark.memory.fraction": "0.6", "spark.memory.storageFraction": "0.5", # Parallelism "spark.sql.shuffle.partitions": "200", "spark.default.parallelism": "200", # Serialization "spark.serializer": "org.apache.spark.serializer.KryoSerializer", "spark.sql.execution.arrow.pyspark.enabled": "true", # Compression "spark.io.compression.codec": "lz4", "spark.shuffle.compress": "true", # Broadcast "spark.sql.autoBroadcastJoinThreshold": "50MB", # File handling "spark.sql.files.maxPartitionBytes": "128MB", "spark.sql.files.openCostInBytes": "4MB", } ``` ## Best Practices ### Do's - **Enable AQE** - Adaptive query execution handles many issues - **Use Parquet/Delta** - Columnar formats with compression - **Broadcast small tables** - Avoid shuffle for small joins - **Monitor Spark UI** - Check for skew, spills, GC - **Right-size partitions** - 128MB - 256MB per partition ### Don'ts - **Don't collect large data** - Keep data distributed - **Don't use UDFs unnecessarily** - Use built-in functions - **Don't over-cache** - Memory is limited - **Don't ignore data skew** - It dominates job time - **Don't use `.count()` for existence** - Use `.take(1)` or `.isEmpty()` ## Resources - [Spark Performance Tuning](https://spark.apache.org/docs/latest/sql-performance-tuning.html) - [Spark Configuration](https://spark.apache.org/docs/latest/configuration.html) - [Databricks Optimization Guide](https://docs.databricks.com/en/optimizations/index.html)
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

postmortem-writing

Write effective blameless postmortems with root cause analysis,

coding
โญ1
# Postmortem Writing Comprehensive guide to writing effective, blameless postmortems that drive organizational learning and prevent incident recurrence. ## When to Use This Skill - Conducting post-incident reviews - Writing postmortem documents - Facilitating blameless postmortem meetings - Identifying root causes and contributing factors - Creating actionable follow-up items - Building organizational learning culture ## Core Concepts ### 1. Blameless Culture | Blame-Focused | Blameless | | ------------------------ | --------------------------------- | | "Who caused this?" | "What conditions allowed this?" | | "Someone made a mistake" | "The system allowed this mistake" | | Punish individuals | Improve systems | | Hide information | Share learnings | | Fear of speaking up | Psychological safety | ### 2. Postmortem Triggers - SEV1 or SEV2 incidents - Customer-facing outages > 15 minutes - Data loss or security incidents - Near-misses that could have been severe - Novel failure modes - Incidents requiring unusual intervention ## Quick Start ### Postmortem Timeline ``` Day 0: Incident occurs Day 1-2: Draft postmortem document Day 3-5: Postmortem meeting Day 5-7: Finalize document, create tickets Week 2+: Action item completion Quarterly: Review patterns across incidents ``` ## Templates ### Template 1: Standard Postmortem ```markdown # Postmortem: [Incident Title] **Date**: 2024-01-15 **Authors**: @alice, @bob **Status**: Draft | In Review | Final **Incident Severity**: SEV2 **Incident Duration**: 47 minutes ## Executive Summary On January 15, 2024, the payment processing service experienced a 47-minute outage affecting approximately 12,000 customers. The root cause was a database connection pool exhaustion triggered by a configuration change in deployment v2.3.4. The incident was resolved by rolling back to v2.3.3 and increasing connection pool limits. **Impact**: - 12,000 customers unable to complete purchases - Estimated revenue loss: $45,000 - 847 support tickets created - No data loss or security implications ## Timeline (All times UTC) | Time | Event | | ----- | ----------------------------------------------- | | 14:23 | Deployment v2.3.4 completed to production | | 14:31 | First alert: `payment_error_rate > 5%` | | 14:33 | On-call engineer @alice acknowledges alert | | 14:35 | Initial investigation begins, error rate at 23% | | 14:41 | Incident declared SEV2, @bob joins | | 14:45 | Database connection exhaustion identified | | 14:52 | Decision to rollback deployment | | 14:58 | Rollback to v2.3.3 initiated | | 15:10 | Rollback complete, error rate dropping | | 15:18 | Service fully recovered, incident resolved | ## Root Cause Analysis ### What Happened The v2.3.4 deployment included a change to the database query pattern that inadvertently removed connection pooling for a frequently-called endpoint. Each request opened a new database connection instead of reusing pooled connections. ### Why It Happened 1. **Proximate Cause**: Code change in `PaymentRepository.java` replaced pooled `DataSource` with direct `DriverManager.getConnection()` calls. 2. **Contributing Factors**: - Code review did not catch the connection handling change - No integration tests specifically for connection pool behavior - Staging environment has lower traffic, masking the issue - Database connection metrics alert threshold was too high (90%) 3. **5 Whys Analysis**: - Why did the service fail? โ†’ Database connections exhausted - Why were connections exhausted? โ†’ Each request opened new connection - Why did each request open new connection? โ†’ Code bypassed connection pool - Why did code bypass connection pool? โ†’ Developer unfamiliar with codebase patterns - Why was developer unfamiliar? โ†’ No documentation on connection management patterns ### System Diagram ``` [Client] โ†’ [Load Balancer] โ†’ [Payment Service] โ†’ [Database] โ†“ Connection Pool (broken) โ†“ Direct connections (cause) ``` ## Detection ### What Worked - Error rate alert fired within 8 minutes of deployment - Grafana dashboard clearly showed connection spike - On-call response was swift (2 minute acknowledgment) ### What Didn't Work - Database connection metric alert threshold too high - No deployment-correlated alerting - Canary deployment would have caught this earlier ### Detection Gap The deployment completed at 14:23, but the first alert didn't fire until 14:31 (8 minutes). A deployment-aware alert could have detected the issue faster. ## Response ### What Worked - On-call engineer quickly identified database as the issue - Rollback decision was made decisively - Clear communication in incident channel ### What Could Be Improved - Took 10 minutes to correlate issue with recent deployment - Had to manually check deployment history - Rollback took 12 minutes (could be faster) ## Impact ### Customer Impact - 12,000 unique customers affected - Average impact duration: 35 minutes - 847 support tickets (23% of affected users) - Customer satisfaction score dropped 12 points ### Business Impact - Estimated revenue loss: $45,000 - Support cost: ~$2,500 (agent time) - Engineering time: ~8 person-hours ### Technical Impact - Database primary experienced elevated load - Some replica lag during incident - No permanent damage to systems ## Lessons Learned ### What Went Well 1. Alerting detected the issue before customer reports 2. Team collaborated effectively under pressure 3. Rollback procedure worked smoothly 4. Communication was clear and timely ### What Went Wrong 1. Code review missed critical change 2. Test coverage gap for connection pooling 3. Staging environment doesn't reflect production traffic 4. Alert thresholds were not tuned properly ### Where We Got Lucky 1. Incident occurred during business hours with full team available 2. Database handled the load without failing completely 3. No other incidents occurred simultaneously ## Action Items | Priority | Action | Owner | Due Date | Ticket | |----------|--------|-------|----------|--------| | P0 | Add integration test for connection pool behavior | @alice | 2024-01-22 | ENG-1234 | | P0 | Lower database connection alert threshold to 70% | @bob | 2024-01-17 | OPS-567 | | P1 | Document connection management patterns | @alice | 2024-01-29 | DOC-89 | | P1 | Implement deployment-correlated alerting | @bob | 2024-02-05 | OPS-568 | | P2 | Evaluate canary deployment strategy | @charlie | 2024-02-15 | ENG-1235 | | P2 | Load test staging with production-like traffic | @dave | 2024-02-28 | QA-123 | ## Appendix ### Supporting Data #### Error Rate Graph [Link to Grafana dashboard snapshot] #### Database Connection Graph [Link to metrics] ### Related Incidents - 2023-11-02: Similar connection issue in User Service (POSTMORTEM-42) ### References - [Connection Pool Best Practices](internal-wiki/connection-pools) - [Deployment Runbook](internal-wiki/deployment-runbook) ``` ### Template 2: 5 Whys Analysis ```markdown # 5 Whys Analysis: [Incident] ## Problem Statement Payment service experienced 47-minute outage due to database connection exhaustion. ## Analysis ### Why #1: Why did the service fail? **Answer**: Database connections were exhausted, causing all new requests to fail. **Evidence**: Metrics showed connection count at 100/100 (max), with 500+ pending requests. --- ### Why #2: Why were database connections exhausted? **Answer**: Each incoming request opened a new database connection instead of using the connection pool. **Evidence**: Code diff shows direct `DriverManager.getConnection()` instead of pooled `DataSource`. --- ### Why #3: Why did the code bypass the connection pool? **Answer**: A developer refactored the repository class and inadvertently changed the connection acquisition method. **Evidence**: PR #1234 shows the change, made while fixing a different bug. --- ### Why #4: Why wasn't this caught in code review? **Answer**: The reviewer focused on the functional change (the bug fix) and didn't notice the infrastructure change. **Evidence**: Review comments only discuss business logic. --- ### Why #5: Why isn't there a safety net for this type of change? **Answer**: We lack automated tests that verify connection pool behavior and lack documentation about our connection patterns. **Evidence**: Test suite has no tests for connection handling; wiki has no article on database connections. ## Root Causes Identified 1. **Primary**: Missing automated tests for infrastructure behavior 2. **Secondary**: Insufficient documentation of architectural patterns 3. **Tertiary**: Code review checklist doesn't include infrastructure considerations ## Systemic Improvements | Root Cause | Improvement | Type | | ------------- | --------------------------------- | ---------- | | Missing tests | Add infrastructure behavior tests | Prevention | | Missing docs | Document connection patterns | Prevention | | Review gaps | Update review checklist | Detection | | No canary | Implement canary deployments | Mitigation | ``` ### Template 3: Quick Postmortem (Minor Incidents) ```markdown # Quick Postmortem: [Brief Title] **Date**: 2024-01-15 | **Duration**: 12 min | **Severity**: SEV3 ## What Happened API latency spiked to 5s due to cache miss storm after cache flush. ## Timeline - 10:00 - Cache flush initiated for config update - 10:02 - Latency alerts fire - 10:05 - Identified as cache miss storm - 10:08 - Enabled cache warming - 10:12 - Latency normalized ## Root Cause Full cache flush for minor config update caused thundering herd. ## Fix - Immediate: Enabled cache warming - Long-term: Implement partial cache invalidation (ENG-999) ## Lessons Don't full-flush cache in production; use targeted invalidation. ``` ## Facilitation Guide ### Running a Postmortem Meeting ```markdown ## Meeting Structure (60 minutes) ### 1. Opening (5 min) - Remind everyone of blameless culture - "We're here to learn, not to blame" - Review meeting norms ### 2. Timeline Review (15 min) - Walk through events chronologically - Ask clarifying questions - Identify gaps in timeline ### 3. Analysis Discussion (20 min) - What failed? - Why did it fail? - What conditions allowed this? - What would have prevented it? ### 4. Action Items (15 min) - Brainstorm improvements - Prioritize by impact and effort - Assign owners and due dates ### 5. Closing (5 min) - Summarize key learnings - Confirm action item owners - Schedule follow-up if needed ## Facilitation Tips - Keep discussion on track - Redirect blame to systems - Encourage quiet participants - Document dissenting views - Time-box tangents ``` ## Anti-Patterns to Avoid | Anti-Pattern | Problem | Better Approach | | ----------------------- | -------------------------- | ------------------------------- | | **Blame game** | Shuts down learning | Focus on systems | | **Shallow analysis** | Doesn't prevent recurrence | Ask "why" 5 times | | **No action items** | Waste of time | Always have concrete next steps | | **Unrealistic actions** | Never completed | Scope to achievable tasks | | **No follow-up** | Actions forgotten | Track in ticketing system | ## Best Practices ### Do's - **Start immediately** - Memory fades fast - **Be specific** - Exact times, exact errors - **Include graphs** - Visual evidence - **Assign owners** - No orphan action items - **Share widely** - Organizational learning ### Don'ts - **Don't name and shame** - Ever - **Don't skip small incidents** - They reveal patterns - **Don't make it a blame doc** - That kills learning - **Don't create busywork** - Actions should be meaningful - **Don't skip follow-up** - Verify actions completed ## Resources - [Google SRE - Postmortem Culture](https://sre.google/sre-book/postmortem-culture/) - [Etsy's Blameless Postmortems](https://codeascraft.com/2012/05/22/blameless-postmortems/) - [PagerDuty Postmortem Guide](https://postmortems.pagerduty.com/)
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

embedding-strategies

Select and optimize embedding models for semantic search and RAG

coding
โญ1
# Embedding Strategies Guide to selecting and optimizing embedding models for vector search applications. ## When to Use This Skill - Choosing embedding models for RAG - Optimizing chunking strategies - Fine-tuning embeddings for domains - Comparing embedding model performance - Reducing embedding dimensions - Handling multilingual content ## Core Concepts ### 1. Embedding Model Comparison (2026) | Model | Dimensions | Max Tokens | Best For | | -------------------------- | ---------- | ---------- | ----------------------------------- | | **voyage-3-large** | 1024 | 32000 | Claude apps (Anthropic recommended) | | **voyage-3** | 1024 | 32000 | Claude apps, cost-effective | | **voyage-code-3** | 1024 | 32000 | Code search | | **voyage-finance-2** | 1024 | 32000 | Financial documents | | **voyage-law-2** | 1024 | 32000 | Legal documents | | **text-embedding-3-large** | 3072 | 8191 | OpenAI apps, high accuracy | | **text-embedding-3-small** | 1536 | 8191 | OpenAI apps, cost-effective | | **bge-large-en-v1.5** | 1024 | 512 | Open source, local deployment | | **all-MiniLM-L6-v2** | 384 | 256 | Fast, lightweight | | **multilingual-e5-large** | 1024 | 512 | Multi-language | ### 2. Embedding Pipeline ``` Document โ†’ Chunking โ†’ Preprocessing โ†’ Embedding Model โ†’ Vector โ†“ [Overlap, Size] [Clean, Normalize] [API/Local] ``` ## Templates ### Template 1: Voyage AI Embeddings (Recommended for Claude) ```python from langchain_voyageai import VoyageAIEmbeddings from typing import List import os # Initialize Voyage AI embeddings (recommended by Anthropic for Claude) embeddings = VoyageAIEmbeddings( model="voyage-3-large", voyage_api_key=os.environ.get("VOYAGE_API_KEY") ) def get_embeddings(texts: List[str]) -> List[List[float]]: """Get embeddings from Voyage AI.""" return embeddings.embed_documents(texts) def get_query_embedding(query: str) -> List[float]: """Get single query embedding.""" return embeddings.embed_query(query) # Specialized models for domains code_embeddings = VoyageAIEmbeddings(model="voyage-code-3") finance_embeddings = VoyageAIEmbeddings(model="voyage-finance-2") legal_embeddings = VoyageAIEmbeddings(model="voyage-law-2") ``` ### Template 2: OpenAI Embeddings ```python from openai import OpenAI from typing import List import numpy as np client = OpenAI() def get_embeddings( texts: List[str], model: str = "text-embedding-3-small", dimensions: int = None ) -> List[List[float]]: """Get embeddings from OpenAI with optional dimension reduction.""" # Handle batching for large lists batch_size = 100 all_embeddings = [] for i in range(0, len(texts), batch_size): batch = texts[i:i + batch_size] kwargs = {"input": batch, "model": model} if dimensions: # Matryoshka dimensionality reduction kwargs["dimensions"] = dimensions response = client.embeddings.create(**kwargs) embeddings = [item.embedding for item in response.data] all_embeddings.extend(embeddings) return all_embeddings def get_embedding(text: str, **kwargs) -> List[float]: """Get single embedding.""" return get_embeddings([text], **kwargs)[0] # Dimension reduction with Matryoshka embeddings def get_reduced_embedding(text: str, dimensions: int = 512) -> List[float]: """Get embedding with reduced dimensions (Matryoshka).""" return get_embedding( text, model="text-embedding-3-small", dimensions=dimensions ) ``` ### Template 3: Local Embeddings with Sentence Transformers ```python from sentence_transformers import SentenceTransformer from typing import List, Optional import numpy as np class LocalEmbedder: """Local embedding with sentence-transformers.""" def __init__( self, model_name: str = "BAAI/bge-large-en-v1.5", device: str = "cuda" ): self.model = SentenceTransformer(model_name, device=device) self.model_name = model_name def embed( self, texts: List[str], normalize: bool = True, show_progress: bool = False ) -> np.ndarray: """Embed texts with optional normalization.""" embeddings = self.model.encode( texts, normalize_embeddings=normalize, show_progress_bar=show_progress, convert_to_numpy=True ) return embeddings def embed_query(self, query: str) -> np.ndarray: """Embed a query with appropriate prefix for retrieval models.""" # BGE and similar models benefit from query prefix if "bge" in self.model_name.lower(): query = f"Represent this sentence for searching relevant passages: {query}" return self.embed([query])[0] def embed_documents(self, documents: List[str]) -> np.ndarray: """Embed documents for indexing.""" return self.embed(documents) # E5 model with instructions class E5Embedder: def __init__(self, model_name: str = "intfloat/multilingual-e5-large"): self.model = SentenceTransformer(model_name) def embed_query(self, query: str) -> np.ndarray: """E5 requires 'query:' prefix for queries.""" return self.model.encode(f"query: {query}") def embed_document(self, document: str) -> np.ndarray: """E5 requires 'passage:' prefix for documents.""" return self.model.encode(f"passage: {document}") ``` ### Template 4: Chunking Strategies ```python from typing import List, Tuple import re def chunk_by_tokens( text: str, chunk_size: int = 512, chunk_overlap: int = 50, tokenizer=None ) -> List[str]: """Chunk text by token count.""" import tiktoken tokenizer = tokenizer or tiktoken.get_encoding("cl100k_base") tokens = tokenizer.encode(text) chunks = [] start = 0 while start < len(tokens): end = start + chunk_size chunk_tokens = tokens[start:end] chunk_text = tokenizer.decode(chunk_tokens) chunks.append(chunk_text) start = end - chunk_overlap return chunks def chunk_by_sentences( text: str, max_chunk_size: int = 1000, min_chunk_size: int = 100 ) -> List[str]: """Chunk text by sentences, respecting size limits.""" import nltk sentences = nltk.sent_tokenize(text) chunks = [] current_chunk = [] current_size = 0 for sentence in sentences: sentence_size = len(sentence) if current_size + sentence_size > max_chunk_size and current_chunk: chunks.append(" ".join(current_chunk)) current_chunk = [] current_size = 0 current_chunk.append(sentence) current_size += sentence_size if current_chunk: chunks.append(" ".join(current_chunk)) return chunks def chunk_by_semantic_sections( text: str, headers_pattern: str = r'^#{1,3}\s+.+$' ) -> List[Tuple[str, str]]: """Chunk markdown by headers, preserving hierarchy.""" lines = text.split('\n') chunks = [] current_header = "" current_content = [] for line in lines: if re.match(headers_pattern, line, re.MULTILINE): if current_content: chunks.append((current_header, '\n'.join(current_content))) current_header = line current_content = [] else: current_content.append(line) if current_content: chunks.append((current_header, '\n'.join(current_content))) return chunks def recursive_character_splitter( text: str, chunk_size: int = 1000, chunk_overlap: int = 200, separators: List[str] = None ) -> List[str]: """LangChain-style recursive splitter.""" separators = separators or ["\n\n", "\n", ". ", " ", ""] def split_text(text: str, separators: List[str]) -> List[str]: if not text: return [] separator = separators[0] remaining_separators = separators[1:] if separator == "": # Character-level split return [text[i:i+chunk_size] for i in range(0, len(text), chunk_size - chunk_overlap)] splits = text.split(separator) chunks = [] current_chunk = [] current_length = 0 for split in splits: split_length = len(split) + len(separator) if current_length + split_length > chunk_size and current_chunk: chunk_text = separator.join(current_chunk) # Recursively split if still too large if len(chunk_text) > chunk_size and remaining_separators: chunks.extend(split_text(chunk_text, remaining_separators)) else: chunks.append(chunk_text) # Start new chunk with overlap overlap_splits = [] overlap_length = 0 for s in reversed(current_chunk): if overlap_length + len(s) <= chunk_overlap: overlap_splits.insert(0, s) overlap_length += len(s) else: break current_chunk = overlap_splits current_length = overlap_length current_chunk.append(split) current_length += split_length if current_chunk: chunks.append(separator.join(current_chunk)) return chunks return split_text(text, separators) ``` ### Template 5: Domain-Specific Embedding Pipeline ```python import re from typing import List, Optional from dataclasses import dataclass @dataclass class EmbeddedDocument: id: str document_id: str chunk_index: int text: str embedding: List[float] metadata: dict class DomainEmbeddingPipeline: """Pipeline for domain-specific embeddings.""" def __init__( self, embedding_model: str = "voyage-3-large", chunk_size: int = 512, chunk_overlap: int = 50, preprocessing_fn=None ): self.embeddings = VoyageAIEmbeddings(model=embedding_model) self.chunk_size = chunk_size self.chunk_overlap = chunk_overlap self.preprocess = preprocessing_fn or self._default_preprocess def _default_preprocess(self, text: str) -> str: """Default preprocessing.""" # Remove excessive whitespace text = re.sub(r'\s+', ' ', text) # Remove special characters (customize for your domain) text = re.sub(r'[^\w\s.,!?-]', '', text) return text.strip() async def process_documents( self, documents: List[dict], id_field: str = "id", content_field: str = "content", metadata_fields: Optional[List[str]] = None ) -> List[EmbeddedDocument]: """Process documents for vector storage.""" processed = [] for doc in documents: content = doc[content_field] doc_id = doc[id_field] # Preprocess cleaned = self.preprocess(content) # Chunk chunks = chunk_by_tokens( cleaned, self.chunk_size, self.chunk_overlap ) # Create embeddings embeddings = await self.embeddings.aembed_documents(chunks) # Create records for i, (chunk, embedding) in enumerate(zip(chunks, embeddings)): metadata = {"document_id": doc_id, "chunk_index": i} # Add specified metadata fields if metadata_fields: for field in metadata_fields: if field in doc: metadata[field] = doc[field] processed.append(EmbeddedDocument( id=f"{doc_id}_chunk_{i}", document_id=doc_id, chunk_index=i, text=chunk, embedding=embedding, metadata=metadata )) return processed # Code-specific pipeline class CodeEmbeddingPipeline: """Specialized pipeline for code embeddings.""" def __init__(self): # Use Voyage's code-specific model self.embeddings = VoyageAIEmbeddings(model="voyage-code-3") def chunk_code(self, code: str, language: str) -> List[dict]: """Chunk code by functions/classes using tree-sitter.""" try: import tree_sitter_languages parser = tree_sitter_languages.get_parser(language) tree = parser.parse(bytes(code, "utf8")) chunks = [] # Extract function and class definitions self._extract_nodes(tree.root_node, code, chunks) return chunks except ImportError: # Fallback to simple chunking return [{"text": code, "type": "module"}] def _extract_nodes(self, node, source_code: str, chunks: list): """Recursively extract function/class definitions.""" if node.type in ['function_definition', 'class_definition', 'method_definition']: text = source_code[node.start_byte:node.end_byte] chunks.append({ "text": text, "type": node.type, "name": self._get_name(node), "start_line": node.start_point[0], "end_line": node.end_point[0] }) for child in node.children: self._extract_nodes(child, source_code, chunks) def _get_name(self, node) -> str: """Extract name from function/class node.""" for child in node.children: if child.type == 'identifier' or child.type == 'name': return child.text.decode('utf8') return "unknown" async def embed_with_context( self, chunk: str, context: str = "" ) -> List[float]: """Embed code with surrounding context.""" if context: combined = f"Context: {context}\n\nCode:\n{chunk}" else: combined = chunk return await self.embeddings.aembed_query(combined) ``` ### Template 6: Embedding Quality Evaluation ```python import numpy as np from typing import List, Dict def evaluate_retrieval_quality( queries: List[str], relevant_docs: List[List[str]], # List of relevant doc IDs per query retrieved_docs: List[List[str]], # List of retrieved doc IDs per query k: int = 10 ) -> Dict[str, float]: """Evaluate embedding quality for retrieval.""" def precision_at_k(relevant: set, retrieved: List[str], k: int) -> float: retrieved_k = retrieved[:k] relevant_retrieved = len(set(retrieved_k) & relevant) return relevant_retrieved / k if k > 0 else 0 def recall_at_k(relevant: set, retrieved: List[str], k: int) -> float: retrieved_k = retrieved[:k] relevant_retrieved = len(set(retrieved_k) & relevant) return relevant_retrieved / len(relevant) if relevant else 0 def mrr(relevant: set, retrieved: List[str]) -> float: for i, doc in enumerate(retrieved): if doc in relevant: return 1 / (i + 1) return 0 def ndcg_at_k(relevant: set, retrieved: List[str], k: int) -> float: dcg = sum( 1 / np.log2(i + 2) if doc in relevant else 0 for i, doc in enumerate(retrieved[:k]) ) ideal_dcg = sum(1 / np.log2(i + 2) for i in range(min(len(relevant), k))) return dcg / ideal_dcg if ideal_dcg > 0 else 0 metrics = { f"precision@{k}": [], f"recall@{k}": [], "mrr": [], f"ndcg@{k}": [] } for relevant, retrieved in zip(relevant_docs, retrieved_docs): relevant_set = set(relevant) metrics[f"precision@{k}"].append(precision_at_k(relevant_set, retrieved, k)) metrics[f"recall@{k}"].append(recall_at_k(relevant_set, retrieved, k)) metrics["mrr"].append(mrr(relevant_set, retrieved)) metrics[f"ndcg@{k}"].append(ndcg_at_k(relevant_set, retrieved, k)) return {name: np.mean(values) for name, values in metrics.items()} def compute_embedding_similarity( embeddings1: np.ndarray, embeddings2: np.ndarray, metric: str = "cosine" ) -> np.ndarray: """Compute similarity matrix between embedding sets.""" if metric == "cosine": # Normalize and compute dot product norm1 = embeddings1 / np.linalg.norm(embeddings1, axis=1, keepdims=True) norm2 = embeddings2 / np.linalg.norm(embeddings2, axis=1, keepdims=True) return norm1 @ norm2.T elif metric == "euclidean": from scipy.spatial.distance import cdist return -cdist(embeddings1, embeddings2, metric='euclidean') elif metric == "dot": return embeddings1 @ embeddings2.T else: raise ValueError(f"Unknown metric: {metric}") def compare_embedding_models( texts: List[str], models: Dict[str, callable], queries: List[str], relevant_indices: List[List[int]], k: int = 5 ) -> Dict[str, Dict[str, float]]: """Compare multiple embedding models on retrieval quality.""" results = {} for model_name, embed_fn in models.items(): # Embed all texts doc_embeddings = np.array(embed_fn(texts)) retrieved_per_query = [] for query in queries: query_embedding = np.array(embed_fn([query])[0]) # Compute similarities similarities = compute_embedding_similarity( query_embedding.reshape(1, -1), doc_embeddings, metric="cosine" )[0] # Get top-k indices top_k_indices = np.argsort(similarities)[::-1][:k] retrieved_per_query.append([str(i) for i in top_k_indices]) # Convert relevant indices to string IDs relevant_docs = [[str(i) for i in indices] for indices in relevant_indices] results[model_name] = evaluate_retrieval_quality( queries, relevant_docs, retrieved_per_query, k ) return results ``` ## Best Practices ### Do's - **Match model to use case**: Code vs prose vs multilingual - **Chunk thoughtfully**: Preserve semantic boundaries - **Normalize embeddings**: For cosine similarity search - **Batch requests**: More efficient than one-by-one - **Cache embeddings**: Avoid recomputing for static content - **Use Voyage AI for Claude apps**: Recommended by Anthropic ### Don'ts - **Don't ignore token limits**: Truncation loses information - **Don't mix embedding models**: Incompatible vector spaces - **Don't skip preprocessing**: Garbage in, garbage out - **Don't over-chunk**: Lose important context - **Don't forget metadata**: Essential for filtering and debugging ## Resources - [Voyage AI Documentation](https://docs.voyageai.com/) - [OpenAI Embeddings Guide](https://platform.openai.com/docs/guides/embeddings) - [Sentence Transformers](https://www.sbert.net/) - [MTEB Benchmar
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

hybrid-search-implementation

Combine vector and keyword search for improved retrieval. Use when

coding
โญ1
# Hybrid Search Implementation Patterns for combining vector similarity and keyword-based search. ## When to Use This Skill - Building RAG systems with improved recall - Combining semantic understanding with exact matching - Handling queries with specific terms (names, codes) - Improving search for domain-specific vocabulary - When pure vector search misses keyword matches ## Core Concepts ### 1. Hybrid Search Architecture ``` Query โ†’ โ”ฌโ”€โ–บ Vector Search โ”€โ”€โ–บ Candidates โ”€โ” โ”‚ โ”‚ โ””โ”€โ–บ Keyword Search โ”€โ–บ Candidates โ”€โ”ดโ”€โ–บ Fusion โ”€โ–บ Results ``` ### 2. Fusion Methods | Method | Description | Best For | | ----------------- | ------------------------ | --------------- | | **RRF** | Reciprocal Rank Fusion | General purpose | | **Linear** | Weighted sum of scores | Tunable balance | | **Cross-encoder** | Rerank with neural model | Highest quality | | **Cascade** | Filter then rerank | Efficiency | ## Templates ### Template 1: Reciprocal Rank Fusion ```python from typing import List, Dict, Tuple from collections import defaultdict def reciprocal_rank_fusion( result_lists: List[List[Tuple[str, float]]], k: int = 60, weights: List[float] = None ) -> List[Tuple[str, float]]: """ Combine multiple ranked lists using RRF. Args: result_lists: List of (doc_id, score) tuples per search method k: RRF constant (higher = more weight to lower ranks) weights: Optional weights per result list Returns: Fused ranking as (doc_id, score) tuples """ if weights is None: weights = [1.0] * len(result_lists) scores = defaultdict(float) for result_list, weight in zip(result_lists, weights): for rank, (doc_id, _) in enumerate(result_list): # RRF formula: 1 / (k + rank) scores[doc_id] += weight * (1.0 / (k + rank + 1)) # Sort by fused score return sorted(scores.items(), key=lambda x: x[1], reverse=True) def linear_combination( vector_results: List[Tuple[str, float]], keyword_results: List[Tuple[str, float]], alpha: float = 0.5 ) -> List[Tuple[str, float]]: """ Combine results with linear interpolation. Args: vector_results: (doc_id, similarity_score) from vector search keyword_results: (doc_id, bm25_score) from keyword search alpha: Weight for vector search (1-alpha for keyword) """ # Normalize scores to [0, 1] def normalize(results): if not results: return {} scores = [s for _, s in results] min_s, max_s = min(scores), max(scores) range_s = max_s - min_s if max_s != min_s else 1 return {doc_id: (score - min_s) / range_s for doc_id, score in results} vector_scores = normalize(vector_results) keyword_scores = normalize(keyword_results) # Combine all_docs = set(vector_scores.keys()) | set(keyword_scores.keys()) combined = {} for doc_id in all_docs: v_score = vector_scores.get(doc_id, 0) k_score = keyword_scores.get(doc_id, 0) combined[doc_id] = alpha * v_score + (1 - alpha) * k_score return sorted(combined.items(), key=lambda x: x[1], reverse=True) ``` ### Template 2: PostgreSQL Hybrid Search ```python import asyncpg from typing import List, Dict, Optional import numpy as np class PostgresHybridSearch: """Hybrid search with pgvector and full-text search.""" def __init__(self, pool: asyncpg.Pool): self.pool = pool async def setup_schema(self): """Create tables and indexes.""" async with self.pool.acquire() as conn: await conn.execute(""" CREATE EXTENSION IF NOT EXISTS vector; CREATE TABLE IF NOT EXISTS documents ( id TEXT PRIMARY KEY, content TEXT NOT NULL, embedding vector(1536), metadata JSONB DEFAULT '{}', ts_content tsvector GENERATED ALWAYS AS ( to_tsvector('english', content) ) STORED ); -- Vector index (HNSW) CREATE INDEX IF NOT EXISTS documents_embedding_idx ON documents USING hnsw (embedding vector_cosine_ops); -- Full-text index (GIN) CREATE INDEX IF NOT EXISTS documents_fts_idx ON documents USING gin (ts_content); """) async def hybrid_search( self, query: str, query_embedding: List[float], limit: int = 10, vector_weight: float = 0.5, filter_metadata: Optional[Dict] = None ) -> List[Dict]: """ Perform hybrid search combining vector and full-text. Uses RRF fusion for combining results. """ async with self.pool.acquire() as conn: # Build filter clause where_clause = "1=1" params = [query_embedding, query, limit * 3] if filter_metadata: for key, value in filter_metadata.items(): params.append(value) where_clause += f" AND metadata->>'{key}' = ${len(params)}" results = await conn.fetch(f""" WITH vector_search AS ( SELECT id, content, metadata, ROW_NUMBER() OVER (ORDER BY embedding <=> $1::vector) as vector_rank, 1 - (embedding <=> $1::vector) as vector_score FROM documents WHERE {where_clause} ORDER BY embedding <=> $1::vector LIMIT $3 ), keyword_search AS ( SELECT id, content, metadata, ROW_NUMBER() OVER (ORDER BY ts_rank(ts_content, websearch_to_tsquery('english', $2)) DESC) as keyword_rank, ts_rank(ts_content, websearch_to_tsquery('english', $2)) as keyword_score FROM documents WHERE ts_content @@ websearch_to_tsquery('english', $2) AND {where_clause} ORDER BY ts_rank(ts_content, websearch_to_tsquery('english', $2)) DESC LIMIT $3 ) SELECT COALESCE(v.id, k.id) as id, COALESCE(v.content, k.content) as content, COALESCE(v.metadata, k.metadata) as metadata, v.vector_score, k.keyword_score, -- RRF fusion COALESCE(1.0 / (60 + v.vector_rank), 0) * $4::float + COALESCE(1.0 / (60 + k.keyword_rank), 0) * (1 - $4::float) as rrf_score FROM vector_search v FULL OUTER JOIN keyword_search k ON v.id = k.id ORDER BY rrf_score DESC LIMIT $3 / 3 """, *params, vector_weight) return [dict(row) for row in results] async def search_with_rerank( self, query: str, query_embedding: List[float], limit: int = 10, rerank_candidates: int = 50 ) -> List[Dict]: """Hybrid search with cross-encoder reranking.""" from sentence_transformers import CrossEncoder # Get candidates candidates = await self.hybrid_search( query, query_embedding, limit=rerank_candidates ) if not candidates: return [] # Rerank with cross-encoder model = CrossEncoder('cross-encoder/ms-marco-MiniLM-L-6-v2') pairs = [(query, c["content"]) for c in candidates] scores = model.predict(pairs) for candidate, score in zip(candidates, scores): candidate["rerank_score"] = float(score) # Sort by rerank score and return top results reranked = sorted(candidates, key=lambda x: x["rerank_score"], reverse=True) return reranked[:limit] ``` ### Template 3: Elasticsearch Hybrid Search ```python from elasticsearch import Elasticsearch from typing import List, Dict, Optional class ElasticsearchHybridSearch: """Hybrid search with Elasticsearch and dense vectors.""" def __init__( self, es_client: Elasticsearch, index_name: str = "documents" ): self.es = es_client self.index_name = index_name def create_index(self, vector_dims: int = 1536): """Create index with dense vector and text fields.""" mapping = { "mappings": { "properties": { "content": { "type": "text", "analyzer": "english" }, "embedding": { "type": "dense_vector", "dims": vector_dims, "index": True, "similarity": "cosine" }, "metadata": { "type": "object", "enabled": True } } } } self.es.indices.create(index=self.index_name, body=mapping, ignore=400) def hybrid_search( self, query: str, query_embedding: List[float], limit: int = 10, boost_vector: float = 1.0, boost_text: float = 1.0, filter: Optional[Dict] = None ) -> List[Dict]: """ Hybrid search using Elasticsearch's built-in capabilities. """ # Build the hybrid query search_body = { "size": limit, "query": { "bool": { "should": [ # Vector search (kNN) { "script_score": { "query": {"match_all": {}}, "script": { "source": f"cosineSimilarity(params.query_vector, 'embedding') * {boost_vector} + 1.0", "params": {"query_vector": query_embedding} } } }, # Text search (BM25) { "match": { "content": { "query": query, "boost": boost_text } } } ], "minimum_should_match": 1 } } } # Add filter if provided if filter: search_body["query"]["bool"]["filter"] = filter response = self.es.search(index=self.index_name, body=search_body) return [ { "id": hit["_id"], "content": hit["_source"]["content"], "metadata": hit["_source"].get("metadata", {}), "score": hit["_score"] } for hit in response["hits"]["hits"] ] def hybrid_search_rrf( self, query: str, query_embedding: List[float], limit: int = 10, window_size: int = 100 ) -> List[Dict]: """ Hybrid search using Elasticsearch 8.x RRF. """ search_body = { "size": limit, "sub_searches": [ { "query": { "match": { "content": query } } }, { "query": { "knn": { "field": "embedding", "query_vector": query_embedding, "k": window_size, "num_candidates": window_size * 2 } } } ], "rank": { "rrf": { "window_size": window_size, "rank_constant": 60 } } } response = self.es.search(index=self.index_name, body=search_body) return [ { "id": hit["_id"], "content": hit["_source"]["content"], "score": hit["_score"] } for hit in response["hits"]["hits"] ] ``` ### Template 4: Custom Hybrid RAG Pipeline ```python from typing import List, Dict, Optional, Callable from dataclasses import dataclass @dataclass class SearchResult: id: str content: str score: float source: str # "vector", "keyword", "hybrid" metadata: Dict = None class HybridRAGPipeline: """Complete hybrid search pipeline for RAG.""" def __init__( self, vector_store, keyword_store, embedder, reranker=None, fusion_method: str = "rrf", vector_weight: float = 0.5 ): self.vector_store = vector_store self.keyword_store = keyword_store self.embedder = embedder self.reranker = reranker self.fusion_method = fusion_method self.vector_weight = vector_weight async def search( self, query: str, top_k: int = 10, filter: Optional[Dict] = None, use_rerank: bool = True ) -> List[SearchResult]: """Execute hybrid search pipeline.""" # Step 1: Get query embedding query_embedding = self.embedder.embed(query) # Step 2: Execute parallel searches vector_results, keyword_results = await asyncio.gather( self._vector_search(query_embedding, top_k * 3, filter), self._keyword_search(query, top_k * 3, filter) ) # Step 3: Fuse results if self.fusion_method == "rrf": fused = self._rrf_fusion(vector_results, keyword_results) else: fused = self._linear_fusion(vector_results, keyword_results) # Step 4: Rerank if enabled if use_rerank and self.reranker: fused = await self._rerank(query, fused[:top_k * 2]) return fused[:top_k] async def _vector_search( self, embedding: List[float], limit: int, filter: Dict ) -> List[SearchResult]: results = await self.vector_store.search(embedding, limit, filter) return [ SearchResult( id=r["id"], content=r["content"], score=r["score"], source="vector", metadata=r.get("metadata") ) for r in results ] async def _keyword_search( self, query: str, limit: int, filter: Dict ) -> List[SearchResult]: results = await self.keyword_store.search(query, limit, filter) return [ SearchResult( id=r["id"], content=r["content"], score=r["score"], source="keyword", metadata=r.get("metadata") ) for r in results ] def _rrf_fusion( self, vector_results: List[SearchResult], keyword_results: List[SearchResult] ) -> List[SearchResult]: """Fuse with RRF.""" k = 60 scores = {} content_map = {} for rank, result in enumerate(vector_results): scores[result.id] = scores.get(result.id, 0) + 1 / (k + rank + 1) content_map[result.id] = result for rank, result in enumerate(keyword_results): scores[result.id] = scores.get(result.id, 0) + 1 / (k + rank + 1) if result.id not in content_map: content_map[result.id] = result sorted_ids = sorted(scores.keys(), key=lambda x: scores[x], reverse=True) return [ SearchResult( id=doc_id, content=content_map[doc_id].content, score=scores[doc_id], source="hybrid", metadata=content_map[doc_id].metadata ) for doc_id in sorted_ids ] async def _rerank( self, query: str, results: List[SearchResult] ) -> List[SearchResult]: """Rerank with cross-encoder.""" if not results: return results pairs = [(query, r.content) for r in results] scores = self.reranker.predict(pairs) for result, score in zip(results, scores): result.score = float(score) return sorted(results, key=lambda x: x.score, reverse=True) ``` ## Best Practices ### Do's - **Tune weights empirically** - Test on your data - **Use RRF for simplicity** - Works well without tuning - **Add reranking** - Significant quality improvement - **Log both scores** - Helps with debugging - **A/B test** - Measure real user impact ### Don'ts - **Don't assume one size fits all** - Different queries need different weights - **Don't skip keyword search** - Handles exact matches better - **Don't over-fetch** - Balance recall vs latency - **Don't ignore edge cases** - Empty results, single word queries ## Resources - [RRF Paper](https://plg.uwaterloo.ca/~gvcormac/cormacksigir09-rrf.pdf) - [Vespa Hybrid Search](https://blog.vespa.ai/improving-text-ranking-with-few-shot-prompting/) - [Cohere Rerank](https://docs.cohere.com/docs/reranking)
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

similarity-search-patterns

Implement efficient similarity search with vector databases. Use

coding
โญ1
# Similarity Search Patterns Patterns for implementing efficient similarity search in production systems. ## When to Use This Skill - Building semantic search systems - Implementing RAG retrieval - Creating recommendation engines - Optimizing search latency - Scaling to millions of vectors - Combining semantic and keyword search ## Core Concepts ### 1. Distance Metrics | Metric | Formula | Best For | | ------------------ | ------------------ | --------------------- | --- | -------------- | | **Cosine** | 1 - (AยทB)/(โ€–Aโ€–โ€–Bโ€–) | Normalized embeddings | | **Euclidean (L2)** | โˆšฮฃ(a-b)ยฒ | Raw embeddings | | **Dot Product** | AยทB | Magnitude matters | | **Manhattan (L1)** | ฮฃ | a-b | | Sparse vectors | ### 2. Index Types ``` โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ” โ”‚ Index Types โ”‚ โ”œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ฌโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ฌโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ค โ”‚ Flat โ”‚ HNSW โ”‚ IVF+PQ โ”‚ โ”‚ (Exact) โ”‚ (Graph-based) โ”‚ (Quantized) โ”‚ โ”œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ผโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ผโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ค โ”‚ O(n) search โ”‚ O(log n) โ”‚ O(โˆšn) โ”‚ โ”‚ 100% recall โ”‚ ~95-99% โ”‚ ~90-95% โ”‚ โ”‚ Small data โ”‚ Medium-Large โ”‚ Very Large โ”‚ โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ดโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ดโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜ ``` ## Templates ### Template 1: Pinecone Implementation ```python from pinecone import Pinecone, ServerlessSpec from typing import List, Dict, Optional import hashlib class PineconeVectorStore: def __init__( self, api_key: str, index_name: str, dimension: int = 1536, metric: str = "cosine" ): self.pc = Pinecone(api_key=api_key) # Create index if not exists if index_name not in self.pc.list_indexes().names(): self.pc.create_index( name=index_name, dimension=dimension, metric=metric, spec=ServerlessSpec(cloud="aws", region="us-east-1") ) self.index = self.pc.Index(index_name) def upsert( self, vectors: List[Dict], namespace: str = "" ) -> int: """ Upsert vectors. vectors: [{"id": str, "values": List[float], "metadata": dict}] """ # Batch upsert batch_size = 100 total = 0 for i in range(0, len(vectors), batch_size): batch = vectors[i:i + batch_size] self.index.upsert(vectors=batch, namespace=namespace) total += len(batch) return total def search( self, query_vector: List[float], top_k: int = 10, namespace: str = "", filter: Optional[Dict] = None, include_metadata: bool = True ) -> List[Dict]: """Search for similar vectors.""" results = self.index.query( vector=query_vector, top_k=top_k, namespace=namespace, filter=filter, include_metadata=include_metadata ) return [ { "id": match.id, "score": match.score, "metadata": match.metadata } for match in results.matches ] def search_with_rerank( self, query: str, query_vector: List[float], top_k: int = 10, rerank_top_n: int = 50, namespace: str = "" ) -> List[Dict]: """Search and rerank results.""" # Over-fetch for reranking initial_results = self.search( query_vector, top_k=rerank_top_n, namespace=namespace ) # Rerank with cross-encoder or LLM reranked = self._rerank(query, initial_results) return reranked[:top_k] def _rerank(self, query: str, results: List[Dict]) -> List[Dict]: """Rerank results using cross-encoder.""" from sentence_transformers import CrossEncoder model = CrossEncoder('cross-encoder/ms-marco-MiniLM-L-6-v2') pairs = [(query, r["metadata"]["text"]) for r in results] scores = model.predict(pairs) for result, score in zip(results, scores): result["rerank_score"] = float(score) return sorted(results, key=lambda x: x["rerank_score"], reverse=True) def delete(self, ids: List[str], namespace: str = ""): """Delete vectors by ID.""" self.index.delete(ids=ids, namespace=namespace) def delete_by_filter(self, filter: Dict, namespace: str = ""): """Delete vectors matching filter.""" self.index.delete(filter=filter, namespace=namespace) ``` ### Template 2: Qdrant Implementation ```python from qdrant_client import QdrantClient from qdrant_client.http import models from typing import List, Dict, Optional class QdrantVectorStore: def __init__( self, url: str = "localhost", port: int = 6333, collection_name: str = "documents", vector_size: int = 1536 ): self.client = QdrantClient(url=url, port=port) self.collection_name = collection_name # Create collection if not exists collections = self.client.get_collections().collections if collection_name not in [c.name for c in collections]: self.client.create_collection( collection_name=collection_name, vectors_config=models.VectorParams( size=vector_size, distance=models.Distance.COSINE ), # Optional: enable quantization for memory efficiency quantization_config=models.ScalarQuantization( scalar=models.ScalarQuantizationConfig( type=models.ScalarType.INT8, quantile=0.99, always_ram=True ) ) ) def upsert(self, points: List[Dict]) -> int: """ Upsert points. points: [{"id": str/int, "vector": List[float], "payload": dict}] """ qdrant_points = [ models.PointStruct( id=p["id"], vector=p["vector"], payload=p.get("payload", {}) ) for p in points ] self.client.upsert( collection_name=self.collection_name, points=qdrant_points ) return len(points) def search( self, query_vector: List[float], limit: int = 10, filter: Optional[models.Filter] = None, score_threshold: Optional[float] = None ) -> List[Dict]: """Search for similar vectors.""" results = self.client.search( collection_name=self.collection_name, query_vector=query_vector, limit=limit, query_filter=filter, score_threshold=score_threshold ) return [ { "id": r.id, "score": r.score, "payload": r.payload } for r in results ] def search_with_filter( self, query_vector: List[float], must_conditions: List[Dict] = None, should_conditions: List[Dict] = None, must_not_conditions: List[Dict] = None, limit: int = 10 ) -> List[Dict]: """Search with complex filters.""" conditions = [] if must_conditions: conditions.extend([ models.FieldCondition( key=c["key"], match=models.MatchValue(value=c["value"]) ) for c in must_conditions ]) filter = models.Filter(must=conditions) if conditions else None return self.search(query_vector, limit=limit, filter=filter) def search_with_sparse( self, dense_vector: List[float], sparse_vector: Dict[int, float], limit: int = 10, dense_weight: float = 0.7 ) -> List[Dict]: """Hybrid search with dense and sparse vectors.""" # Requires collection with named vectors results = self.client.search( collection_name=self.collection_name, query_vector=models.NamedVector( name="dense", vector=dense_vector ), limit=limit ) return [{"id": r.id, "score": r.score, "payload": r.payload} for r in results] ``` ### Template 3: pgvector with PostgreSQL ```python import asyncpg from typing import List, Dict, Optional import numpy as np class PgVectorStore: def __init__(self, connection_string: str): self.connection_string = connection_string async def init(self): """Initialize connection pool and extension.""" self.pool = await asyncpg.create_pool(self.connection_string) async with self.pool.acquire() as conn: # Enable extension await conn.execute("CREATE EXTENSION IF NOT EXISTS vector") # Create table await conn.execute(""" CREATE TABLE IF NOT EXISTS documents ( id TEXT PRIMARY KEY, content TEXT, metadata JSONB, embedding vector(1536) ) """) # Create index (HNSW for better performance) await conn.execute(""" CREATE INDEX IF NOT EXISTS documents_embedding_idx ON documents USING hnsw (embedding vector_cosine_ops) WITH (m = 16, ef_construction = 64) """) async def upsert(self, documents: List[Dict]): """Upsert documents with embeddings.""" async with self.pool.acquire() as conn: await conn.executemany( """ INSERT INTO documents (id, content, metadata, embedding) VALUES ($1, $2, $3, $4) ON CONFLICT (id) DO UPDATE SET content = EXCLUDED.content, metadata = EXCLUDED.metadata, embedding = EXCLUDED.embedding """, [ ( doc["id"], doc["content"], doc.get("metadata", {}), np.array(doc["embedding"]).tolist() ) for doc in documents ] ) async def search( self, query_embedding: List[float], limit: int = 10, filter_metadata: Optional[Dict] = None ) -> List[Dict]: """Search for similar documents.""" query = """ SELECT id, content, metadata, 1 - (embedding <=> $1::vector) as similarity FROM documents """ params = [query_embedding] if filter_metadata: conditions = [] for key, value in filter_metadata.items(): params.append(value) conditions.append(f"metadata->>'{key}' = ${len(params)}") query += " WHERE " + " AND ".join(conditions) query += f" ORDER BY embedding <=> $1::vector LIMIT ${len(params) + 1}" params.append(limit) async with self.pool.acquire() as conn: rows = await conn.fetch(query, *params) return [ { "id": row["id"], "content": row["content"], "metadata": row["metadata"], "score": row["similarity"] } for row in rows ] async def hybrid_search( self, query_embedding: List[float], query_text: str, limit: int = 10, vector_weight: float = 0.5 ) -> List[Dict]: """Hybrid search combining vector and full-text.""" async with self.pool.acquire() as conn: rows = await conn.fetch( """ WITH vector_results AS ( SELECT id, content, metadata, 1 - (embedding <=> $1::vector) as vector_score FROM documents ORDER BY embedding <=> $1::vector LIMIT $3 * 2 ), text_results AS ( SELECT id, content, metadata, ts_rank(to_tsvector('english', content), plainto_tsquery('english', $2)) as text_score FROM documents WHERE to_tsvector('english', content) @@ plainto_tsquery('english', $2) LIMIT $3 * 2 ) SELECT COALESCE(v.id, t.id) as id, COALESCE(v.content, t.content) as content, COALESCE(v.metadata, t.metadata) as metadata, COALESCE(v.vector_score, 0) * $4 + COALESCE(t.text_score, 0) * (1 - $4) as combined_score FROM vector_results v FULL OUTER JOIN text_results t ON v.id = t.id ORDER BY combined_score DESC LIMIT $3 """, query_embedding, query_text, limit, vector_weight ) return [dict(row) for row in rows] ``` ### Template 4: Weaviate Implementation ```python import weaviate from weaviate.util import generate_uuid5 from typing import List, Dict, Optional class WeaviateVectorStore: def __init__( self, url: str = "http://localhost:8080", class_name: str = "Document" ): self.client = weaviate.Client(url=url) self.class_name = class_name self._ensure_schema() def _ensure_schema(self): """Create schema if not exists.""" schema = { "class": self.class_name, "vectorizer": "none", # We provide vectors "properties": [ {"name": "content", "dataType": ["text"]}, {"name": "source", "dataType": ["string"]}, {"name": "chunk_id", "dataType": ["int"]} ] } if not self.client.schema.exists(self.class_name): self.client.schema.create_class(schema) def upsert(self, documents: List[Dict]): """Batch upsert documents.""" with self.client.batch as batch: batch.batch_size = 100 for doc in documents: batch.add_data_object( data_object={ "content": doc["content"], "source": doc.get("source", ""), "chunk_id": doc.get("chunk_id", 0) }, class_name=self.class_name, uuid=generate_uuid5(doc["id"]), vector=doc["embedding"] ) def search( self, query_vector: List[float], limit: int = 10, where_filter: Optional[Dict] = None ) -> List[Dict]: """Vector search.""" query = ( self.client.query .get(self.class_name, ["content", "source", "chunk_id"]) .with_near_vector({"vector": query_vector}) .with_limit(limit) .with_additional(["distance", "id"]) ) if where_filter: query = query.with_where(where_filter) results = query.do() return [ { "id": item["_additional"]["id"], "content": item["content"], "source": item["source"], "score": 1 - item["_additional"]["distance"] } for item in results["data"]["Get"][self.class_name] ] def hybrid_search( self, query: str, query_vector: List[float], limit: int = 10, alpha: float = 0.5 # 0 = keyword, 1 = vector ) -> List[Dict]: """Hybrid search combining BM25 and vector.""" results = ( self.client.query .get(self.class_name, ["content", "source"]) .with_hybrid(query=query, vector=query_vector, alpha=alpha) .with_limit(limit) .with_additional(["score"]) .do() ) return [ { "content": item["content"], "source": item["source"], "score": item["_additional"]["score"] } for item in results["data"]["Get"][self.class_name] ] ``` ## Best Practices ### Do's - **Use appropriate index** - HNSW for most cases - **Tune parameters** - ef_search, nprobe for recall/speed - **Implement hybrid search** - Combine with keyword search - **Monitor recall** - Measure search quality - **Pre-filter when possible** - Reduce search space ### Don'ts - **Don't skip evaluation** - Measure before optimizing - **Don't over-index** - Start with flat, scale up - **Don't ignore latency** - P99 matters for UX - **Don't forget costs** - Vector storage adds up ## Resources - [Pinecone Docs](https://docs.pinecone.io/) - [Qdrant Docs](https://qdrant.tech/documentation/) - [pgvector](https://github.com/pgvector/pgvector) - [Weaviate Docs](https://weaviate.io/developers/weaviate)
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

vector-index-tuning

Optimize vector index performance for latency, recall, and memory.

coding
โญ1
# Vector Index Tuning Guide to optimizing vector indexes for production performance. ## When to Use This Skill - Tuning HNSW parameters - Implementing quantization - Optimizing memory usage - Reducing search latency - Balancing recall vs speed - Scaling to billions of vectors ## Core Concepts ### 1. Index Type Selection ``` Data Size Recommended Index โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ < 10K vectors โ†’ Flat (exact search) 10K - 1M โ†’ HNSW 1M - 100M โ†’ HNSW + Quantization > 100M โ†’ IVF + PQ or DiskANN ``` ### 2. HNSW Parameters | Parameter | Default | Effect | | ------------------ | ------- | ---------------------------------------------------- | | **M** | 16 | Connections per node, โ†‘ = better recall, more memory | | **efConstruction** | 100 | Build quality, โ†‘ = better index, slower build | | **efSearch** | 50 | Search quality, โ†‘ = better recall, slower search | ### 3. Quantization Types ``` Full Precision (FP32): 4 bytes ร— dimensions Half Precision (FP16): 2 bytes ร— dimensions INT8 Scalar: 1 byte ร— dimensions Product Quantization: ~32-64 bytes total Binary: dimensions/8 bytes ``` ## Templates ### Template 1: HNSW Parameter Tuning ```python import numpy as np from typing import List, Tuple import time def benchmark_hnsw_parameters( vectors: np.ndarray, queries: np.ndarray, ground_truth: np.ndarray, m_values: List[int] = [8, 16, 32, 64], ef_construction_values: List[int] = [64, 128, 256], ef_search_values: List[int] = [32, 64, 128, 256] ) -> List[dict]: """Benchmark different HNSW configurations.""" import hnswlib results = [] dim = vectors.shape[1] n = vectors.shape[0] for m in m_values: for ef_construction in ef_construction_values: # Build index index = hnswlib.Index(space='cosine', dim=dim) index.init_index(max_elements=n, M=m, ef_construction=ef_construction) build_start = time.time() index.add_items(vectors) build_time = time.time() - build_start # Get memory usage memory_bytes = index.element_count * ( dim * 4 + # Vector storage m * 2 * 4 # Graph edges (approximate) ) for ef_search in ef_search_values: index.set_ef(ef_search) # Measure search search_start = time.time() labels, distances = index.knn_query(queries, k=10) search_time = time.time() - search_start # Calculate recall recall = calculate_recall(labels, ground_truth, k=10) results.append({ "M": m, "ef_construction": ef_construction, "ef_search": ef_search, "build_time_s": build_time, "search_time_ms": search_time * 1000 / len(queries), "recall@10": recall, "memory_mb": memory_bytes / 1024 / 1024 }) return results def calculate_recall(predictions: np.ndarray, ground_truth: np.ndarray, k: int) -> float: """Calculate recall@k.""" correct = 0 for pred, truth in zip(predictions, ground_truth): correct += len(set(pred[:k]) & set(truth[:k])) return correct / (len(predictions) * k) def recommend_hnsw_params( num_vectors: int, target_recall: float = 0.95, max_latency_ms: float = 10, available_memory_gb: float = 8 ) -> dict: """Recommend HNSW parameters based on requirements.""" # Base recommendations if num_vectors < 100_000: m = 16 ef_construction = 100 elif num_vectors < 1_000_000: m = 32 ef_construction = 200 else: m = 48 ef_construction = 256 # Adjust ef_search based on recall target if target_recall >= 0.99: ef_search = 256 elif target_recall >= 0.95: ef_search = 128 else: ef_search = 64 return { "M": m, "ef_construction": ef_construction, "ef_search": ef_search, "notes": f"Estimated for {num_vectors:,} vectors, {target_recall:.0%} recall" } ``` ### Template 2: Quantization Strategies ```python import numpy as np from typing import Optional class VectorQuantizer: """Quantization strategies for vector compression.""" @staticmethod def scalar_quantize_int8( vectors: np.ndarray, min_val: Optional[float] = None, max_val: Optional[float] = None ) -> Tuple[np.ndarray, dict]: """Scalar quantization to INT8.""" if min_val is None: min_val = vectors.min() if max_val is None: max_val = vectors.max() # Scale to 0-255 range scale = 255.0 / (max_val - min_val) quantized = np.clip( np.round((vectors - min_val) * scale), 0, 255 ).astype(np.uint8) params = {"min_val": min_val, "max_val": max_val, "scale": scale} return quantized, params @staticmethod def dequantize_int8( quantized: np.ndarray, params: dict ) -> np.ndarray: """Dequantize INT8 vectors.""" return quantized.astype(np.float32) / params["scale"] + params["min_val"] @staticmethod def product_quantize( vectors: np.ndarray, n_subvectors: int = 8, n_centroids: int = 256 ) -> Tuple[np.ndarray, dict]: """Product quantization for aggressive compression.""" from sklearn.cluster import KMeans n, dim = vectors.shape assert dim % n_subvectors == 0 subvector_dim = dim // n_subvectors codebooks = [] codes = np.zeros((n, n_subvectors), dtype=np.uint8) for i in range(n_subvectors): start = i * subvector_dim end = (i + 1) * subvector_dim subvectors = vectors[:, start:end] kmeans = KMeans(n_clusters=n_centroids, random_state=42) codes[:, i] = kmeans.fit_predict(subvectors) codebooks.append(kmeans.cluster_centers_) params = { "codebooks": codebooks, "n_subvectors": n_subvectors, "subvector_dim": subvector_dim } return codes, params @staticmethod def binary_quantize(vectors: np.ndarray) -> np.ndarray: """Binary quantization (sign of each dimension).""" # Convert to binary: positive = 1, negative = 0 binary = (vectors > 0).astype(np.uint8) # Pack bits into bytes n, dim = vectors.shape packed_dim = (dim + 7) // 8 packed = np.zeros((n, packed_dim), dtype=np.uint8) for i in range(dim): byte_idx = i // 8 bit_idx = i % 8 packed[:, byte_idx] |= (binary[:, i] << bit_idx) return packed def estimate_memory_usage( num_vectors: int, dimensions: int, quantization: str = "fp32", index_type: str = "hnsw", hnsw_m: int = 16 ) -> dict: """Estimate memory usage for different configurations.""" # Vector storage bytes_per_dimension = { "fp32": 4, "fp16": 2, "int8": 1, "pq": 0.05, # Approximate "binary": 0.125 } vector_bytes = num_vectors * dimensions * bytes_per_dimension[quantization] # Index overhead if index_type == "hnsw": # Each node has ~M*2 edges, each edge is 4 bytes (int32) index_bytes = num_vectors * hnsw_m * 2 * 4 elif index_type == "ivf": # Inverted lists + centroids index_bytes = num_vectors * 8 + 65536 * dimensions * 4 else: index_bytes = 0 total_bytes = vector_bytes + index_bytes return { "vector_storage_mb": vector_bytes / 1024 / 1024, "index_overhead_mb": index_bytes / 1024 / 1024, "total_mb": total_bytes / 1024 / 1024, "total_gb": total_bytes / 1024 / 1024 / 1024 } ``` ### Template 3: Qdrant Index Configuration ```python from qdrant_client import QdrantClient from qdrant_client.http import models def create_optimized_collection( client: QdrantClient, collection_name: str, vector_size: int, num_vectors: int, optimize_for: str = "balanced" # "recall", "speed", "memory" ) -> None: """Create collection with optimized settings.""" # HNSW configuration based on optimization target hnsw_configs = { "recall": models.HnswConfigDiff(m=32, ef_construct=256), "speed": models.HnswConfigDiff(m=16, ef_construct=64), "balanced": models.HnswConfigDiff(m=16, ef_construct=128), "memory": models.HnswConfigDiff(m=8, ef_construct=64) } # Quantization configuration quantization_configs = { "recall": None, # No quantization for max recall "speed": models.ScalarQuantization( scalar=models.ScalarQuantizationConfig( type=models.ScalarType.INT8, quantile=0.99, always_ram=True ) ), "balanced": models.ScalarQuantization( scalar=models.ScalarQuantizationConfig( type=models.ScalarType.INT8, quantile=0.99, always_ram=False ) ), "memory": models.ProductQuantization( product=models.ProductQuantizationConfig( compression=models.CompressionRatio.X16, always_ram=False ) ) } # Optimizer configuration optimizer_configs = { "recall": models.OptimizersConfigDiff( indexing_threshold=10000, memmap_threshold=50000 ), "speed": models.OptimizersConfigDiff( indexing_threshold=5000, memmap_threshold=20000 ), "balanced": models.OptimizersConfigDiff( indexing_threshold=20000, memmap_threshold=50000 ), "memory": models.OptimizersConfigDiff( indexing_threshold=50000, memmap_threshold=10000 # Use disk sooner ) } client.create_collection( collection_name=collection_name, vectors_config=models.VectorParams( size=vector_size, distance=models.Distance.COSINE ), hnsw_config=hnsw_configs[optimize_for], quantization_config=quantization_configs[optimize_for], optimizers_config=optimizer_configs[optimize_for] ) def tune_search_parameters( client: QdrantClient, collection_name: str, target_recall: float = 0.95 ) -> dict: """Tune search parameters for target recall.""" # Search parameter recommendations if target_recall >= 0.99: search_params = models.SearchParams( hnsw_ef=256, exact=False, quantization=models.QuantizationSearchParams( ignore=True, # Don't use quantization for search rescore=True ) ) elif target_recall >= 0.95: search_params = models.SearchParams( hnsw_ef=128, exact=False, quantization=models.QuantizationSearchParams( ignore=False, rescore=True, oversampling=2.0 ) ) else: search_params = models.SearchParams( hnsw_ef=64, exact=False, quantization=models.QuantizationSearchParams( ignore=False, rescore=False ) ) return search_params ``` ### Template 4: Performance Monitoring ```python import time from dataclasses import dataclass from typing import List import numpy as np @dataclass class SearchMetrics: latency_p50_ms: float latency_p95_ms: float latency_p99_ms: float recall: float qps: float class VectorSearchMonitor: """Monitor vector search performance.""" def __init__(self, ground_truth_fn=None): self.latencies = [] self.recalls = [] self.ground_truth_fn = ground_truth_fn def measure_search( self, search_fn, query_vectors: np.ndarray, k: int = 10, num_iterations: int = 100 ) -> SearchMetrics: """Benchmark search performance.""" latencies = [] for _ in range(num_iterations): for query in query_vectors: start = time.perf_counter() results = search_fn(query, k=k) latency = (time.perf_counter() - start) * 1000 latencies.append(latency) latencies = np.array(latencies) total_queries = num_iterations * len(query_vectors) total_time = sum(latencies) / 1000 # seconds return SearchMetrics( latency_p50_ms=np.percentile(latencies, 50), latency_p95_ms=np.percentile(latencies, 95), latency_p99_ms=np.percentile(latencies, 99), recall=self._calculate_recall(search_fn, query_vectors, k) if self.ground_truth_fn else 0, qps=total_queries / total_time ) def _calculate_recall(self, search_fn, queries: np.ndarray, k: int) -> float: """Calculate recall against ground truth.""" if not self.ground_truth_fn: return 0 correct = 0 total = 0 for query in queries: predicted = set(search_fn(query, k=k)) actual = set(self.ground_truth_fn(query, k=k)) correct += len(predicted & actual) total += k return correct / total def profile_index_build( build_fn, vectors: np.ndarray, batch_sizes: List[int] = [1000, 10000, 50000] ) -> dict: """Profile index build performance.""" results = {} for batch_size in batch_sizes: times = [] for i in range(0, len(vectors), batch_size): batch = vectors[i:i + batch_size] start = time.perf_counter() build_fn(batch) times.append(time.perf_counter() - start) results[batch_size] = { "avg_batch_time_s": np.mean(times), "vectors_per_second": batch_size / np.mean(times) } return results ``` ## Best Practices ### Do's - **Benchmark with real queries** - Synthetic may not represent production - **Monitor recall continuously** - Can degrade with data drift - **Start with defaults** - Tune only when needed - **Use quantization** - Significant memory savings - **Consider tiered storage** - Hot/cold data separation ### Don'ts - **Don't over-optimize early** - Profile first - **Don't ignore build time** - Index updates have cost - **Don't forget reindexing** - Plan for maintenance - **Don't skip warming** - Cold indexes are slow ## Resources - [HNSW Paper](https://arxiv.org/abs/1603.09320) - [Faiss Wiki](https://github.com/facebookresearch/faiss/wiki) - [ANN Benchmarks](https://ann-benchmarks.com/)
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

ml-pipeline-workflow

Build end-to-end MLOps pipelines from data preparation through

coding
โญ1
# ML Pipeline Workflow Complete end-to-end MLOps pipeline orchestration from data preparation through model deployment. ## Overview This skill provides comprehensive guidance for building production ML pipelines that handle the full lifecycle: data ingestion โ†’ preparation โ†’ training โ†’ validation โ†’ deployment โ†’ monitoring. ## When to Use This Skill - Building new ML pipelines from scratch - Designing workflow orchestration for ML systems - Implementing data โ†’ model โ†’ deployment automation - Setting up reproducible training workflows - Creating DAG-based ML orchestration - Integrating ML components into production systems ## What This Skill Provides ### Core Capabilities 1. **Pipeline Architecture** - End-to-end workflow design - DAG orchestration patterns (Airflow, Dagster, Kubeflow) - Component dependencies and data flow - Error handling and retry strategies 2. **Data Preparation** - Data validation and quality checks - Feature engineering pipelines - Data versioning and lineage - Train/validation/test splitting strategies 3. **Model Training** - Training job orchestration - Hyperparameter management - Experiment tracking integration - Distributed training patterns 4. **Model Validation** - Validation frameworks and metrics - A/B testing infrastructure - Performance regression detection - Model comparison workflows 5. **Deployment Automation** - Model serving patterns - Canary deployments - Blue-green deployment strategies - Rollback mechanisms ### Reference Documentation See the `references/` directory for detailed guides: - **data-preparation.md** - Data cleaning, validation, and feature engineering - **model-training.md** - Training workflows and best practices - **model-validation.md** - Validation strategies and metrics - **model-deployment.md** - Deployment patterns and serving architectures ### Assets and Templates The `assets/` directory contains: - **pipeline-dag.yaml.template** - DAG template for workflow orchestration - **training-config.yaml** - Training configuration template - **validation-checklist.md** - Pre-deployment validation checklist ## Usage Patterns ### Basic Pipeline Setup ```python # 1. Define pipeline stages stages = [ "data_ingestion", "data_validation", "feature_engineering", "model_training", "model_validation", "model_deployment" ] # 2. Configure dependencies # See assets/pipeline-dag.yaml.template for full example ``` ### Production Workflow 1. **Data Preparation Phase** - Ingest raw data from sources - Run data quality checks - Apply feature transformations - Version processed datasets 2. **Training Phase** - Load versioned training data - Execute training jobs - Track experiments and metrics - Save trained models 3. **Validation Phase** - Run validation test suite - Compare against baseline - Generate performance reports - Approve for deployment 4. **Deployment Phase** - Package model artifacts - Deploy to serving infrastructure - Configure monitoring - Validate production traffic ## Best Practices ### Pipeline Design - **Modularity**: Each stage should be independently testable - **Idempotency**: Re-running stages should be safe - **Observability**: Log metrics at every stage - **Versioning**: Track data, code, and model versions - **Failure Handling**: Implement retry logic and alerting ### Data Management - Use data validation libraries (Great Expectations, TFX) - Version datasets with DVC or similar tools - Document feature engineering transformations - Maintain data lineage tracking ### Model Operations - Separate training and serving infrastructure - Use model registries (MLflow, Weights & Biases) - Implement gradual rollouts for new models - Monitor model performance drift - Maintain rollback capabilities ### Deployment Strategies - Start with shadow deployments - Use canary releases for validation - Implement A/B testing infrastructure - Set up automated rollback triggers - Monitor latency and throughput ## Integration Points ### Orchestration Tools - **Apache Airflow**: DAG-based workflow orchestration - **Dagster**: Asset-based pipeline orchestration - **Kubeflow Pipelines**: Kubernetes-native ML workflows - **Prefect**: Modern dataflow automation ### Experiment Tracking - MLflow for experiment tracking and model registry - Weights & Biases for visualization and collaboration - TensorBoard for training metrics ### Deployment Platforms - AWS SageMaker for managed ML infrastructure - Google Vertex AI for GCP deployments - Azure ML for Azure cloud - Kubernetes + KServe for cloud-agnostic serving ## Progressive Disclosure Start with the basics and gradually add complexity: 1. **Level 1**: Simple linear pipeline (data โ†’ train โ†’ deploy) 2. **Level 2**: Add validation and monitoring stages 3. **Level 3**: Implement hyperparameter tuning 4. **Level 4**: Add A/B testing and gradual rollouts 5. **Level 5**: Multi-model pipelines with ensemble strategies ## Common Patterns ### Batch Training Pipeline ```yaml # See assets/pipeline-dag.yaml.template stages: - name: data_preparation dependencies: [] - name: model_training dependencies: [data_preparation] - name: model_evaluation dependencies: [model_training] - name: model_deployment dependencies: [model_evaluation] ``` ### Real-time Feature Pipeline ```python # Stream processing for real-time features # Combined with batch training # See references/data-preparation.md ``` ### Continuous Training ```python # Automated retraining on schedule # Triggered by data drift detection # See references/model-training.md ``` ## Troubleshooting ### Common Issues - **Pipeline failures**: Check dependencies and data availability - **Training instability**: Review hyperparameters and data quality - **Deployment issues**: Validate model artifacts and serving config - **Performance degradation**: Monitor data drift and model metrics ### Debugging Steps 1. Check pipeline logs for each stage 2. Validate input/output data at boundaries 3. Test components in isolation 4. Review experiment tracking metrics 5. Inspect model artifacts and metadata ## Next Steps After setting up your pipeline: 1. Explore **hyperparameter-tuning** skill for optimization 2. Learn **experiment-tracking-setup** for MLflow/W&B 3. Review **model-deployment-patterns** for serving strategies 4. Implement monitoring with observability tools ## Related Skills - **experiment-tracking-setup**: MLflow and Weights & Biases integration - **hyperparameter-tuning**: Automated hyperparameter optimization - **model-deployment-patterns**: Advanced deployment strategies
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

sast-configuration

Configure Static Application Security Testing (SAST) tools for

security
โญ1
# SAST Configuration Static Application Security Testing (SAST) tool setup, configuration, and custom rule creation for comprehensive security scanning across multiple programming languages. ## Overview This skill provides comprehensive guidance for setting up and configuring SAST tools including Semgrep, SonarQube, and CodeQL. Use this skill when you need to: - Set up SAST scanning in CI/CD pipelines - Create custom security rules for your codebase - Configure quality gates and compliance policies - Optimize scan performance and reduce false positives - Integrate multiple SAST tools for defense-in-depth ## Core Capabilities ### 1. Semgrep Configuration - Custom rule creation with pattern matching - Language-specific security rules (Python, JavaScript, Go, Java, etc.) - CI/CD integration (GitHub Actions, GitLab CI, Jenkins) - False positive tuning and rule optimization - Organizational policy enforcement ### 2. SonarQube Setup - Quality gate configuration - Security hotspot analysis - Code coverage and technical debt tracking - Custom quality profiles for languages - Enterprise integration with LDAP/SAML ### 3. CodeQL Analysis - GitHub Advanced Security integration - Custom query development - Vulnerability variant analysis - Security research workflows - SARIF result processing ## Quick Start ### Initial Assessment 1. Identify primary programming languages in your codebase 2. Determine compliance requirements (PCI-DSS, SOC 2, etc.) 3. Choose SAST tool based on language support and integration needs 4. Review baseline scan to understand current security posture ### Basic Setup ```bash # Semgrep quick start pip install semgrep semgrep --config=auto --error # SonarQube with Docker docker run -d --name sonarqube -p 9000:9000 sonarqube:latest # CodeQL CLI setup gh extension install github/gh-codeql codeql database create mydb --language=python ``` ## Reference Documentation - [Semgrep Rule Creation](references/semgrep-rules.md) - Pattern-based security rule development - [SonarQube Configuration](references/sonarqube-config.md) - Quality gates and profiles - [CodeQL Setup Guide](references/codeql-setup.md) - Query development and workflows ## Templates & Assets - [semgrep-config.yml](assets/semgrep-config.yml) - Production-ready Semgrep configuration - [sonarqube-settings.xml](assets/sonarqube-settings.xml) - SonarQube quality profile template - [run-sast.sh](scripts/run-sast.sh) - Automated SAST execution script ## Integration Patterns ### CI/CD Pipeline Integration ```yaml # GitHub Actions example - name: Run Semgrep uses: returntocorp/semgrep-action@v1 with: config: >- p/security-audit p/owasp-top-ten ``` ### Pre-commit Hook ```bash # .pre-commit-config.yaml - repo: https://github.com/returntocorp/semgrep rev: v1.45.0 hooks: - id: semgrep args: ['--config=auto', '--error'] ``` ## Best Practices 1. **Start with Baseline** - Run initial scan to establish security baseline - Prioritize critical and high severity findings - Create remediation roadmap 2. **Incremental Adoption** - Begin with security-focused rules - Gradually add code quality rules - Implement blocking only for critical issues 3. **False Positive Management** - Document legitimate suppressions - Create allow lists for known safe patterns - Regularly review suppressed findings 4. **Performance Optimization** - Exclude test files and generated code - Use incremental scanning for large codebases - Cache scan results in CI/CD 5. **Team Enablement** - Provide security training for developers - Create internal documentation for common patterns - Establish security champions program ## Common Use Cases ### New Project Setup ```bash ./scripts/run-sast.sh --setup --language python --tools semgrep,sonarqube ``` ### Custom Rule Development ```yaml # See references/semgrep-rules.md for detailed examples rules: - id: hardcoded-jwt-secret pattern: jwt.encode($DATA, "...", ...) message: JWT secret should not be hardcoded severity: ERROR ``` ### Compliance Scanning ```bash # PCI-DSS focused scan semgrep --config p/pci-dss --json -o pci-scan-results.json ``` ## Troubleshooting ### High False Positive Rate - Review and tune rule sensitivity - Add path filters to exclude test files - Use nostmt metadata for noisy patterns - Create organization-specific rule exceptions ### Performance Issues - Enable incremental scanning - Parallelize scans across modules - Optimize rule patterns for efficiency - Cache dependencies and scan results ### Integration Failures - Verify API tokens and credentials - Check network connectivity and proxy settings - Review SARIF output format compatibility - Validate CI/CD runner permissions ## Related Skills - [OWASP Top 10 Checklist](../owasp-top10-checklist/SKILL.md) - [Container Security](../container-security/SKILL.md) - [Dependency Scanning](../dependency-scanning/SKILL.md) ## Tool Comparison | Tool | Best For | Language Support | Cost | Integration | | --------- | ------------------------ | ---------------- | --------------- | ------------- | | Semgrep | Custom rules, fast scans | 30+ languages | Free/Enterprise | Excellent | | SonarQube | Code quality + security | 25+ languages | Free/Commercial | Good | | CodeQL | Deep analysis, research | 10+ languages | Free (OSS) | GitHub native | ## Next Steps 1. Complete initial SAST tool setup 2. Run baseline security scan 3. Create custom rules for organization-specific patterns 4. Integrate into CI/CD pipeline 5. Establish security gate policies 6. Train development team on findings and remediation
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered