Skip to main content
EVOKORE// BROWSE
>

./browse/prompts

29 NODES
πŸ€–system promptβ€’7 months ago

data-storytelling

Transform data into compelling narratives using visualization,

data
⭐1
# Data Storytelling Transform raw data into compelling narratives that drive decisions and inspire action. ## When to Use This Skill - Presenting analytics to executives - Creating quarterly business reviews - Building investor presentations - Writing data-driven reports - Communicating insights to non-technical audiences - Making recommendations based on data ## Core Concepts ### 1. Story Structure ``` Setup β†’ Conflict β†’ Resolution Setup: Context and baseline Conflict: The problem or opportunity Resolution: Insights and recommendations ``` ### 2. Narrative Arc ``` 1. Hook: Grab attention with surprising insight 2. Context: Establish the baseline 3. Rising Action: Build through data points 4. Climax: The key insight 5. Resolution: Recommendations 6. Call to Action: Next steps ``` ### 3. Three Pillars | Pillar | Purpose | Components | | ------------- | -------- | -------------------------------- | | **Data** | Evidence | Numbers, trends, comparisons | | **Narrative** | Meaning | Context, causation, implications | | **Visuals** | Clarity | Charts, diagrams, highlights | ## Story Frameworks ### Framework 1: The Problem-Solution Story ```markdown # Customer Churn Analysis ## The Hook "We're losing $2.4M annually to preventable churn." ## The Context - Current churn rate: 8.5% (industry average: 5%) - Average customer lifetime value: $4,800 - 500 customers churned last quarter ## The Problem Analysis of churned customers reveals a pattern: - 73% churned within first 90 days - Common factor: < 3 support interactions - Low feature adoption in first month ## The Insight [Show engagement curve visualization] Customers who don't engage in the first 14 days are 4x more likely to churn. ## The Solution 1. Implement 14-day onboarding sequence 2. Proactive outreach at day 7 3. Feature adoption tracking ## Expected Impact - Reduce early churn by 40% - Save $960K annually - Payback period: 3 months ## Call to Action Approve $50K budget for onboarding automation. ``` ### Framework 2: The Trend Story ```markdown # Q4 Performance Analysis ## Where We Started Q3 ended with $1.2M MRR, 15% below target. Team morale was low after missed goals. ## What Changed [Timeline visualization] - Oct: Launched self-serve pricing - Nov: Reduced friction in signup - Dec: Added customer success calls ## The Transformation [Before/after comparison chart] | Metric | Q3 | Q4 | Change | |----------------|--------|--------|--------| | Trial β†’ Paid | 8% | 15% | +87% | | Time to Value | 14 days| 5 days | -64% | | Expansion Rate | 2% | 8% | +300% | ## Key Insight Self-serve + high-touch creates compound growth. Customers who self-serve AND get a success call have 3x higher expansion rate. ## Going Forward Double down on hybrid model. Target: $1.8M MRR by Q2. ``` ### Framework 3: The Comparison Story ```markdown # Market Opportunity Analysis ## The Question Should we expand into EMEA or APAC first? ## The Comparison [Side-by-side market analysis] ### EMEA - Market size: $4.2B - Growth rate: 8% - Competition: High - Regulatory: Complex (GDPR) - Language: Multiple ### APAC - Market size: $3.8B - Growth rate: 15% - Competition: Moderate - Regulatory: Varied - Language: Multiple ## The Analysis [Weighted scoring matrix visualization] | Factor | Weight | EMEA Score | APAC Score | | ----------- | ------ | ---------- | ---------- | | Market Size | 25% | 5 | 4 | | Growth | 30% | 3 | 5 | | Competition | 20% | 2 | 4 | | Ease | 25% | 2 | 3 | | **Total** | | **2.9** | **4.1** | ## The Recommendation APAC first. Higher growth, less competition. Start with Singapore hub (English, business-friendly). Enter EMEA in Year 2 with localization ready. ## Risk Mitigation - Timezone coverage: Hire 24/7 support - Cultural fit: Local partnerships - Payment: Multi-currency from day 1 ``` ## Visualization Techniques ### Technique 1: Progressive Reveal ```markdown Start simple, add layers: Slide 1: "Revenue is growing" [single line chart] Slide 2: "But growth is slowing" [add growth rate overlay] Slide 3: "Driven by one segment" [add segment breakdown] Slide 4: "Which is saturating" [add market share] Slide 5: "We need new segments" [add opportunity zones] ``` ### Technique 2: Contrast and Compare ```markdown Before/After: β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ BEFORE β”‚ AFTER β”‚ β”‚ β”‚ β”‚ β”‚ Process: 5 daysβ”‚ Process: 1 day β”‚ β”‚ Errors: 15% β”‚ Errors: 2% β”‚ β”‚ Cost: $50/unit β”‚ Cost: $20/unit β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ This/That (emphasize difference): β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ CUSTOMER A vs B β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ β”‚ β”‚ β–ˆβ–ˆ β”‚ β”‚ β”‚ β”‚ $45,000 β”‚ β”‚ $8,000 β”‚ β”‚ β”‚ β”‚ LTV β”‚ β”‚ LTV β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ Onboarded No onboarding β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### Technique 3: Annotation and Highlight ```python import matplotlib.pyplot as plt import pandas as pd fig, ax = plt.subplots(figsize=(12, 6)) # Plot the main data ax.plot(dates, revenue, linewidth=2, color='#2E86AB') # Add annotation for key events ax.annotate( 'Product Launch\n+32% spike', xy=(launch_date, launch_revenue), xytext=(launch_date, launch_revenue * 1.2), fontsize=10, arrowprops=dict(arrowstyle='->', color='#E63946'), color='#E63946' ) # Highlight a region ax.axvspan(growth_start, growth_end, alpha=0.2, color='green', label='Growth Period') # Add threshold line ax.axhline(y=target, color='gray', linestyle='--', label=f'Target: ${target:,.0f}') ax.set_title('Revenue Growth Story', fontsize=14, fontweight='bold') ax.legend() ``` ## Presentation Templates ### Template 1: Executive Summary Slide ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ KEY INSIGHT β”‚ β”‚ ══════════════════════════════════════════════════════════│ β”‚ β”‚ β”‚ "Customers who complete onboarding in week 1 β”‚ β”‚ have 3x higher lifetime value" β”‚ β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ β”‚ β”‚ β”‚ THE DATA β”‚ THE IMPLICATION β”‚ β”‚ β”‚ β”‚ β”‚ Week 1 completers: β”‚ βœ“ Prioritize onboarding UX β”‚ β”‚ β€’ LTV: $4,500 β”‚ βœ“ Add day-1 success milestones β”‚ β”‚ β€’ Retention: 85% β”‚ βœ“ Proactive week-1 outreach β”‚ β”‚ β€’ NPS: 72 β”‚ β”‚ β”‚ β”‚ Investment: $75K β”‚ β”‚ Others: β”‚ Expected ROI: 8x β”‚ β”‚ β€’ LTV: $1,500 β”‚ β”‚ β”‚ β€’ Retention: 45% β”‚ β”‚ β”‚ β€’ NPS: 34 β”‚ β”‚ β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### Template 2: Data Story Flow ``` Slide 1: THE HEADLINE "We can grow 40% faster by fixing onboarding" Slide 2: THE CONTEXT Current state metrics Industry benchmarks Gap analysis Slide 3: THE DISCOVERY What the data revealed Surprising finding Pattern identification Slide 4: THE DEEP DIVE Root cause analysis Segment breakdowns Statistical significance Slide 5: THE RECOMMENDATION Proposed actions Resource requirements Timeline Slide 6: THE IMPACT Expected outcomes ROI calculation Risk assessment Slide 7: THE ASK Specific request Decision needed Next steps ``` ### Template 3: One-Page Dashboard Story ```markdown # Monthly Business Review: January 2024 ## THE HEADLINE Revenue up 15% but CAC increasing faster than LTV ## KEY METRICS AT A GLANCE β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ MRR β”‚ NRR β”‚ CAC β”‚ LTV β”‚ β”‚ $125K β”‚ 108% β”‚ $450 β”‚ $2,200 β”‚ β”‚ β–²15% β”‚ β–²3% β”‚ β–²22% β”‚ β–²8% β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ## WHAT'S WORKING βœ“ Enterprise segment growing 25% MoM βœ“ Referral program driving 30% of new logos βœ“ Support satisfaction at all-time high (94%) ## WHAT NEEDS ATTENTION βœ— SMB acquisition cost up 40% βœ— Trial conversion down 5 points βœ— Time-to-value increased by 3 days ## ROOT CAUSE [Mini chart showing SMB vs Enterprise CAC trend] SMB paid ads becoming less efficient. CPC up 35% while conversion flat. ## RECOMMENDATION 1. Shift $20K/mo from paid to content 2. Launch SMB self-serve trial 3. A/B test shorter onboarding ## NEXT MONTH'S FOCUS - Launch content marketing pilot - Complete self-serve MVP - Reduce time-to-value to < 7 days ``` ## Writing Techniques ### Headlines That Work ```markdown BAD: "Q4 Sales Analysis" GOOD: "Q4 Sales Beat Target by 23% - Here's Why" BAD: "Customer Churn Report" GOOD: "We're Losing $2.4M to Preventable Churn" BAD: "Marketing Performance" GOOD: "Content Marketing Delivers 4x ROI vs. Paid" Formula: [Specific Number] + [Business Impact] + [Actionable Context] ``` ### Transition Phrases ```markdown Building the narrative: β€’ "This leads us to ask..." β€’ "When we dig deeper..." β€’ "The pattern becomes clear when..." β€’ "Contrast this with..." Introducing insights: β€’ "The data reveals..." β€’ "What surprised us was..." β€’ "The inflection point came when..." β€’ "The key finding is..." Moving to action: β€’ "This insight suggests..." β€’ "Based on this analysis..." β€’ "The implication is clear..." β€’ "Our recommendation is..." ``` ### Handling Uncertainty ```markdown Acknowledge limitations: β€’ "With 95% confidence, we can say..." β€’ "The sample size of 500 shows..." β€’ "While correlation is strong, causation requires..." β€’ "This trend holds for [segment], though [caveat]..." Present ranges: β€’ "Impact estimate: $400K-$600K" β€’ "Confidence interval: 15-20% improvement" β€’ "Best case: X, Conservative: Y" ``` ## Best Practices ### Do's - **Start with the "so what"** - Lead with insight - **Use the rule of three** - Three points, three comparisons - **Show, don't tell** - Let data speak - **Make it personal** - Connect to audience goals - **End with action** - Clear next steps ### Don'ts - **Don't data dump** - Curate ruthlessly - **Don't bury the insight** - Front-load key findings - **Don't use jargon** - Match audience vocabulary - **Don't show methodology first** - Context, then method - **Don't forget the narrative** - Numbers need meaning ## Resources - [Storytelling with Data (Cole Nussbaumer)](https://www.storytellingwithdata.com/) - [The Pyramid Principle (Barbara Minto)](https://www.amazon.com/Pyramid-Principle-Logic-Writing-Thinking/dp/0273710516) - [Resonate (Nancy Duarte)](https://www.duarte.com/resonate/)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

kpi-dashboard-design

Design effective KPI dashboards with metrics selection,

coding
⭐1
# KPI Dashboard Design Comprehensive patterns for designing effective Key Performance Indicator (KPI) dashboards that drive business decisions. ## When to Use This Skill - Designing executive dashboards - Selecting meaningful KPIs - Building real-time monitoring displays - Creating department-specific metrics views - Improving existing dashboard layouts - Establishing metric governance ## Core Concepts ### 1. KPI Framework | Level | Focus | Update Frequency | Audience | | --------------- | ---------------- | ----------------- | ---------- | | **Strategic** | Long-term goals | Monthly/Quarterly | Executives | | **Tactical** | Department goals | Weekly/Monthly | Managers | | **Operational** | Day-to-day | Real-time/Daily | Teams | ### 2. SMART KPIs ``` Specific: Clear definition Measurable: Quantifiable Achievable: Realistic targets Relevant: Aligned to goals Time-bound: Defined period ``` ### 3. Dashboard Hierarchy ``` β”œβ”€β”€ Executive Summary (1 page) β”‚ β”œβ”€β”€ 4-6 headline KPIs β”‚ β”œβ”€β”€ Trend indicators β”‚ └── Key alerts β”œβ”€β”€ Department Views β”‚ β”œβ”€β”€ Sales Dashboard β”‚ β”œβ”€β”€ Marketing Dashboard β”‚ β”œβ”€β”€ Operations Dashboard β”‚ └── Finance Dashboard └── Detailed Drilldowns β”œβ”€β”€ Individual metrics └── Root cause analysis ``` ## Common KPIs by Department ### Sales KPIs ```yaml Revenue Metrics: - Monthly Recurring Revenue (MRR) - Annual Recurring Revenue (ARR) - Average Revenue Per User (ARPU) - Revenue Growth Rate Pipeline Metrics: - Sales Pipeline Value - Win Rate - Average Deal Size - Sales Cycle Length Activity Metrics: - Calls/Emails per Rep - Demos Scheduled - Proposals Sent - Close Rate ``` ### Marketing KPIs ```yaml Acquisition: - Cost Per Acquisition (CPA) - Customer Acquisition Cost (CAC) - Lead Volume - Marketing Qualified Leads (MQL) Engagement: - Website Traffic - Conversion Rate - Email Open/Click Rate - Social Engagement ROI: - Marketing ROI - Campaign Performance - Channel Attribution - CAC Payback Period ``` ### Product KPIs ```yaml Usage: - Daily/Monthly Active Users (DAU/MAU) - Session Duration - Feature Adoption Rate - Stickiness (DAU/MAU) Quality: - Net Promoter Score (NPS) - Customer Satisfaction (CSAT) - Bug/Issue Count - Time to Resolution Growth: - User Growth Rate - Activation Rate - Retention Rate - Churn Rate ``` ### Finance KPIs ```yaml Profitability: - Gross Margin - Net Profit Margin - EBITDA - Operating Margin Liquidity: - Current Ratio - Quick Ratio - Cash Flow - Working Capital Efficiency: - Revenue per Employee - Operating Expense Ratio - Days Sales Outstanding - Inventory Turnover ``` ## Dashboard Layout Patterns ### Pattern 1: Executive Summary ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ EXECUTIVE DASHBOARD [Date Range β–Ό] β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ REVENUE β”‚ PROFIT β”‚ CUSTOMERS β”‚ NPS SCORE β”‚ β”‚ $2.4M β”‚ $450K β”‚ 12,450 β”‚ 72 β”‚ β”‚ β–² 12% β”‚ β–² 8% β”‚ β–² 15% β”‚ β–² 5pts β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ β”‚ β”‚ Revenue Trend β”‚ Revenue by Product β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ /\ /\ β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 45% β”‚ β”‚ β”‚ β”‚ / \ / \ /\ β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 32% β”‚ β”‚ β”‚ β”‚ / \/ \ / \ β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆβ–ˆ 18% β”‚ β”‚ β”‚ β”‚ / \/ \ β”‚ β”‚ β”‚ β–ˆβ–ˆ 5% β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ πŸ”΄ Alert: Churn rate exceeded threshold (>5%) β”‚ β”‚ 🟑 Warning: Support ticket volume 20% above average β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### Pattern 2: SaaS Metrics Dashboard ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ SAAS METRICS Jan 2024 [Monthly β–Ό] β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ MRR GROWTH β”‚ β”‚ β”‚ MRR β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ $125,000 β”‚ β”‚ β”‚ /── β”‚ β”‚ β”‚ β”‚ β–² 8% β”‚ β”‚ β”‚ /────/ β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ /────/ β”‚ β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ /────/ β”‚ β”‚ β”‚ β”‚ ARR β”‚ β”‚ β”‚ /────/ β”‚ β”‚ β”‚ β”‚ $1,500,000 β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”‚ β–² 15% β”‚ β”‚ J F M A M J J A S O N D β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ UNIT ECONOMICS β”‚ COHORT RETENTION β”‚ β”‚ β”‚ β”‚ β”‚ CAC: $450 β”‚ Month 1: β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 100% β”‚ β”‚ LTV: $2,700 β”‚ Month 3: β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 85% β”‚ β”‚ LTV/CAC: 6.0x β”‚ Month 6: β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 80% β”‚ β”‚ β”‚ Month 12: β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆ 72% β”‚ β”‚ Payback: 4 months β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ CHURN ANALYSIS β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ Gross β”‚ Net β”‚ Logo β”‚ Expansion β”‚ β”‚ β”‚ β”‚ 4.2% β”‚ 1.8% β”‚ 3.1% β”‚ 2.4% β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### Pattern 3: Real-time Operations ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ OPERATIONS CENTER Live ● Last: 10:42:15 β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ SYSTEM HEALTH β”‚ SERVICE STATUS β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ β”‚ CPU MEM DISK β”‚ β”‚ ● API Gateway Healthy β”‚ β”‚ β”‚ 45% 72% 58% β”‚ β”‚ ● User Service Healthy β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆ β”‚ β”‚ ● Payment Service Degraded β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆ β”‚ β”‚ ● Database Healthy β”‚ β”‚ β”‚ β–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆβ–ˆ β–ˆβ–ˆβ–ˆ β”‚ β”‚ ● Cache Healthy β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ REQUEST THROUGHPUT β”‚ ERROR RATE β”‚ β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ β”‚ β”‚ β–β–‚β–ƒβ–„β–…β–†β–‡β–ˆβ–‡β–†β–…β–„β–ƒβ–‚β–β–‚β–ƒβ–„β–… β”‚ β”‚ β”‚ ▁▁▁▁▁▂▁▁▁▁▁▁▁▁▁▁▁▁▁▁ β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ Current: 12,450 req/s β”‚ Current: 0.02% β”‚ β”‚ Peak: 18,200 req/s β”‚ Threshold: 1.0% β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ RECENT ALERTS β”‚ β”‚ 10:40 🟑 High latency on payment-service (p99 > 500ms) β”‚ β”‚ 10:35 🟒 Resolved: Database connection pool recovered β”‚ β”‚ 10:22 πŸ”΄ Payment service circuit breaker tripped β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ## Implementation Patterns ### SQL for KPI Calculations ```sql -- Monthly Recurring Revenue (MRR) WITH mrr_calculation AS ( SELECT DATE_TRUNC('month', billing_date) AS month, SUM( CASE subscription_interval WHEN 'monthly' THEN amount WHEN 'yearly' THEN amount / 12 WHEN 'quarterly' THEN amount / 3 END ) AS mrr FROM subscriptions WHERE status = 'active' GROUP BY DATE_TRUNC('month', billing_date) ) SELECT month, mrr, LAG(mrr) OVER (ORDER BY month) AS prev_mrr, (mrr - LAG(mrr) OVER (ORDER BY month)) / LAG(mrr) OVER (ORDER BY month) * 100 AS growth_pct FROM mrr_calculation; -- Cohort Retention WITH cohorts AS ( SELECT user_id, DATE_TRUNC('month', created_at) AS cohort_month FROM users ), activity AS ( SELECT user_id, DATE_TRUNC('month', event_date) AS activity_month FROM user_events WHERE event_type = 'active_session' ) SELECT c.cohort_month, EXTRACT(MONTH FROM age(a.activity_month, c.cohort_month)) AS months_since_signup, COUNT(DISTINCT a.user_id) AS active_users, COUNT(DISTINCT a.user_id)::FLOAT / COUNT(DISTINCT c.user_id) * 100 AS retention_rate FROM cohorts c LEFT JOIN activity a ON c.user_id = a.user_id AND a.activity_month >= c.cohort_month GROUP BY c.cohort_month, EXTRACT(MONTH FROM age(a.activity_month, c.cohort_month)) ORDER BY c.cohort_month, months_since_signup; -- Customer Acquisition Cost (CAC) SELECT DATE_TRUNC('month', acquired_date) AS month, SUM(marketing_spend) / NULLIF(COUNT(new_customers), 0) AS cac, SUM(marketing_spend) AS total_spend, COUNT(new_customers) AS customers_acquired FROM ( SELECT DATE_TRUNC('month', u.created_at) AS acquired_date, u.id AS new_customers, m.spend AS marketing_spend FROM users u JOIN marketing_spend m ON DATE_TRUNC('month', u.created_at) = m.month WHERE u.source = 'marketing' ) acquisition GROUP BY DATE_TRUNC('month', acquired_date); ``` ### Python Dashboard Code (Streamlit) ```python import streamlit as st import pandas as pd import plotly.express as px import plotly.graph_objects as go st.set_page_config(page_title="KPI Dashboard", layout="wide") # Header with date filter col1, col2 = st.columns([3, 1]) with col1: st.title("Executive Dashboard") with col2: date_range = st.selectbox( "Period", ["Last 7 Days", "Last 30 Days", "Last Quarter", "YTD"] ) # KPI Cards def metric_card(label, value, delta, prefix="", suffix=""): delta_color = "green" if delta >= 0 else "red" delta_arrow = "β–²" if delta >= 0 else "β–Ό" st.metric( label=label, value=f"{prefix}{value:,.0f}{suffix}", delta=f"{delta_arrow} {abs(delta):.1f}%" ) col1, col2, col3, col4 = st.columns(4) with col1: metric_card("Revenue", 2400000, 12.5, prefix="$") with col2: metric_card("Customers", 12450, 15.2) with col3: metric_card("NPS Score", 72, 5.0) with col4: metric_card("Churn Rate", 4.2, -0.8, suffix="%") # Charts col1, col2 = st.columns(2) with col1: st.subheader("Revenue Trend") revenue_data = pd.DataFrame({ 'Month': pd.date_range('2024-01-01', periods=12, freq='M'), 'Revenue': [180000, 195000, 210000, 225000, 240000, 255000, 270000, 285000, 300000, 315000, 330000, 345000] }) fig = px.line(revenue_data, x='Month', y='Revenue', line_shape='spline', markers=True) fig.update_layout(height=300) st.plotly_chart(fig, use_container_width=True) with col2: st.subheader("Revenue by Product") product_data = pd.DataFrame({ 'Product': ['Enterprise', 'Professional', 'Starter', 'Other'], 'Revenue': [45, 32, 18, 5] }) fig = px.pie(product_data, values='Revenue', names='Product', hole=0.4) fig.update_layout(height=300) st.plotly_chart(fig, use_container_width=True) # Cohort Heatmap st.subheader("Cohort Retention") cohort_data = pd.DataFrame({ 'Cohort': ['Jan', 'Feb', 'Mar', 'Apr', 'May'], 'M0': [100, 100, 100, 100, 100], 'M1': [85, 87, 84, 86, 88], 'M2': [78, 80, 76, 79, None], 'M3': [72, 74, 70, None, None], 'M4': [68, 70, None, None, None], }) fig = go.Figure(data=go.Heatmap( z=cohort_data.iloc[:, 1:].values, x=['M0', 'M1', 'M2', 'M3', 'M4'], y=cohort_data['Cohort'], colorscale='Blues', text=cohort_data.iloc[:, 1:].values, texttemplate='%{text}%', textfont={"size": 12}, )) fig.update_layout(height=250) st.plotly_chart(fig, use_container_width=True) # Alerts Section st.subheader("Alerts") alerts = [ {"level": "error", "message": "Churn rate exceeded threshold (>5%)"}, {"level": "warning", "message": "Support ticket volume 20% above average"}, ] for alert in alerts: if alert["level"] == "error": st.error(f"πŸ”΄ {alert['message']}") elif alert["level"] == "warning": st.warning(f"🟑 {alert['message']}") ``` ## Best Practices ### Do's - **Limit to 5-7 KPIs** - Focus on what matters - **Show context** - Comparisons, trends, targets - **Use consistent colors** - Red=bad, green=good - **Enable drilldown** - From summary to detail - **Update appropriately** - Match metric frequency ### Don'ts - **Don't show vanity metrics** - Focus on actionable data - **Don't overcrowd** - White space aids comprehension - **Don't use 3D charts** - They distort perception - **Don't hide methodology** - Document calculations - **Don't ignore mobile** - Ensure responsive design ## Resources - [Stephen Few's Dashboard Design](https://www.perceptualedge.com/articles/visual_business_intelligence/rules_for_using_color.pdf) - [Edward Tufte's Principles](https://www.edwardtufte.com/tufte/) - [Google Data Studio Gallery](https://datastudio.google.com/gallery)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

deployment-pipeline-design

Design multi-stage CI/CD pipelines with approval gates, security

coding
⭐1
# Deployment Pipeline Design Architecture patterns for multi-stage CI/CD pipelines with approval gates and deployment strategies. ## Purpose Design robust, secure deployment pipelines that balance speed with safety through proper stage organization and approval workflows. ## When to Use - Design CI/CD architecture - Implement deployment gates - Configure multi-environment pipelines - Establish deployment best practices - Implement progressive delivery ## Pipeline Stages ### Standard Pipeline Flow ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Build β”‚ β†’ β”‚ Test β”‚ β†’ β”‚ Staging β”‚ β†’ β”‚ Approveβ”‚ β†’ β”‚Productionβ”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### Detailed Stage Breakdown 1. **Source** - Code checkout 2. **Build** - Compile, package, containerize 3. **Test** - Unit, integration, security scans 4. **Staging Deploy** - Deploy to staging environment 5. **Integration Tests** - E2E, smoke tests 6. **Approval Gate** - Manual approval required 7. **Production Deploy** - Canary, blue-green, rolling 8. **Verification** - Health checks, monitoring 9. **Rollback** - Automated rollback on failure ## Approval Gate Patterns ### Pattern 1: Manual Approval ```yaml # GitHub Actions production-deploy: needs: staging-deploy environment: name: production url: https://app.example.com runs-on: ubuntu-latest steps: - name: Deploy to production run: | # Deployment commands ``` ### Pattern 2: Time-Based Approval ```yaml # GitLab CI deploy:production: stage: deploy script: - deploy.sh production environment: name: production when: delayed start_in: 30 minutes only: - main ``` ### Pattern 3: Multi-Approver ```yaml # Azure Pipelines stages: - stage: Production dependsOn: Staging jobs: - deployment: Deploy environment: name: production resourceType: Kubernetes strategy: runOnce: preDeploy: steps: - task: ManualValidation@0 inputs: notifyUsers: "team-leads@example.com" instructions: "Review staging metrics before approving" ``` **Reference:** See `assets/approval-gate-template.yml` ## Deployment Strategies ### 1. Rolling Deployment ```yaml apiVersion: apps/v1 kind: Deployment metadata: name: my-app spec: replicas: 10 strategy: type: RollingUpdate rollingUpdate: maxSurge: 2 maxUnavailable: 1 ``` **Characteristics:** - Gradual rollout - Zero downtime - Easy rollback - Best for most applications ### 2. Blue-Green Deployment ```yaml # Blue (current) kubectl apply -f blue-deployment.yaml kubectl label service my-app version=blue # Green (new) kubectl apply -f green-deployment.yaml # Test green environment kubectl label service my-app version=green # Rollback if needed kubectl label service my-app version=blue ``` **Characteristics:** - Instant switchover - Easy rollback - Doubles infrastructure cost temporarily - Good for high-risk deployments ### 3. Canary Deployment ```yaml apiVersion: argoproj.io/v1alpha1 kind: Rollout metadata: name: my-app spec: replicas: 10 strategy: canary: steps: - setWeight: 10 - pause: { duration: 5m } - setWeight: 25 - pause: { duration: 5m } - setWeight: 50 - pause: { duration: 5m } - setWeight: 100 ``` **Characteristics:** - Gradual traffic shift - Risk mitigation - Real user testing - Requires service mesh or similar ### 4. Feature Flags ```python from flagsmith import Flagsmith flagsmith = Flagsmith(environment_key="API_KEY") if flagsmith.has_feature("new_checkout_flow"): # New code path process_checkout_v2() else: # Existing code path process_checkout_v1() ``` **Characteristics:** - Deploy without releasing - A/B testing - Instant rollback - Granular control ## Pipeline Orchestration ### Multi-Stage Pipeline Example ```yaml name: Production Pipeline on: push: branches: [main] jobs: build: runs-on: ubuntu-latest steps: - uses: actions/checkout@v4 - name: Build application run: make build - name: Build Docker image run: docker build -t myapp:${{ github.sha }} . - name: Push to registry run: docker push myapp:${{ github.sha }} test: needs: build runs-on: ubuntu-latest steps: - name: Unit tests run: make test - name: Security scan run: trivy image myapp:${{ github.sha }} deploy-staging: needs: test runs-on: ubuntu-latest environment: name: staging steps: - name: Deploy to staging run: kubectl apply -f k8s/staging/ integration-test: needs: deploy-staging runs-on: ubuntu-latest steps: - name: Run E2E tests run: npm run test:e2e deploy-production: needs: integration-test runs-on: ubuntu-latest environment: name: production steps: - name: Canary deployment run: | kubectl apply -f k8s/production/ kubectl argo rollouts promote my-app verify: needs: deploy-production runs-on: ubuntu-latest steps: - name: Health check run: curl -f https://app.example.com/health - name: Notify team run: | curl -X POST ${{ secrets.SLACK_WEBHOOK }} \ -d '{"text":"Production deployment successful!"}' ``` ## Pipeline Best Practices 1. **Fail fast** - Run quick tests first 2. **Parallel execution** - Run independent jobs concurrently 3. **Caching** - Cache dependencies between runs 4. **Artifact management** - Store build artifacts 5. **Environment parity** - Keep environments consistent 6. **Secrets management** - Use secret stores (Vault, etc.) 7. **Deployment windows** - Schedule deployments appropriately 8. **Monitoring integration** - Track deployment metrics 9. **Rollback automation** - Auto-rollback on failures 10. **Documentation** - Document pipeline stages ## Rollback Strategies ### Automated Rollback ```yaml deploy-and-verify: steps: - name: Deploy new version run: kubectl apply -f k8s/ - name: Wait for rollout run: kubectl rollout status deployment/my-app - name: Health check id: health run: | for i in {1..10}; do if curl -sf https://app.example.com/health; then exit 0 fi sleep 10 done exit 1 - name: Rollback on failure if: failure() run: kubectl rollout undo deployment/my-app ``` ### Manual Rollback ```bash # List revision history kubectl rollout history deployment/my-app # Rollback to previous version kubectl rollout undo deployment/my-app # Rollback to specific revision kubectl rollout undo deployment/my-app --to-revision=3 ``` ## Monitoring and Metrics ### Key Pipeline Metrics - **Deployment Frequency** - How often deployments occur - **Lead Time** - Time from commit to production - **Change Failure Rate** - Percentage of failed deployments - **Mean Time to Recovery (MTTR)** - Time to recover from failure - **Pipeline Success Rate** - Percentage of successful runs - **Average Pipeline Duration** - Time to complete pipeline ### Integration with Monitoring ```yaml - name: Post-deployment verification run: | # Wait for metrics stabilization sleep 60 # Check error rate ERROR_RATE=$(curl -s "$PROMETHEUS_URL/api/v1/query?query=rate(http_errors_total[5m])" | jq '.data.result[0].value[1]') if (( $(echo "$ERROR_RATE > 0.01" | bc -l) )); then echo "Error rate too high: $ERROR_RATE" exit 1 fi ``` ## Reference Files - `references/pipeline-orchestration.md` - Complex pipeline patterns - `assets/approval-gate-template.yml` - Approval workflow templates ## Related Skills - `github-actions-templates` - For GitHub Actions implementation - `gitlab-ci-patterns` - For GitLab CI implementation - `secrets-management` - For secrets handling
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

cost-optimization

Optimize cloud costs through resource rightsizing, tagging

architecture
⭐1
# Cloud Cost Optimization Strategies and patterns for optimizing cloud costs across AWS, Azure, and GCP. ## Purpose Implement systematic cost optimization strategies to reduce cloud spending while maintaining performance and reliability. ## When to Use - Reduce cloud spending - Right-size resources - Implement cost governance - Optimize multi-cloud costs - Meet budget constraints ## Cost Optimization Framework ### 1. Visibility - Implement cost allocation tags - Use cloud cost management tools - Set up budget alerts - Create cost dashboards ### 2. Right-Sizing - Analyze resource utilization - Downsize over-provisioned resources - Use auto-scaling - Remove idle resources ### 3. Pricing Models - Use reserved capacity - Leverage spot/preemptible instances - Implement savings plans - Use committed use discounts ### 4. Architecture Optimization - Use managed services - Implement caching - Optimize data transfer - Use lifecycle policies ## AWS Cost Optimization ### Reserved Instances ``` Savings: 30-72% vs On-Demand Term: 1 or 3 years Payment: All/Partial/No upfront Flexibility: Standard or Convertible ``` ### Savings Plans ``` Compute Savings Plans: 66% savings EC2 Instance Savings Plans: 72% savings Applies to: EC2, Fargate, Lambda Flexible across: Instance families, regions, OS ``` ### Spot Instances ``` Savings: Up to 90% vs On-Demand Best for: Batch jobs, CI/CD, stateless workloads Risk: 2-minute interruption notice Strategy: Mix with On-Demand for resilience ``` ### S3 Cost Optimization ```hcl resource "aws_s3_bucket_lifecycle_configuration" "example" { bucket = aws_s3_bucket.example.id rule { id = "transition-to-ia" status = "Enabled" transition { days = 30 storage_class = "STANDARD_IA" } transition { days = 90 storage_class = "GLACIER" } expiration { days = 365 } } } ``` ## Azure Cost Optimization ### Reserved VM Instances - 1 or 3 year terms - Up to 72% savings - Flexible sizing - Exchangeable ### Azure Hybrid Benefit - Use existing Windows Server licenses - Up to 80% savings with RI - Available for Windows and SQL Server ### Azure Advisor Recommendations - Right-size VMs - Delete unused resources - Use reserved capacity - Optimize storage ## GCP Cost Optimization ### Committed Use Discounts - 1 or 3 year commitment - Up to 57% savings - Applies to vCPUs and memory - Resource-based or spend-based ### Sustained Use Discounts - Automatic discounts - Up to 30% for running instances - No commitment required - Applies to Compute Engine, GKE ### Preemptible VMs - Up to 80% savings - 24-hour maximum runtime - Best for batch workloads ## Tagging Strategy ### AWS Tagging ```hcl locals { common_tags = { Environment = "production" Project = "my-project" CostCenter = "engineering" Owner = "team@example.com" ManagedBy = "terraform" } } resource "aws_instance" "example" { ami = "ami-12345678" instance_type = "t3.medium" tags = merge( local.common_tags, { Name = "web-server" } ) } ``` **Reference:** See `references/tagging-standards.md` ## Cost Monitoring ### Budget Alerts ```hcl # AWS Budget resource "aws_budgets_budget" "monthly" { name = "monthly-budget" budget_type = "COST" limit_amount = "1000" limit_unit = "USD" time_period_start = "2024-01-01_00:00" time_unit = "MONTHLY" notification { comparison_operator = "GREATER_THAN" threshold = 80 threshold_type = "PERCENTAGE" notification_type = "ACTUAL" subscriber_email_addresses = ["team@example.com"] } } ``` ### Cost Anomaly Detection - AWS Cost Anomaly Detection - Azure Cost Management alerts - GCP Budget alerts ## Architecture Patterns ### Pattern 1: Serverless First - Use Lambda/Functions for event-driven - Pay only for execution time - Auto-scaling included - No idle costs ### Pattern 2: Right-Sized Databases ``` Development: t3.small RDS Staging: t3.large RDS Production: r6g.2xlarge RDS with read replicas ``` ### Pattern 3: Multi-Tier Storage ``` Hot data: S3 Standard Warm data: S3 Standard-IA (30 days) Cold data: S3 Glacier (90 days) Archive: S3 Deep Archive (365 days) ``` ### Pattern 4: Auto-Scaling ```hcl resource "aws_autoscaling_policy" "scale_up" { name = "scale-up" scaling_adjustment = 2 adjustment_type = "ChangeInCapacity" cooldown = 300 autoscaling_group_name = aws_autoscaling_group.main.name } resource "aws_cloudwatch_metric_alarm" "cpu_high" { alarm_name = "cpu-high" comparison_operator = "GreaterThanThreshold" evaluation_periods = "2" metric_name = "CPUUtilization" namespace = "AWS/EC2" period = "60" statistic = "Average" threshold = "80" alarm_actions = [aws_autoscaling_policy.scale_up.arn] } ``` ## Cost Optimization Checklist - [ ] Implement cost allocation tags - [ ] Delete unused resources (EBS, EIPs, snapshots) - [ ] Right-size instances based on utilization - [ ] Use reserved capacity for steady workloads - [ ] Implement auto-scaling - [ ] Optimize storage classes - [ ] Use lifecycle policies - [ ] Enable cost anomaly detection - [ ] Set budget alerts - [ ] Review costs weekly - [ ] Use spot/preemptible instances - [ ] Optimize data transfer costs - [ ] Implement caching layers - [ ] Use managed services - [ ] Monitor and optimize continuously ## Tools - **AWS:** Cost Explorer, Cost Anomaly Detection, Compute Optimizer - **Azure:** Cost Management, Advisor - **GCP:** Cost Management, Recommender - **Multi-cloud:** CloudHealth, Cloudability, Kubecost ## Reference Files - `references/tagging-standards.md` - Tagging conventions - `assets/cost-analysis-template.xlsx` - Cost analysis spreadsheet ## Related Skills - `terraform-module-library` - For resource provisioning - `multi-cloud-architecture` - For cloud selection
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

hybrid-cloud-networking

Configure secure, high-performance connectivity between on-premises

architecture
⭐1
# Hybrid Cloud Networking Configure secure, high-performance connectivity between on-premises and cloud environments using VPN, Direct Connect, and ExpressRoute. ## Purpose Establish secure, reliable network connectivity between on-premises data centers and cloud providers (AWS, Azure, GCP). ## When to Use - Connect on-premises to cloud - Extend datacenter to cloud - Implement hybrid active-active setups - Meet compliance requirements - Migrate to cloud gradually ## Connection Options ### AWS Connectivity #### 1. Site-to-Site VPN - IPSec VPN over internet - Up to 1.25 Gbps per tunnel - Cost-effective for moderate bandwidth - Higher latency, internet-dependent ```hcl resource "aws_vpn_gateway" "main" { vpc_id = aws_vpc.main.id tags = { Name = "main-vpn-gateway" } } resource "aws_customer_gateway" "main" { bgp_asn = 65000 ip_address = "203.0.113.1" type = "ipsec.1" } resource "aws_vpn_connection" "main" { vpn_gateway_id = aws_vpn_gateway.main.id customer_gateway_id = aws_customer_gateway.main.id type = "ipsec.1" static_routes_only = false } ``` #### 2. AWS Direct Connect - Dedicated network connection - 1 Gbps to 100 Gbps - Lower latency, consistent bandwidth - More expensive, setup time required **Reference:** See `references/direct-connect.md` ### Azure Connectivity #### 1. Site-to-Site VPN ```hcl resource "azurerm_virtual_network_gateway" "vpn" { name = "vpn-gateway" location = azurerm_resource_group.main.location resource_group_name = azurerm_resource_group.main.name type = "Vpn" vpn_type = "RouteBased" sku = "VpnGw1" ip_configuration { name = "vnetGatewayConfig" public_ip_address_id = azurerm_public_ip.vpn.id private_ip_address_allocation = "Dynamic" subnet_id = azurerm_subnet.gateway.id } } ``` #### 2. Azure ExpressRoute - Private connection via connectivity provider - Up to 100 Gbps - Low latency, high reliability - Premium for global connectivity ### GCP Connectivity #### 1. Cloud VPN - IPSec VPN (Classic or HA VPN) - HA VPN: 99.99% SLA - Up to 3 Gbps per tunnel #### 2. Cloud Interconnect - Dedicated (10 Gbps, 100 Gbps) - Partner (50 Mbps to 50 Gbps) - Lower latency than VPN ## Hybrid Network Patterns ### Pattern 1: Hub-and-Spoke ``` On-Premises Datacenter ↓ VPN/Direct Connect ↓ Transit Gateway (AWS) / vWAN (Azure) ↓ β”œβ”€ Production VPC/VNet β”œβ”€ Staging VPC/VNet └─ Development VPC/VNet ``` ### Pattern 2: Multi-Region Hybrid ``` On-Premises β”œβ”€ Direct Connect β†’ us-east-1 └─ Direct Connect β†’ us-west-2 ↓ Cross-Region Peering ``` ### Pattern 3: Multi-Cloud Hybrid ``` On-Premises Datacenter β”œβ”€ Direct Connect β†’ AWS β”œβ”€ ExpressRoute β†’ Azure └─ Interconnect β†’ GCP ``` ## Routing Configuration ### BGP Configuration ``` On-Premises Router: - AS Number: 65000 - Advertise: 10.0.0.0/8 Cloud Router: - AS Number: 64512 (AWS), 65515 (Azure) - Advertise: Cloud VPC/VNet CIDRs ``` ### Route Propagation - Enable route propagation on route tables - Use BGP for dynamic routing - Implement route filtering - Monitor route advertisements ## Security Best Practices 1. **Use private connectivity** (Direct Connect/ExpressRoute) 2. **Implement encryption** for VPN tunnels 3. **Use VPC endpoints** to avoid internet routing 4. **Configure network ACLs** and security groups 5. **Enable VPC Flow Logs** for monitoring 6. **Implement DDoS protection** 7. **Use PrivateLink/Private Endpoints** 8. **Monitor connections** with CloudWatch/Monitor 9. **Implement redundancy** (dual tunnels) 10. **Regular security audits** ## High Availability ### Dual VPN Tunnels ```hcl resource "aws_vpn_connection" "primary" { vpn_gateway_id = aws_vpn_gateway.main.id customer_gateway_id = aws_customer_gateway.primary.id type = "ipsec.1" } resource "aws_vpn_connection" "secondary" { vpn_gateway_id = aws_vpn_gateway.main.id customer_gateway_id = aws_customer_gateway.secondary.id type = "ipsec.1" } ``` ### Active-Active Configuration - Multiple connections from different locations - BGP for automatic failover - Equal-cost multi-path (ECMP) routing - Monitor health of all connections ## Monitoring and Troubleshooting ### Key Metrics - Tunnel status (up/down) - Bytes in/out - Packet loss - Latency - BGP session status ### Troubleshooting ```bash # AWS VPN aws ec2 describe-vpn-connections aws ec2 get-vpn-connection-telemetry # Azure VPN az network vpn-connection show az network vpn-connection show-device-config-script ``` ## Cost Optimization 1. **Right-size connections** based on traffic 2. **Use VPN for low-bandwidth** workloads 3. **Consolidate traffic** through fewer connections 4. **Minimize data transfer** costs 5. **Use Direct Connect** for high bandwidth 6. **Implement caching** to reduce traffic ## Reference Files - `references/vpn-setup.md` - VPN configuration guide - `references/direct-connect.md` - Direct Connect setup ## Related Skills - `multi-cloud-architecture` - For architecture decisions - `terraform-module-library` - For IaC implementation
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

multi-cloud-architecture

Design multi-cloud architectures using a decision framework to

architecture
⭐1
# Multi-Cloud Architecture Decision framework and patterns for architecting applications across AWS, Azure, and GCP. ## Purpose Design cloud-agnostic architectures and make informed decisions about service selection across cloud providers. ## When to Use - Design multi-cloud strategies - Migrate between cloud providers - Select cloud services for specific workloads - Implement cloud-agnostic architectures - Optimize costs across providers ## Cloud Service Comparison ### Compute Services | AWS | Azure | GCP | Use Case | | ------- | ------------------- | --------------- | ------------------ | | EC2 | Virtual Machines | Compute Engine | IaaS VMs | | ECS | Container Instances | Cloud Run | Containers | | EKS | AKS | GKE | Kubernetes | | Lambda | Functions | Cloud Functions | Serverless | | Fargate | Container Apps | Cloud Run | Managed containers | ### Storage Services | AWS | Azure | GCP | Use Case | | ------- | --------------- | --------------- | -------------- | | S3 | Blob Storage | Cloud Storage | Object storage | | EBS | Managed Disks | Persistent Disk | Block storage | | EFS | Azure Files | Filestore | File storage | | Glacier | Archive Storage | Archive Storage | Cold storage | ### Database Services | AWS | Azure | GCP | Use Case | | ----------- | ---------------- | ------------- | --------------- | | RDS | SQL Database | Cloud SQL | Managed SQL | | DynamoDB | Cosmos DB | Firestore | NoSQL | | Aurora | PostgreSQL/MySQL | Cloud Spanner | Distributed SQL | | ElastiCache | Cache for Redis | Memorystore | Caching | **Reference:** See `references/service-comparison.md` for complete comparison ## Multi-Cloud Patterns ### Pattern 1: Single Provider with DR - Primary workload in one cloud - Disaster recovery in another - Database replication across clouds - Automated failover ### Pattern 2: Best-of-Breed - Use best service from each provider - AI/ML on GCP - Enterprise apps on Azure - General compute on AWS ### Pattern 3: Geographic Distribution - Serve users from nearest cloud region - Data sovereignty compliance - Global load balancing - Regional failover ### Pattern 4: Cloud-Agnostic Abstraction - Kubernetes for compute - PostgreSQL for database - S3-compatible storage (MinIO) - Open source tools ## Cloud-Agnostic Architecture ### Use Cloud-Native Alternatives - **Compute:** Kubernetes (EKS/AKS/GKE) - **Database:** PostgreSQL/MySQL (RDS/SQL Database/Cloud SQL) - **Message Queue:** Apache Kafka (MSK/Event Hubs/Confluent) - **Cache:** Redis (ElastiCache/Azure Cache/Memorystore) - **Object Storage:** S3-compatible API - **Monitoring:** Prometheus/Grafana - **Service Mesh:** Istio/Linkerd ### Abstraction Layers ``` Application Layer ↓ Infrastructure Abstraction (Terraform) ↓ Cloud Provider APIs ↓ AWS / Azure / GCP ``` ## Cost Comparison ### Compute Pricing Factors - **AWS:** On-demand, Reserved, Spot, Savings Plans - **Azure:** Pay-as-you-go, Reserved, Spot - **GCP:** On-demand, Committed use, Preemptible ### Cost Optimization Strategies 1. Use reserved/committed capacity (30-70% savings) 2. Leverage spot/preemptible instances 3. Right-size resources 4. Use serverless for variable workloads 5. Optimize data transfer costs 6. Implement lifecycle policies 7. Use cost allocation tags 8. Monitor with cloud cost tools **Reference:** See `references/multi-cloud-patterns.md` ## Migration Strategy ### Phase 1: Assessment - Inventory current infrastructure - Identify dependencies - Assess cloud compatibility - Estimate costs ### Phase 2: Pilot - Select pilot workload - Implement in target cloud - Test thoroughly - Document learnings ### Phase 3: Migration - Migrate workloads incrementally - Maintain dual-run period - Monitor performance - Validate functionality ### Phase 4: Optimization - Right-size resources - Implement cloud-native services - Optimize costs - Enhance security ## Best Practices 1. **Use infrastructure as code** (Terraform/OpenTofu) 2. **Implement CI/CD pipelines** for deployments 3. **Design for failure** across clouds 4. **Use managed services** when possible 5. **Implement comprehensive monitoring** 6. **Automate cost optimization** 7. **Follow security best practices** 8. **Document cloud-specific configurations** 9. **Test disaster recovery** procedures 10. **Train teams** on multiple clouds ## Reference Files - `references/service-comparison.md` - Complete service comparison - `references/multi-cloud-patterns.md` - Architecture patterns ## Related Skills - `terraform-module-library` - For IaC implementation - `cost-optimization` - For cost management - `hybrid-cloud-networking` - For connectivity
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

service-mesh-observability

Implement comprehensive observability for service meshes including

architecture
⭐1
# Service Mesh Observability Complete guide to observability patterns for Istio, Linkerd, and service mesh deployments. ## When to Use This Skill - Setting up distributed tracing across services - Implementing service mesh metrics and dashboards - Debugging latency and error issues - Defining SLOs for service communication - Visualizing service dependencies - Troubleshooting mesh connectivity ## Core Concepts ### 1. Three Pillars of Observability ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Observability β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ Metrics β”‚ Traces β”‚ Logs β”‚ β”‚ β”‚ β”‚ β”‚ β”‚ β€’ Request rate β”‚ β€’ Span context β”‚ β€’ Access logs β”‚ β”‚ β€’ Error rate β”‚ β€’ Latency β”‚ β€’ Error details β”‚ β”‚ β€’ Latency P50 β”‚ β€’ Dependencies β”‚ β€’ Debug info β”‚ β”‚ β€’ Saturation β”‚ β€’ Bottlenecks β”‚ β€’ Audit trail β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### 2. Golden Signals for Mesh | Signal | Description | Alert Threshold | | -------------- | ------------------------- | ----------------- | | **Latency** | Request duration P50, P99 | P99 > 500ms | | **Traffic** | Requests per second | Anomaly detection | | **Errors** | 5xx error rate | > 1% | | **Saturation** | Resource utilization | > 80% | ## Templates ### Template 1: Istio with Prometheus & Grafana ```yaml # Install Prometheus apiVersion: v1 kind: ConfigMap metadata: name: prometheus namespace: istio-system data: prometheus.yml: | global: scrape_interval: 15s scrape_configs: - job_name: 'istio-mesh' kubernetes_sd_configs: - role: endpoints namespaces: names: - istio-system relabel_configs: - source_labels: [__meta_kubernetes_service_name] action: keep regex: istio-telemetry --- # ServiceMonitor for Prometheus Operator apiVersion: monitoring.coreos.com/v1 kind: ServiceMonitor metadata: name: istio-mesh namespace: istio-system spec: selector: matchLabels: app: istiod endpoints: - port: http-monitoring interval: 15s ``` ### Template 2: Key Istio Metrics Queries ```promql # Request rate by service sum(rate(istio_requests_total{reporter="destination"}[5m])) by (destination_service_name) # Error rate (5xx) sum(rate(istio_requests_total{reporter="destination", response_code=~"5.."}[5m])) / sum(rate(istio_requests_total{reporter="destination"}[5m])) * 100 # P99 latency histogram_quantile(0.99, sum(rate(istio_request_duration_milliseconds_bucket{reporter="destination"}[5m])) by (le, destination_service_name)) # TCP connections sum(istio_tcp_connections_opened_total{reporter="destination"}) by (destination_service_name) # Request size histogram_quantile(0.99, sum(rate(istio_request_bytes_bucket{reporter="destination"}[5m])) by (le, destination_service_name)) ``` ### Template 3: Jaeger Distributed Tracing ```yaml # Jaeger installation for Istio apiVersion: install.istio.io/v1alpha1 kind: IstioOperator spec: meshConfig: enableTracing: true defaultConfig: tracing: sampling: 100.0 # 100% in dev, lower in prod zipkin: address: jaeger-collector.istio-system:9411 --- # Jaeger deployment apiVersion: apps/v1 kind: Deployment metadata: name: jaeger namespace: istio-system spec: selector: matchLabels: app: jaeger template: metadata: labels: app: jaeger spec: containers: - name: jaeger image: jaegertracing/all-in-one:1.50 ports: - containerPort: 5775 # UDP - containerPort: 6831 # Thrift - containerPort: 6832 # Thrift - containerPort: 5778 # Config - containerPort: 16686 # UI - containerPort: 14268 # HTTP - containerPort: 14250 # gRPC - containerPort: 9411 # Zipkin env: - name: COLLECTOR_ZIPKIN_HOST_PORT value: ":9411" ``` ### Template 4: Linkerd Viz Dashboard ```bash # Install Linkerd viz extension linkerd viz install | kubectl apply -f - # Access dashboard linkerd viz dashboard # CLI commands for observability # Top requests linkerd viz top deploy/my-app # Per-route metrics linkerd viz routes deploy/my-app --to deploy/backend # Live traffic inspection linkerd viz tap deploy/my-app --to deploy/backend # Service edges (dependencies) linkerd viz edges deployment -n my-namespace ``` ### Template 5: Grafana Dashboard JSON ```json { "dashboard": { "title": "Service Mesh Overview", "panels": [ { "title": "Request Rate", "type": "graph", "targets": [ { "expr": "sum(rate(istio_requests_total{reporter=\"destination\"}[5m])) by (destination_service_name)", "legendFormat": "{{destination_service_name}}" } ] }, { "title": "Error Rate", "type": "gauge", "targets": [ { "expr": "sum(rate(istio_requests_total{response_code=~\"5..\"}[5m])) / sum(rate(istio_requests_total[5m])) * 100" } ], "fieldConfig": { "defaults": { "thresholds": { "steps": [ { "value": 0, "color": "green" }, { "value": 1, "color": "yellow" }, { "value": 5, "color": "red" } ] } } } }, { "title": "P99 Latency", "type": "graph", "targets": [ { "expr": "histogram_quantile(0.99, sum(rate(istio_request_duration_milliseconds_bucket{reporter=\"destination\"}[5m])) by (le, destination_service_name))", "legendFormat": "{{destination_service_name}}" } ] }, { "title": "Service Topology", "type": "nodeGraph", "targets": [ { "expr": "sum(rate(istio_requests_total{reporter=\"destination\"}[5m])) by (source_workload, destination_service_name)" } ] } ] } } ``` ### Template 6: Kiali Service Mesh Visualization ```yaml # Kiali installation apiVersion: kiali.io/v1alpha1 kind: Kiali metadata: name: kiali namespace: istio-system spec: auth: strategy: anonymous # or openid, token deployment: accessible_namespaces: - "**" external_services: prometheus: url: http://prometheus.istio-system:9090 tracing: url: http://jaeger-query.istio-system:16686 grafana: url: http://grafana.istio-system:3000 ``` ### Template 7: OpenTelemetry Integration ```yaml # OpenTelemetry Collector for mesh apiVersion: v1 kind: ConfigMap metadata: name: otel-collector-config data: config.yaml: | receivers: otlp: protocols: grpc: endpoint: 0.0.0.0:4317 http: endpoint: 0.0.0.0:4318 zipkin: endpoint: 0.0.0.0:9411 processors: batch: timeout: 10s exporters: jaeger: endpoint: jaeger-collector:14250 tls: insecure: true prometheus: endpoint: 0.0.0.0:8889 service: pipelines: traces: receivers: [otlp, zipkin] processors: [batch] exporters: [jaeger] metrics: receivers: [otlp] processors: [batch] exporters: [prometheus] --- # Istio Telemetry v2 with OTel apiVersion: telemetry.istio.io/v1alpha1 kind: Telemetry metadata: name: mesh-default namespace: istio-system spec: tracing: - providers: - name: otel randomSamplingPercentage: 10 ``` ## Alerting Rules ```yaml apiVersion: monitoring.coreos.com/v1 kind: PrometheusRule metadata: name: mesh-alerts namespace: istio-system spec: groups: - name: mesh.rules rules: - alert: HighErrorRate expr: | sum(rate(istio_requests_total{response_code=~"5.."}[5m])) by (destination_service_name) / sum(rate(istio_requests_total[5m])) by (destination_service_name) > 0.05 for: 5m labels: severity: critical annotations: summary: "High error rate for {{ $labels.destination_service_name }}" - alert: HighLatency expr: | histogram_quantile(0.99, sum(rate(istio_request_duration_milliseconds_bucket[5m])) by (le, destination_service_name)) > 1000 for: 5m labels: severity: warning annotations: summary: "High P99 latency for {{ $labels.destination_service_name }}" - alert: MeshCertExpiring expr: | (certmanager_certificate_expiration_timestamp_seconds - time()) / 86400 < 7 labels: severity: warning annotations: summary: "Mesh certificate expiring in less than 7 days" ``` ## Best Practices ### Do's - **Sample appropriately** - 100% in dev, 1-10% in prod - **Use trace context** - Propagate headers consistently - **Set up alerts** - For golden signals - **Correlate metrics/traces** - Use exemplars - **Retain strategically** - Hot/cold storage tiers ### Don'ts - **Don't over-sample** - Storage costs add up - **Don't ignore cardinality** - Limit label values - **Don't skip dashboards** - Visualize dependencies - **Don't forget costs** - Monitor observability costs ## Resources - [Istio Observability](https://istio.io/latest/docs/tasks/observability/) - [Linkerd Observability](https://linkerd.io/2.14/features/dashboard/) - [OpenTelemetry](https://opentelemetry.io/) - [Kiali](https://kiali.io/)
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

terraform-module-library

Build reusable Terraform modules for AWS, Azure, and GCP

architecture
⭐1
# Terraform Module Library Production-ready Terraform module patterns for AWS, Azure, and GCP infrastructure. ## Purpose Create reusable, well-tested Terraform modules for common cloud infrastructure patterns across multiple cloud providers. ## When to Use - Build reusable infrastructure components - Standardize cloud resource provisioning - Implement infrastructure as code best practices - Create multi-cloud compatible modules - Establish organizational Terraform standards ## Module Structure ``` terraform-modules/ β”œβ”€β”€ aws/ β”‚ β”œβ”€β”€ vpc/ β”‚ β”œβ”€β”€ eks/ β”‚ β”œβ”€β”€ rds/ β”‚ └── s3/ β”œβ”€β”€ azure/ β”‚ β”œβ”€β”€ vnet/ β”‚ β”œβ”€β”€ aks/ β”‚ └── storage/ └── gcp/ β”œβ”€β”€ vpc/ β”œβ”€β”€ gke/ └── cloud-sql/ ``` ## Standard Module Pattern ``` module-name/ β”œβ”€β”€ main.tf # Main resources β”œβ”€β”€ variables.tf # Input variables β”œβ”€β”€ outputs.tf # Output values β”œβ”€β”€ versions.tf # Provider versions β”œβ”€β”€ README.md # Documentation β”œβ”€β”€ examples/ # Usage examples β”‚ └── complete/ β”‚ β”œβ”€β”€ main.tf β”‚ └── variables.tf └── tests/ # Terratest files └── module_test.go ``` ## AWS VPC Module Example **main.tf:** ```hcl resource "aws_vpc" "main" { cidr_block = var.cidr_block enable_dns_hostnames = var.enable_dns_hostnames enable_dns_support = var.enable_dns_support tags = merge( { Name = var.name }, var.tags ) } resource "aws_subnet" "private" { count = length(var.private_subnet_cidrs) vpc_id = aws_vpc.main.id cidr_block = var.private_subnet_cidrs[count.index] availability_zone = var.availability_zones[count.index] tags = merge( { Name = "${var.name}-private-${count.index + 1}" Tier = "private" }, var.tags ) } resource "aws_internet_gateway" "main" { count = var.create_internet_gateway ? 1 : 0 vpc_id = aws_vpc.main.id tags = merge( { Name = "${var.name}-igw" }, var.tags ) } ``` **variables.tf:** ```hcl variable "name" { description = "Name of the VPC" type = string } variable "cidr_block" { description = "CIDR block for VPC" type = string validation { condition = can(regex("^([0-9]{1,3}\\.){3}[0-9]{1,3}/[0-9]{1,2}$", var.cidr_block)) error_message = "CIDR block must be valid IPv4 CIDR notation." } } variable "availability_zones" { description = "List of availability zones" type = list(string) } variable "private_subnet_cidrs" { description = "CIDR blocks for private subnets" type = list(string) default = [] } variable "enable_dns_hostnames" { description = "Enable DNS hostnames in VPC" type = bool default = true } variable "tags" { description = "Additional tags" type = map(string) default = {} } ``` **outputs.tf:** ```hcl output "vpc_id" { description = "ID of the VPC" value = aws_vpc.main.id } output "private_subnet_ids" { description = "IDs of private subnets" value = aws_subnet.private[*].id } output "vpc_cidr_block" { description = "CIDR block of VPC" value = aws_vpc.main.cidr_block } ``` ## Best Practices 1. **Use semantic versioning** for modules 2. **Document all variables** with descriptions 3. **Provide examples** in examples/ directory 4. **Use validation blocks** for input validation 5. **Output important attributes** for module composition 6. **Pin provider versions** in versions.tf 7. **Use locals** for computed values 8. **Implement conditional resources** with count/for_each 9. **Test modules** with Terratest 10. **Tag all resources** consistently ## Module Composition ```hcl module "vpc" { source = "../../modules/aws/vpc" name = "production" cidr_block = "10.0.0.0/16" availability_zones = ["us-west-2a", "us-west-2b", "us-west-2c"] private_subnet_cidrs = [ "10.0.1.0/24", "10.0.2.0/24", "10.0.3.0/24" ] tags = { Environment = "production" ManagedBy = "terraform" } } module "rds" { source = "../../modules/aws/rds" identifier = "production-db" engine = "postgres" engine_version = "15.3" instance_class = "db.t3.large" vpc_id = module.vpc.vpc_id subnet_ids = module.vpc.private_subnet_ids tags = { Environment = "production" } } ``` ## Reference Files - `assets/vpc-module/` - Complete VPC module example - `assets/rds-module/` - RDS module example - `references/aws-modules.md` - AWS module patterns - `references/azure-modules.md` - Azure module patterns - `references/gcp-modules.md` - GCP module patterns ## Testing ```go // tests/vpc_test.go package test import ( "testing" "github.com/gruntwork-io/terratest/modules/terraform" "github.com/stretchr/testify/assert" ) func TestVPCModule(t *testing.T) { terraformOptions := &terraform.Options{ TerraformDir: "../examples/complete", } defer terraform.Destroy(t, terraformOptions) terraform.InitAndApply(t, terraformOptions) vpcID := terraform.Output(t, terraformOptions, "vpc_id") assert.NotEmpty(t, vpcID) } ``` ## Related Skills - `multi-cloud-architecture` - For architectural decisions - `cost-optimization` - For cost-effective designs
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–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

sql-optimization-patterns

Master SQL query optimization, indexing strategies, and EXPLAIN

coding
⭐1
# SQL Optimization Patterns Transform slow database queries into lightning-fast operations through systematic optimization, proper indexing, and query plan analysis. ## When to Use This Skill - Debugging slow-running queries - Designing performant database schemas - Optimizing application response times - Reducing database load and costs - Improving scalability for growing datasets - Analyzing EXPLAIN query plans - Implementing efficient indexes - Resolving N+1 query problems ## Core Concepts ### 1. Query Execution Plans (EXPLAIN) Understanding EXPLAIN output is fundamental to optimization. **PostgreSQL EXPLAIN:** ```sql -- Basic explain EXPLAIN SELECT * FROM users WHERE email = 'user@example.com'; -- With actual execution stats EXPLAIN ANALYZE SELECT * FROM users WHERE email = 'user@example.com'; -- Verbose output with more details EXPLAIN (ANALYZE, BUFFERS, VERBOSE) SELECT u.*, o.order_total FROM users u JOIN orders o ON u.id = o.user_id WHERE u.created_at > NOW() - INTERVAL '30 days'; ``` **Key Metrics to Watch:** - **Seq Scan**: Full table scan (usually slow for large tables) - **Index Scan**: Using index (good) - **Index Only Scan**: Using index without touching table (best) - **Nested Loop**: Join method (okay for small datasets) - **Hash Join**: Join method (good for larger datasets) - **Merge Join**: Join method (good for sorted data) - **Cost**: Estimated query cost (lower is better) - **Rows**: Estimated rows returned - **Actual Time**: Real execution time ### 2. Index Strategies Indexes are the most powerful optimization tool. **Index Types:** - **B-Tree**: Default, good for equality and range queries - **Hash**: Only for equality (=) comparisons - **GIN**: Full-text search, array queries, JSONB - **GiST**: Geometric data, full-text search - **BRIN**: Block Range INdex for very large tables with correlation ```sql -- Standard B-Tree index CREATE INDEX idx_users_email ON users(email); -- Composite index (order matters!) CREATE INDEX idx_orders_user_status ON orders(user_id, status); -- Partial index (index subset of rows) CREATE INDEX idx_active_users ON users(email) WHERE status = 'active'; -- Expression index CREATE INDEX idx_users_lower_email ON users(LOWER(email)); -- Covering index (include additional columns) CREATE INDEX idx_users_email_covering ON users(email) INCLUDE (name, created_at); -- Full-text search index CREATE INDEX idx_posts_search ON posts USING GIN(to_tsvector('english', title || ' ' || body)); -- JSONB index CREATE INDEX idx_metadata ON events USING GIN(metadata); ``` ### 3. Query Optimization Patterns **Avoid SELECT \*:** ```sql -- Bad: Fetches unnecessary columns SELECT * FROM users WHERE id = 123; -- Good: Fetch only what you need SELECT id, email, name FROM users WHERE id = 123; ``` **Use WHERE Clause Efficiently:** ```sql -- Bad: Function prevents index usage SELECT * FROM users WHERE LOWER(email) = 'user@example.com'; -- Good: Create functional index or use exact match CREATE INDEX idx_users_email_lower ON users(LOWER(email)); -- Then: SELECT * FROM users WHERE LOWER(email) = 'user@example.com'; -- Or store normalized data SELECT * FROM users WHERE email = 'user@example.com'; ``` **Optimize JOINs:** ```sql -- Bad: Cartesian product then filter SELECT u.name, o.total FROM users u, orders o WHERE u.id = o.user_id AND u.created_at > '2024-01-01'; -- Good: Filter before join SELECT u.name, o.total FROM users u JOIN orders o ON u.id = o.user_id WHERE u.created_at > '2024-01-01'; -- Better: Filter both tables SELECT u.name, o.total FROM (SELECT * FROM users WHERE created_at > '2024-01-01') u JOIN orders o ON u.id = o.user_id; ``` ## Optimization Patterns ### Pattern 1: Eliminate N+1 Queries **Problem: N+1 Query Anti-Pattern** ```python # Bad: Executes N+1 queries users = db.query("SELECT * FROM users LIMIT 10") for user in users: orders = db.query("SELECT * FROM orders WHERE user_id = ?", user.id) # Process orders ``` **Solution: Use JOINs or Batch Loading** ```sql -- Solution 1: JOIN SELECT u.id, u.name, o.id as order_id, o.total FROM users u LEFT JOIN orders o ON u.id = o.user_id WHERE u.id IN (1, 2, 3, 4, 5); -- Solution 2: Batch query SELECT * FROM orders WHERE user_id IN (1, 2, 3, 4, 5); ``` ```python # Good: Single query with JOIN or batch load # Using JOIN results = db.query(""" SELECT u.id, u.name, o.id as order_id, o.total FROM users u LEFT JOIN orders o ON u.id = o.user_id WHERE u.id IN (1, 2, 3, 4, 5) """) # Or batch load users = db.query("SELECT * FROM users LIMIT 10") user_ids = [u.id for u in users] orders = db.query( "SELECT * FROM orders WHERE user_id IN (?)", user_ids ) # Group orders by user_id orders_by_user = {} for order in orders: orders_by_user.setdefault(order.user_id, []).append(order) ``` ### Pattern 2: Optimize Pagination **Bad: OFFSET on Large Tables** ```sql -- Slow for large offsets SELECT * FROM users ORDER BY created_at DESC LIMIT 20 OFFSET 100000; -- Very slow! ``` **Good: Cursor-Based Pagination** ```sql -- Much faster: Use cursor (last seen ID) SELECT * FROM users WHERE created_at < '2024-01-15 10:30:00' -- Last cursor ORDER BY created_at DESC LIMIT 20; -- With composite sorting SELECT * FROM users WHERE (created_at, id) < ('2024-01-15 10:30:00', 12345) ORDER BY created_at DESC, id DESC LIMIT 20; -- Requires index CREATE INDEX idx_users_cursor ON users(created_at DESC, id DESC); ``` ### Pattern 3: Aggregate Efficiently **Optimize COUNT Queries:** ```sql -- Bad: Counts all rows SELECT COUNT(*) FROM orders; -- Slow on large tables -- Good: Use estimates for approximate counts SELECT reltuples::bigint AS estimate FROM pg_class WHERE relname = 'orders'; -- Good: Filter before counting SELECT COUNT(*) FROM orders WHERE created_at > NOW() - INTERVAL '7 days'; -- Better: Use index-only scan CREATE INDEX idx_orders_created ON orders(created_at); SELECT COUNT(*) FROM orders WHERE created_at > NOW() - INTERVAL '7 days'; ``` **Optimize GROUP BY:** ```sql -- Bad: Group by then filter SELECT user_id, COUNT(*) as order_count FROM orders GROUP BY user_id HAVING COUNT(*) > 10; -- Better: Filter first, then group (if possible) SELECT user_id, COUNT(*) as order_count FROM orders WHERE status = 'completed' GROUP BY user_id HAVING COUNT(*) > 10; -- Best: Use covering index CREATE INDEX idx_orders_user_status ON orders(user_id, status); ``` ### Pattern 4: Subquery Optimization **Transform Correlated Subqueries:** ```sql -- Bad: Correlated subquery (runs for each row) SELECT u.name, u.email, (SELECT COUNT(*) FROM orders o WHERE o.user_id = u.id) as order_count FROM users u; -- Good: JOIN with aggregation SELECT u.name, u.email, COUNT(o.id) as order_count FROM users u LEFT JOIN orders o ON o.user_id = u.id GROUP BY u.id, u.name, u.email; -- Better: Use window functions SELECT DISTINCT ON (u.id) u.name, u.email, COUNT(o.id) OVER (PARTITION BY u.id) as order_count FROM users u LEFT JOIN orders o ON o.user_id = u.id; ``` **Use CTEs for Clarity:** ```sql -- Using Common Table Expressions WITH recent_users AS ( SELECT id, name, email FROM users WHERE created_at > NOW() - INTERVAL '30 days' ), user_order_counts AS ( SELECT user_id, COUNT(*) as order_count FROM orders WHERE created_at > NOW() - INTERVAL '30 days' GROUP BY user_id ) SELECT ru.name, ru.email, COALESCE(uoc.order_count, 0) as orders FROM recent_users ru LEFT JOIN user_order_counts uoc ON ru.id = uoc.user_id; ``` ### Pattern 5: Batch Operations **Batch INSERT:** ```sql -- Bad: Multiple individual inserts INSERT INTO users (name, email) VALUES ('Alice', 'alice@example.com'); INSERT INTO users (name, email) VALUES ('Bob', 'bob@example.com'); INSERT INTO users (name, email) VALUES ('Carol', 'carol@example.com'); -- Good: Batch insert INSERT INTO users (name, email) VALUES ('Alice', 'alice@example.com'), ('Bob', 'bob@example.com'), ('Carol', 'carol@example.com'); -- Better: Use COPY for bulk inserts (PostgreSQL) COPY users (name, email) FROM '/tmp/users.csv' CSV HEADER; ``` **Batch UPDATE:** ```sql -- Bad: Update in loop UPDATE users SET status = 'active' WHERE id = 1; UPDATE users SET status = 'active' WHERE id = 2; -- ... repeat for many IDs -- Good: Single UPDATE with IN clause UPDATE users SET status = 'active' WHERE id IN (1, 2, 3, 4, 5, ...); -- Better: Use temporary table for large batches CREATE TEMP TABLE temp_user_updates (id INT, new_status VARCHAR); INSERT INTO temp_user_updates VALUES (1, 'active'), (2, 'active'), ...; UPDATE users u SET status = t.new_status FROM temp_user_updates t WHERE u.id = t.id; ``` ## Advanced Techniques ### Materialized Views Pre-compute expensive queries. ```sql -- Create materialized view CREATE MATERIALIZED VIEW user_order_summary AS SELECT u.id, u.name, COUNT(o.id) as total_orders, SUM(o.total) as total_spent, MAX(o.created_at) as last_order_date FROM users u LEFT JOIN orders o ON u.id = o.user_id GROUP BY u.id, u.name; -- Add index to materialized view CREATE INDEX idx_user_summary_spent ON user_order_summary(total_spent DESC); -- Refresh materialized view REFRESH MATERIALIZED VIEW user_order_summary; -- Concurrent refresh (PostgreSQL) REFRESH MATERIALIZED VIEW CONCURRENTLY user_order_summary; -- Query materialized view (very fast) SELECT * FROM user_order_summary WHERE total_spent > 1000 ORDER BY total_spent DESC; ``` ### Partitioning Split large tables for better performance. ```sql -- Range partitioning by date (PostgreSQL) CREATE TABLE orders ( id SERIAL, user_id INT, total DECIMAL, created_at TIMESTAMP ) PARTITION BY RANGE (created_at); -- Create partitions CREATE TABLE orders_2024_q1 PARTITION OF orders FOR VALUES FROM ('2024-01-01') TO ('2024-04-01'); CREATE TABLE orders_2024_q2 PARTITION OF orders FOR VALUES FROM ('2024-04-01') TO ('2024-07-01'); -- Queries automatically use appropriate partition SELECT * FROM orders WHERE created_at BETWEEN '2024-02-01' AND '2024-02-28'; -- Only scans orders_2024_q1 partition ``` ### Query Hints and Optimization ```sql -- Force index usage (MySQL) SELECT * FROM users USE INDEX (idx_users_email) WHERE email = 'user@example.com'; -- Parallel query (PostgreSQL) SET max_parallel_workers_per_gather = 4; SELECT * FROM large_table WHERE condition; -- Join hints (PostgreSQL) SET enable_nestloop = OFF; -- Force hash or merge join ``` ## Best Practices 1. **Index Selectively**: Too many indexes slow down writes 2. **Monitor Query Performance**: Use slow query logs 3. **Keep Statistics Updated**: Run ANALYZE regularly 4. **Use Appropriate Data Types**: Smaller types = better performance 5. **Normalize Thoughtfully**: Balance normalization vs performance 6. **Cache Frequently Accessed Data**: Use application-level caching 7. **Connection Pooling**: Reuse database connections 8. **Regular Maintenance**: VACUUM, ANALYZE, rebuild indexes ```sql -- Update statistics ANALYZE users; ANALYZE VERBOSE orders; -- Vacuum (PostgreSQL) VACUUM ANALYZE users; VACUUM FULL users; -- Reclaim space (locks table) -- Reindex REINDEX INDEX idx_users_email; REINDEX TABLE users; ``` ## Common Pitfalls - **Over-Indexing**: Each index slows down INSERT/UPDATE/DELETE - **Unused Indexes**: Waste space and slow writes - **Missing Indexes**: Slow queries, full table scans - **Implicit Type Conversion**: Prevents index usage - **OR Conditions**: Can't use indexes efficiently - **LIKE with Leading Wildcard**: `LIKE '%abc'` can't use index - **Function in WHERE**: Prevents index usage unless functional index exists ## Monitoring Queries ```sql -- Find slow queries (PostgreSQL) SELECT query, calls, total_time, mean_time FROM pg_stat_statements ORDER BY mean_time DESC LIMIT 10; -- Find missing indexes (PostgreSQL) SELECT schemaname, tablename, seq_scan, seq_tup_read, idx_scan, seq_tup_read / seq_scan AS avg_seq_tup_read FROM pg_stat_user_tables WHERE seq_scan > 0 ORDER BY seq_tup_read DESC LIMIT 10; -- Find unused indexes (PostgreSQL) SELECT schemaname, tablename, indexname, idx_scan, idx_tup_read, idx_tup_fetch FROM pg_stat_user_indexes WHERE idx_scan = 0 ORDER BY pg_relation_size(indexrelid) DESC; ``` ## Resources - **references/postgres-optimization-guide.md**: PostgreSQL-specific optimization - **references/mysql-optimization-guide.md**: MySQL/MariaDB optimization - **references/query-plan-analysis.md**: Deep dive into EXPLAIN plans - **assets/index-strategy-checklist.md**: When and how to create indexes - **assets/query-optimization-checklist.md**: Step-by-step optimization guide - **scripts/analyze-slow-queries.sql**: Identify slow queries in your database - **scripts/index-recommendations.sql**: Generate index recommendations
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

architecture-decision-records

Write and maintain Architecture Decision Records (ADRs) following

coding
⭐1
# Architecture Decision Records Comprehensive patterns for creating, maintaining, and managing Architecture Decision Records (ADRs) that capture the context and rationale behind significant technical decisions. ## When to Use This Skill - Making significant architectural decisions - Documenting technology choices - Recording design trade-offs - Onboarding new team members - Reviewing historical decisions - Establishing decision-making processes ## Core Concepts ### 1. What is an ADR? An Architecture Decision Record captures: - **Context**: Why we needed to make a decision - **Decision**: What we decided - **Consequences**: What happens as a result ### 2. When to Write an ADR | Write ADR | Skip ADR | | -------------------------- | ---------------------- | | New framework adoption | Minor version upgrades | | Database technology choice | Bug fixes | | API design patterns | Implementation details | | Security architecture | Routine maintenance | | Integration patterns | Configuration changes | ### 3. ADR Lifecycle ``` Proposed β†’ Accepted β†’ Deprecated β†’ Superseded ↓ Rejected ``` ## Templates ### Template 1: Standard ADR (MADR Format) ```markdown # ADR-0001: Use PostgreSQL as Primary Database ## Status Accepted ## Context We need to select a primary database for our new e-commerce platform. The system will handle: - ~10,000 concurrent users - Complex product catalog with hierarchical categories - Transaction processing for orders and payments - Full-text search for products - Geospatial queries for store locator The team has experience with MySQL, PostgreSQL, and MongoDB. We need ACID compliance for financial transactions. ## Decision Drivers - **Must have ACID compliance** for payment processing - **Must support complex queries** for reporting - **Should support full-text search** to reduce infrastructure complexity - **Should have good JSON support** for flexible product attributes - **Team familiarity** reduces onboarding time ## Considered Options ### Option 1: PostgreSQL - **Pros**: ACID compliant, excellent JSON support (JSONB), built-in full-text search, PostGIS for geospatial, team has experience - **Cons**: Slightly more complex replication setup than MySQL ### Option 2: MySQL - **Pros**: Very familiar to team, simple replication, large community - **Cons**: Weaker JSON support, no built-in full-text search (need Elasticsearch), no geospatial without extensions ### Option 3: MongoDB - **Pros**: Flexible schema, native JSON, horizontal scaling - **Cons**: No ACID for multi-document transactions (at decision time), team has limited experience, requires schema design discipline ## Decision We will use **PostgreSQL 15** as our primary database. ## Rationale PostgreSQL provides the best balance of: 1. **ACID compliance** essential for e-commerce transactions 2. **Built-in capabilities** (full-text search, JSONB, PostGIS) reduce infrastructure complexity 3. **Team familiarity** with SQL databases reduces learning curve 4. **Mature ecosystem** with excellent tooling and community support The slight complexity in replication is outweighed by the reduction in additional services (no separate Elasticsearch needed). ## Consequences ### Positive - Single database handles transactions, search, and geospatial queries - Reduced operational complexity (fewer services to manage) - Strong consistency guarantees for financial data - Team can leverage existing SQL expertise ### Negative - Need to learn PostgreSQL-specific features (JSONB, full-text search syntax) - Vertical scaling limits may require read replicas sooner - Some team members need PostgreSQL-specific training ### Risks - Full-text search may not scale as well as dedicated search engines - Mitigation: Design for potential Elasticsearch addition if needed ## Implementation Notes - Use JSONB for flexible product attributes - Implement connection pooling with PgBouncer - Set up streaming replication for read replicas - Use pg_trgm extension for fuzzy search ## Related Decisions - ADR-0002: Caching Strategy (Redis) - complements database choice - ADR-0005: Search Architecture - may supersede if Elasticsearch needed ## References - [PostgreSQL JSON Documentation](https://www.postgresql.org/docs/current/datatype-json.html) - [PostgreSQL Full Text Search](https://www.postgresql.org/docs/current/textsearch.html) - Internal: Performance benchmarks in `/docs/benchmarks/database-comparison.md` ``` ### Template 2: Lightweight ADR ```markdown # ADR-0012: Adopt TypeScript for Frontend Development **Status**: Accepted **Date**: 2024-01-15 **Deciders**: @alice, @bob, @charlie ## Context Our React codebase has grown to 50+ components with increasing bug reports related to prop type mismatches and undefined errors. PropTypes provide runtime-only checking. ## Decision Adopt TypeScript for all new frontend code. Migrate existing code incrementally. ## Consequences **Good**: Catch type errors at compile time, better IDE support, self-documenting code. **Bad**: Learning curve for team, initial slowdown, build complexity increase. **Mitigations**: TypeScript training sessions, allow gradual adoption with `allowJs: true`. ``` ### Template 3: Y-Statement Format ```markdown # ADR-0015: API Gateway Selection In the context of **building a microservices architecture**, facing **the need for centralized API management, authentication, and rate limiting**, we decided for **Kong Gateway** and against **AWS API Gateway and custom Nginx solution**, to achieve **vendor independence, plugin extensibility, and team familiarity with Lua**, accepting that **we need to manage Kong infrastructure ourselves**. ``` ### Template 4: ADR for Deprecation ```markdown # ADR-0020: Deprecate MongoDB in Favor of PostgreSQL ## Status Accepted (Supersedes ADR-0003) ## Context ADR-0003 (2021) chose MongoDB for user profile storage due to schema flexibility needs. Since then: - MongoDB's multi-document transactions remain problematic for our use case - Our schema has stabilized and rarely changes - We now have PostgreSQL expertise from other services - Maintaining two databases increases operational burden ## Decision Deprecate MongoDB and migrate user profiles to PostgreSQL. ## Migration Plan 1. **Phase 1** (Week 1-2): Create PostgreSQL schema, dual-write enabled 2. **Phase 2** (Week 3-4): Backfill historical data, validate consistency 3. **Phase 3** (Week 5): Switch reads to PostgreSQL, monitor 4. **Phase 4** (Week 6): Remove MongoDB writes, decommission ## Consequences ### Positive - Single database technology reduces operational complexity - ACID transactions for user data - Team can focus PostgreSQL expertise ### Negative - Migration effort (~4 weeks) - Risk of data issues during migration - Lose some schema flexibility ## Lessons Learned Document from ADR-0003 experience: - Schema flexibility benefits were overestimated - Operational cost of multiple databases was underestimated - Consider long-term maintenance in technology decisions ``` ### Template 5: Request for Comments (RFC) Style ```markdown # RFC-0025: Adopt Event Sourcing for Order Management ## Summary Propose adopting event sourcing pattern for the order management domain to improve auditability, enable temporal queries, and support business analytics. ## Motivation Current challenges: 1. Audit requirements need complete order history 2. "What was the order state at time X?" queries are impossible 3. Analytics team needs event stream for real-time dashboards 4. Order state reconstruction for customer support is manual ## Detailed Design ### Event Store ``` OrderCreated { orderId, customerId, items[], timestamp } OrderItemAdded { orderId, item, timestamp } OrderItemRemoved { orderId, itemId, timestamp } PaymentReceived { orderId, amount, paymentId, timestamp } OrderShipped { orderId, trackingNumber, timestamp } ``` ### Projections - **CurrentOrderState**: Materialized view for queries - **OrderHistory**: Complete timeline for audit - **DailyOrderMetrics**: Analytics aggregation ### Technology - Event Store: EventStoreDB (purpose-built, handles projections) - Alternative considered: Kafka + custom projection service ## Drawbacks - Learning curve for team - Increased complexity vs. CRUD - Need to design events carefully (immutable once stored) - Storage growth (events never deleted) ## Alternatives 1. **Audit tables**: Simpler but doesn't enable temporal queries 2. **CDC from existing DB**: Complex, doesn't change data model 3. **Hybrid**: Event source only for order state changes ## Unresolved Questions - [ ] Event schema versioning strategy - [ ] Retention policy for events - [ ] Snapshot frequency for performance ## Implementation Plan 1. Prototype with single order type (2 weeks) 2. Team training on event sourcing (1 week) 3. Full implementation and migration (4 weeks) 4. Monitoring and optimization (ongoing) ## References - [Event Sourcing by Martin Fowler](https://martinfowler.com/eaaDev/EventSourcing.html) - [EventStoreDB Documentation](https://www.eventstore.com/docs) ``` ## ADR Management ### Directory Structure ``` docs/ β”œβ”€β”€ adr/ β”‚ β”œβ”€β”€ README.md # Index and guidelines β”‚ β”œβ”€β”€ template.md # Team's ADR template β”‚ β”œβ”€β”€ 0001-use-postgresql.md β”‚ β”œβ”€β”€ 0002-caching-strategy.md β”‚ β”œβ”€β”€ 0003-mongodb-user-profiles.md # [DEPRECATED] β”‚ └── 0020-deprecate-mongodb.md # Supersedes 0003 ``` ### ADR Index (README.md) ```markdown # Architecture Decision Records This directory contains Architecture Decision Records (ADRs) for [Project Name]. ## Index | ADR | Title | Status | Date | | ------------------------------------- | ---------------------------------- | ---------- | ---------- | | [0001](0001-use-postgresql.md) | Use PostgreSQL as Primary Database | Accepted | 2024-01-10 | | [0002](0002-caching-strategy.md) | Caching Strategy with Redis | Accepted | 2024-01-12 | | [0003](0003-mongodb-user-profiles.md) | MongoDB for User Profiles | Deprecated | 2023-06-15 | | [0020](0020-deprecate-mongodb.md) | Deprecate MongoDB | Accepted | 2024-01-15 | ## Creating a New ADR 1. Copy `template.md` to `NNNN-title-with-dashes.md` 2. Fill in the template 3. Submit PR for review 4. Update this index after approval ## ADR Status - **Proposed**: Under discussion - **Accepted**: Decision made, implementing - **Deprecated**: No longer relevant - **Superseded**: Replaced by another ADR - **Rejected**: Considered but not adopted ``` ### Automation (adr-tools) ```bash # Install adr-tools brew install adr-tools # Initialize ADR directory adr init docs/adr # Create new ADR adr new "Use PostgreSQL as Primary Database" # Supersede an ADR adr new -s 3 "Deprecate MongoDB in Favor of PostgreSQL" # Generate table of contents adr generate toc > docs/adr/README.md # Link related ADRs adr link 2 "Complements" 1 "Is complemented by" ``` ## Review Process ```markdown ## ADR Review Checklist ### Before Submission - [ ] Context clearly explains the problem - [ ] All viable options considered - [ ] Pros/cons balanced and honest - [ ] Consequences (positive and negative) documented - [ ] Related ADRs linked ### During Review - [ ] At least 2 senior engineers reviewed - [ ] Affected teams consulted - [ ] Security implications considered - [ ] Cost implications documented - [ ] Reversibility assessed ### After Acceptance - [ ] ADR index updated - [ ] Team notified - [ ] Implementation tickets created - [ ] Related documentation updated ``` ## Best Practices ### Do's - **Write ADRs early** - Before implementation starts - **Keep them short** - 1-2 pages maximum - **Be honest about trade-offs** - Include real cons - **Link related decisions** - Build decision graph - **Update status** - Deprecate when superseded ### Don'ts - **Don't change accepted ADRs** - Write new ones to supersede - **Don't skip context** - Future readers need background - **Don't hide failures** - Rejected decisions are valuable - **Don't be vague** - Specific decisions, specific consequences - **Don't forget implementation** - ADR without action is waste ## Resources - [Documenting Architecture Decisions (Michael Nygard)](https://cognitect.com/blog/2011/11/15/documenting-architecture-decisions) - [MADR Template](https://adr.github.io/madr/) - [ADR GitHub Organization](https://adr.github.io/) - [adr-tools](https://github.com/npryce/adr-tools)
πŸ‘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

rag-implementation

Build Retrieval-Augmented Generation (RAG) systems for LLM

coding
⭐1
# RAG Implementation Master Retrieval-Augmented Generation (RAG) to build LLM applications that provide accurate, grounded responses using external knowledge sources. ## When to Use This Skill - Building Q&A systems over proprietary documents - Creating chatbots with current, factual information - Implementing semantic search with natural language queries - Reducing hallucinations with grounded responses - Enabling LLMs to access domain-specific knowledge - Building documentation assistants - Creating research tools with source citation ## Core Components ### 1. Vector Databases **Purpose**: Store and retrieve document embeddings efficiently **Options:** - **Pinecone**: Managed, scalable, serverless - **Weaviate**: Open-source, hybrid search, GraphQL - **Milvus**: High performance, on-premise - **Chroma**: Lightweight, easy to use, local development - **Qdrant**: Fast, filtered search, Rust-based - **pgvector**: PostgreSQL extension, SQL integration ### 2. Embeddings **Purpose**: Convert text to numerical vectors for similarity search **Models (2026):** | Model | Dimensions | Best For | |-------|------------|----------| | **voyage-3-large** | 1024 | Claude apps (Anthropic recommended) | | **voyage-code-3** | 1024 | Code search | | **text-embedding-3-large** | 3072 | OpenAI apps, high accuracy | | **text-embedding-3-small** | 1536 | OpenAI apps, cost-effective | | **bge-large-en-v1.5** | 1024 | Open source, local deployment | | **multilingual-e5-large** | 1024 | Multi-language support | ### 3. Retrieval Strategies **Approaches:** - **Dense Retrieval**: Semantic similarity via embeddings - **Sparse Retrieval**: Keyword matching (BM25, TF-IDF) - **Hybrid Search**: Combine dense + sparse with weighted fusion - **Multi-Query**: Generate multiple query variations - **HyDE**: Generate hypothetical documents for better retrieval ### 4. Reranking **Purpose**: Improve retrieval quality by reordering results **Methods:** - **Cross-Encoders**: BERT-based reranking (ms-marco-MiniLM) - **Cohere Rerank**: API-based reranking - **Maximal Marginal Relevance (MMR)**: Diversity + relevance - **LLM-based**: Use LLM to score relevance ## Quick Start with LangGraph ```python from langgraph.graph import StateGraph, START, END from langchain_anthropic import ChatAnthropic from langchain_voyageai import VoyageAIEmbeddings from langchain_pinecone import PineconeVectorStore from langchain_core.documents import Document from langchain_core.prompts import ChatPromptTemplate from langchain_text_splitters import RecursiveCharacterTextSplitter from typing import TypedDict, Annotated class RAGState(TypedDict): question: str context: list[Document] answer: str # Initialize components llm = ChatAnthropic(model="claude-sonnet-4-6") embeddings = VoyageAIEmbeddings(model="voyage-3-large") vectorstore = PineconeVectorStore(index_name="docs", embedding=embeddings) retriever = vectorstore.as_retriever(search_kwargs={"k": 4}) # RAG prompt rag_prompt = ChatPromptTemplate.from_template( """Answer based on the context below. If you cannot answer, say so. Context: {context} Question: {question} Answer:""" ) async def retrieve(state: RAGState) -> RAGState: """Retrieve relevant documents.""" docs = await retriever.ainvoke(state["question"]) return {"context": docs} async def generate(state: RAGState) -> RAGState: """Generate answer from context.""" context_text = "\n\n".join(doc.page_content for doc in state["context"]) messages = rag_prompt.format_messages( context=context_text, question=state["question"] ) response = await llm.ainvoke(messages) return {"answer": response.content} # Build RAG graph builder = StateGraph(RAGState) builder.add_node("retrieve", retrieve) builder.add_node("generate", generate) builder.add_edge(START, "retrieve") builder.add_edge("retrieve", "generate") builder.add_edge("generate", END) rag_chain = builder.compile() # Use result = await rag_chain.ainvoke({"question": "What are the main features?"}) print(result["answer"]) ``` ## Advanced RAG Patterns ### Pattern 1: Hybrid Search with RRF ```python from langchain_community.retrievers import BM25Retriever from langchain.retrievers import EnsembleRetriever # Sparse retriever (BM25 for keyword matching) bm25_retriever = BM25Retriever.from_documents(documents) bm25_retriever.k = 10 # Dense retriever (embeddings for semantic search) dense_retriever = vectorstore.as_retriever(search_kwargs={"k": 10}) # Combine with Reciprocal Rank Fusion weights ensemble_retriever = EnsembleRetriever( retrievers=[bm25_retriever, dense_retriever], weights=[0.3, 0.7] # 30% keyword, 70% semantic ) ``` ### Pattern 2: Multi-Query Retrieval ```python from langchain.retrievers.multi_query import MultiQueryRetriever # Generate multiple query perspectives for better recall multi_query_retriever = MultiQueryRetriever.from_llm( retriever=vectorstore.as_retriever(search_kwargs={"k": 5}), llm=llm ) # Single query β†’ multiple variations β†’ combined results results = await multi_query_retriever.ainvoke("What is the main topic?") ``` ### Pattern 3: Contextual Compression ```python from langchain.retrievers import ContextualCompressionRetriever from langchain.retrievers.document_compressors import LLMChainExtractor # Compressor extracts only relevant portions compressor = LLMChainExtractor.from_llm(llm) compression_retriever = ContextualCompressionRetriever( base_compressor=compressor, base_retriever=vectorstore.as_retriever(search_kwargs={"k": 10}) ) # Returns only relevant parts of documents compressed_docs = await compression_retriever.ainvoke("specific query") ``` ### Pattern 4: Parent Document Retriever ```python from langchain.retrievers import ParentDocumentRetriever from langchain.storage import InMemoryStore from langchain_text_splitters import RecursiveCharacterTextSplitter # Small chunks for precise retrieval, large chunks for context child_splitter = RecursiveCharacterTextSplitter(chunk_size=400, chunk_overlap=50) parent_splitter = RecursiveCharacterTextSplitter(chunk_size=2000, chunk_overlap=200) # Store for parent documents docstore = InMemoryStore() parent_retriever = ParentDocumentRetriever( vectorstore=vectorstore, docstore=docstore, child_splitter=child_splitter, parent_splitter=parent_splitter ) # Add documents (splits children, stores parents) await parent_retriever.aadd_documents(documents) # Retrieval returns parent documents with full context results = await parent_retriever.ainvoke("query") ``` ### Pattern 5: HyDE (Hypothetical Document Embeddings) ```python from langchain_core.prompts import ChatPromptTemplate class HyDEState(TypedDict): question: str hypothetical_doc: str context: list[Document] answer: str hyde_prompt = ChatPromptTemplate.from_template( """Write a detailed passage that would answer this question: Question: {question} Passage:""" ) async def generate_hypothetical(state: HyDEState) -> HyDEState: """Generate hypothetical document for better retrieval.""" messages = hyde_prompt.format_messages(question=state["question"]) response = await llm.ainvoke(messages) return {"hypothetical_doc": response.content} async def retrieve_with_hyde(state: HyDEState) -> HyDEState: """Retrieve using hypothetical document.""" # Use hypothetical doc for retrieval instead of original query docs = await retriever.ainvoke(state["hypothetical_doc"]) return {"context": docs} # Build HyDE RAG graph builder = StateGraph(HyDEState) builder.add_node("hypothetical", generate_hypothetical) builder.add_node("retrieve", retrieve_with_hyde) builder.add_node("generate", generate) builder.add_edge(START, "hypothetical") builder.add_edge("hypothetical", "retrieve") builder.add_edge("retrieve", "generate") builder.add_edge("generate", END) hyde_rag = builder.compile() ``` ## Document Chunking Strategies ### Recursive Character Text Splitter ```python from langchain_text_splitters import RecursiveCharacterTextSplitter splitter = RecursiveCharacterTextSplitter( chunk_size=1000, chunk_overlap=200, length_function=len, separators=["\n\n", "\n", ". ", " ", ""] # Try in order ) chunks = splitter.split_documents(documents) ``` ### Token-Based Splitting ```python from langchain_text_splitters import TokenTextSplitter splitter = TokenTextSplitter( chunk_size=512, chunk_overlap=50, encoding_name="cl100k_base" # OpenAI tiktoken encoding ) ``` ### Semantic Chunking ```python from langchain_experimental.text_splitter import SemanticChunker splitter = SemanticChunker( embeddings=embeddings, breakpoint_threshold_type="percentile", breakpoint_threshold_amount=95 ) ``` ### Markdown Header Splitter ```python from langchain_text_splitters import MarkdownHeaderTextSplitter headers_to_split_on = [ ("#", "Header 1"), ("##", "Header 2"), ("###", "Header 3"), ] splitter = MarkdownHeaderTextSplitter( headers_to_split_on=headers_to_split_on, strip_headers=False ) ``` ## Vector Store Configurations ### Pinecone (Serverless) ```python from pinecone import Pinecone, ServerlessSpec from langchain_pinecone import PineconeVectorStore # Initialize Pinecone client pc = Pinecone(api_key=os.environ["PINECONE_API_KEY"]) # Create index if needed if "my-index" not in pc.list_indexes().names(): pc.create_index( name="my-index", dimension=1024, # voyage-3-large dimensions metric="cosine", spec=ServerlessSpec(cloud="aws", region="us-east-1") ) # Create vector store index = pc.Index("my-index") vectorstore = PineconeVectorStore(index=index, embedding=embeddings) ``` ### Weaviate ```python import weaviate from langchain_weaviate import WeaviateVectorStore client = weaviate.connect_to_local() # or connect_to_weaviate_cloud() vectorstore = WeaviateVectorStore( client=client, index_name="Documents", text_key="content", embedding=embeddings ) ``` ### Chroma (Local Development) ```python from langchain_chroma import Chroma vectorstore = Chroma( collection_name="my_collection", embedding_function=embeddings, persist_directory="./chroma_db" ) ``` ### pgvector (PostgreSQL) ```python from langchain_postgres.vectorstores import PGVector connection_string = "postgresql+psycopg://user:pass@localhost:5432/vectordb" vectorstore = PGVector( embeddings=embeddings, collection_name="documents", connection=connection_string, ) ``` ## Retrieval Optimization ### 1. Metadata Filtering ```python from langchain_core.documents import Document # Add metadata during indexing docs_with_metadata = [] for doc in documents: doc.metadata.update({ "source": doc.metadata.get("source", "unknown"), "category": determine_category(doc.page_content), "date": datetime.now().isoformat() }) docs_with_metadata.append(doc) # Filter during retrieval results = await vectorstore.asimilarity_search( "query", filter={"category": "technical"}, k=5 ) ``` ### 2. Maximal Marginal Relevance (MMR) ```python # Balance relevance with diversity results = await vectorstore.amax_marginal_relevance_search( "query", k=5, fetch_k=20, # Fetch 20, return top 5 diverse lambda_mult=0.5 # 0=max diversity, 1=max relevance ) ``` ### 3. Reranking with Cross-Encoder ```python from sentence_transformers import CrossEncoder reranker = CrossEncoder('cross-encoder/ms-marco-MiniLM-L-6-v2') async def retrieve_and_rerank(query: str, k: int = 5) -> list[Document]: # Get initial results candidates = await vectorstore.asimilarity_search(query, k=20) # Rerank pairs = [[query, doc.page_content] for doc in candidates] scores = reranker.predict(pairs) # Sort by score and take top k ranked = sorted(zip(candidates, scores), key=lambda x: x[1], reverse=True) return [doc for doc, score in ranked[:k]] ``` ### 4. Cohere Rerank ```python from langchain.retrievers import CohereRerank from langchain_cohere import CohereRerank reranker = CohereRerank(model="rerank-english-v3.0", top_n=5) # Wrap retriever with reranking reranked_retriever = ContextualCompressionRetriever( base_compressor=reranker, base_retriever=vectorstore.as_retriever(search_kwargs={"k": 20}) ) ``` ## Prompt Engineering for RAG ### Contextual Prompt with Citations ```python rag_prompt = ChatPromptTemplate.from_template( """Answer the question based on the context below. Include citations using [1], [2], etc. If you cannot answer based on the context, say "I don't have enough information." Context: {context} Question: {question} Instructions: 1. Use only information from the context 2. Cite sources with [1], [2] format 3. If uncertain, express uncertainty Answer (with citations):""" ) ``` ### Structured Output for RAG ```python from pydantic import BaseModel, Field class RAGResponse(BaseModel): answer: str = Field(description="The answer based on context") confidence: float = Field(description="Confidence score 0-1") sources: list[str] = Field(description="Source document IDs used") reasoning: str = Field(description="Brief reasoning for the answer") # Use with structured output structured_llm = llm.with_structured_output(RAGResponse) ``` ## Evaluation Metrics ```python from typing import TypedDict class RAGEvalMetrics(TypedDict): retrieval_precision: float # Relevant docs / retrieved docs retrieval_recall: float # Retrieved relevant / total relevant answer_relevance: float # Answer addresses question faithfulness: float # Answer grounded in context context_relevance: float # Context relevant to question async def evaluate_rag_system( rag_chain, test_cases: list[dict] ) -> RAGEvalMetrics: """Evaluate RAG system on test cases.""" metrics = {k: [] for k in RAGEvalMetrics.__annotations__} for test in test_cases: result = await rag_chain.ainvoke({"question": test["question"]}) # Retrieval metrics retrieved_ids = {doc.metadata["id"] for doc in result["context"]} relevant_ids = set(test["relevant_doc_ids"]) precision = len(retrieved_ids & relevant_ids) / len(retrieved_ids) recall = len(retrieved_ids & relevant_ids) / len(relevant_ids) metrics["retrieval_precision"].append(precision) metrics["retrieval_recall"].append(recall) # Use LLM-as-judge for quality metrics quality = await evaluate_answer_quality( question=test["question"], answer=result["answer"], context=result["context"], expected=test.get("expected_answer") ) metrics["answer_relevance"].append(quality["relevance"]) metrics["faithfulness"].append(quality["faithfulness"]) metrics["context_relevance"].append(quality["context_relevance"]) return {k: sum(v) / len(v) for k, v in metrics.items()} ``` ## Resources - [LangChain RAG Tutorial](https://python.langchain.com/docs/tutorials/rag/) - [LangGraph RAG Examples](https://langchain-ai.github.io/langgraph/tutorials/rag/) - [Pinecone Best Practices](https://docs.pinecone.io/guides/get-started/overview) - [Voyage AI Embeddings](https://docs.voyageai.com/) - [RAG Evaluation Guide](https://docs.ragas.io/) ## Best Practices 1. **Chunk Size**: Balance between context (larger) and specificity (smaller) - typically 500-1000 tokens 2. **Overlap**: Use 10-20% overlap to preserve context at boundaries 3. **Metadata**: Include source, page, timestamp for filtering and debugging 4. **Hybrid Search**: Combine semantic and keyword search for best recall 5. **Reranking**: Use cross-encoder reranking for precision-critical applications 6. **Citations**: Always return source documents for transparency 7. **Evaluation**: Continuously test retrieval quality and answer accuracy 8. **Monitoring**: Track retrieval metrics and latency in production ## Common Issues - **Poor Retrieval**: Check embedding quality, chunk size, query formulation - **Irrelevant Results**: Add metadata filtering, use hybrid search, rerank - **Missing Information**: Ensure documents are properly indexed, check chunking - **Slow Queries**: Optimize vector store, use caching, reduce k - **Hallucinations**: Improve grounding prompt, add verification step - **Context Too Long**: Use compression or parent document retriever
πŸ‘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

slo-implementation

Define and implement Service Level Indicators (SLIs) and Service

coding
⭐1
# SLO Implementation Framework for defining and implementing Service Level Indicators (SLIs), Service Level Objectives (SLOs), and error budgets. ## Purpose Implement measurable reliability targets using SLIs, SLOs, and error budgets to balance reliability with innovation velocity. ## When to Use - Define service reliability targets - Measure user-perceived reliability - Implement error budgets - Create SLO-based alerts - Track reliability goals ## SLI/SLO/SLA Hierarchy ``` SLA (Service Level Agreement) ↓ Contract with customers SLO (Service Level Objective) ↓ Internal reliability target SLI (Service Level Indicator) ↓ Actual measurement ``` ## Defining SLIs ### Common SLI Types #### 1. Availability SLI ```promql # Successful requests / Total requests sum(rate(http_requests_total{status!~"5.."}[28d])) / sum(rate(http_requests_total[28d])) ``` #### 2. Latency SLI ```promql # Requests below latency threshold / Total requests sum(rate(http_request_duration_seconds_bucket{le="0.5"}[28d])) / sum(rate(http_request_duration_seconds_count[28d])) ``` #### 3. Durability SLI ``` # Successful writes / Total writes sum(storage_writes_successful_total) / sum(storage_writes_total) ``` **Reference:** See `references/slo-definitions.md` ## Setting SLO Targets ### Availability SLO Examples | SLO % | Downtime/Month | Downtime/Year | | ------ | -------------- | ------------- | | 99% | 7.2 hours | 3.65 days | | 99.9% | 43.2 minutes | 8.76 hours | | 99.95% | 21.6 minutes | 4.38 hours | | 99.99% | 4.32 minutes | 52.56 minutes | ### Choose Appropriate SLOs **Consider:** - User expectations - Business requirements - Current performance - Cost of reliability - Competitor benchmarks **Example SLOs:** ```yaml slos: - name: api_availability target: 99.9 window: 28d sli: | sum(rate(http_requests_total{status!~"5.."}[28d])) / sum(rate(http_requests_total[28d])) - name: api_latency_p95 target: 99 window: 28d sli: | sum(rate(http_request_duration_seconds_bucket{le="0.5"}[28d])) / sum(rate(http_request_duration_seconds_count[28d])) ``` ## Error Budget Calculation ### Error Budget Formula ``` Error Budget = 1 - SLO Target ``` **Example:** - SLO: 99.9% availability - Error Budget: 0.1% = 43.2 minutes/month - Current Error: 0.05% = 21.6 minutes/month - Remaining Budget: 50% ### Error Budget Policy ```yaml error_budget_policy: - remaining_budget: 100% action: Normal development velocity - remaining_budget: 50% action: Consider postponing risky changes - remaining_budget: 10% action: Freeze non-critical changes - remaining_budget: 0% action: Feature freeze, focus on reliability ``` **Reference:** See `references/error-budget.md` ## SLO Implementation ### Prometheus Recording Rules ```yaml # SLI Recording Rules groups: - name: sli_rules interval: 30s rules: # Availability SLI - record: sli:http_availability:ratio expr: | sum(rate(http_requests_total{status!~"5.."}[28d])) / sum(rate(http_requests_total[28d])) # Latency SLI (requests < 500ms) - record: sli:http_latency:ratio expr: | sum(rate(http_request_duration_seconds_bucket{le="0.5"}[28d])) / sum(rate(http_request_duration_seconds_count[28d])) - name: slo_rules interval: 5m rules: # SLO compliance (1 = meeting SLO, 0 = violating) - record: slo:http_availability:compliance expr: sli:http_availability:ratio >= bool 0.999 - record: slo:http_latency:compliance expr: sli:http_latency:ratio >= bool 0.99 # Error budget remaining (percentage) - record: slo:http_availability:error_budget_remaining expr: | (sli:http_availability:ratio - 0.999) / (1 - 0.999) * 100 # Error budget burn rate - record: slo:http_availability:burn_rate_5m expr: | (1 - ( sum(rate(http_requests_total{status!~"5.."}[5m])) / sum(rate(http_requests_total[5m])) )) / (1 - 0.999) ``` ### SLO Alerting Rules ```yaml groups: - name: slo_alerts interval: 1m rules: # Fast burn: 14.4x rate, 1 hour window # Consumes 2% error budget in 1 hour - alert: SLOErrorBudgetBurnFast expr: | slo:http_availability:burn_rate_1h > 14.4 and slo:http_availability:burn_rate_5m > 14.4 for: 2m labels: severity: critical annotations: summary: "Fast error budget burn detected" description: "Error budget burning at {{ $value }}x rate" # Slow burn: 6x rate, 6 hour window # Consumes 5% error budget in 6 hours - alert: SLOErrorBudgetBurnSlow expr: | slo:http_availability:burn_rate_6h > 6 and slo:http_availability:burn_rate_30m > 6 for: 15m labels: severity: warning annotations: summary: "Slow error budget burn detected" description: "Error budget burning at {{ $value }}x rate" # Error budget exhausted - alert: SLOErrorBudgetExhausted expr: slo:http_availability:error_budget_remaining < 0 for: 5m labels: severity: critical annotations: summary: "SLO error budget exhausted" description: "Error budget remaining: {{ $value }}%" ``` ## SLO Dashboard **Grafana Dashboard Structure:** ``` β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ SLO Compliance (Current) β”‚ β”‚ βœ“ 99.95% (Target: 99.9%) β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ Error Budget Remaining: 65% β”‚ β”‚ β–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–ˆβ–‘β–‘ 65% β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ SLI Trend (28 days) β”‚ β”‚ [Time series graph] β”‚ β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€ β”‚ Burn Rate Analysis β”‚ β”‚ [Burn rate by time window] β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` **Example Queries:** ```promql # Current SLO compliance sli:http_availability:ratio * 100 # Error budget remaining slo:http_availability:error_budget_remaining # Days until error budget exhausted (at current burn rate) (slo:http_availability:error_budget_remaining / 100) * 28 / (1 - sli:http_availability:ratio) * (1 - 0.999) ``` ## Multi-Window Burn Rate Alerts ```yaml # Combination of short and long windows reduces false positives rules: - alert: SLOBurnRateHigh expr: | ( slo:http_availability:burn_rate_1h > 14.4 and slo:http_availability:burn_rate_5m > 14.4 ) or ( slo:http_availability:burn_rate_6h > 6 and slo:http_availability:burn_rate_30m > 6 ) labels: severity: critical ``` ## SLO Review Process ### Weekly Review - Current SLO compliance - Error budget status - Trend analysis - Incident impact ### Monthly Review - SLO achievement - Error budget usage - Incident postmortems - SLO adjustments ### Quarterly Review - SLO relevance - Target adjustments - Process improvements - Tooling enhancements ## Best Practices 1. **Start with user-facing services** 2. **Use multiple SLIs** (availability, latency, etc.) 3. **Set achievable SLOs** (don't aim for 100%) 4. **Implement multi-window alerts** to reduce noise 5. **Track error budget** consistently 6. **Review SLOs regularly** 7. **Document SLO decisions** 8. **Align with business goals** 9. **Automate SLO reporting** 10. **Use SLOs for prioritization** ## Reference Files - `assets/slo-template.md` - SLO definition template - `references/slo-definitions.md` - SLO definition patterns - `references/error-budget.md` - Error budget calculations ## Related Skills - `prometheus-configuration` - For metric collection - `grafana-dashboards` - For SLO visualization
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

billing-automation

Build automated billing systems for recurring payments, invoicing,

business
⭐1
# Billing Automation Master automated billing systems including recurring billing, invoice generation, dunning management, proration, and tax calculation. ## When to Use This Skill - Implementing SaaS subscription billing - Automating invoice generation and delivery - Managing failed payment recovery (dunning) - Calculating prorated charges for plan changes - Handling sales tax, VAT, and GST - Processing usage-based billing - Managing billing cycles and renewals ## Core Concepts ### 1. Billing Cycles **Common Intervals:** - Monthly (most common for SaaS) - Annual (discounted long-term) - Quarterly - Weekly - Custom (usage-based, per-seat) ### 2. Subscription States ``` trial β†’ active β†’ past_due β†’ canceled β†’ paused β†’ resumed ``` ### 3. Dunning Management Automated process to recover failed payments through: - Retry schedules - Customer notifications - Grace periods - Account restrictions ### 4. Proration Adjusting charges when: - Upgrading/downgrading mid-cycle - Adding/removing seats - Changing billing frequency ## Quick Start ```python from billing import BillingEngine, Subscription # Initialize billing engine billing = BillingEngine() # Create subscription subscription = billing.create_subscription( customer_id="cus_123", plan_id="plan_pro_monthly", billing_cycle_anchor=datetime.now(), trial_days=14 ) # Process billing cycle billing.process_billing_cycle(subscription.id) ``` ## Subscription Lifecycle Management ```python from datetime import datetime, timedelta from enum import Enum class SubscriptionStatus(Enum): TRIAL = "trial" ACTIVE = "active" PAST_DUE = "past_due" CANCELED = "canceled" PAUSED = "paused" class Subscription: def __init__(self, customer_id, plan, billing_cycle_day=None): self.id = generate_id() self.customer_id = customer_id self.plan = plan self.status = SubscriptionStatus.TRIAL self.current_period_start = datetime.now() self.current_period_end = self.current_period_start + timedelta(days=plan.trial_days or 30) self.billing_cycle_day = billing_cycle_day or self.current_period_start.day self.trial_end = datetime.now() + timedelta(days=plan.trial_days) if plan.trial_days else None def start_trial(self, trial_days): """Start trial period.""" self.status = SubscriptionStatus.TRIAL self.trial_end = datetime.now() + timedelta(days=trial_days) self.current_period_end = self.trial_end def activate(self): """Activate subscription after trial or immediately.""" self.status = SubscriptionStatus.ACTIVE self.current_period_start = datetime.now() self.current_period_end = self.calculate_next_billing_date() def mark_past_due(self): """Mark subscription as past due after failed payment.""" self.status = SubscriptionStatus.PAST_DUE # Trigger dunning workflow def cancel(self, at_period_end=True): """Cancel subscription.""" if at_period_end: self.cancel_at_period_end = True # Will cancel when current period ends else: self.status = SubscriptionStatus.CANCELED self.canceled_at = datetime.now() def calculate_next_billing_date(self): """Calculate next billing date based on interval.""" if self.plan.interval == 'month': return self.current_period_start + timedelta(days=30) elif self.plan.interval == 'year': return self.current_period_start + timedelta(days=365) elif self.plan.interval == 'week': return self.current_period_start + timedelta(days=7) ``` ## Billing Cycle Processing ```python class BillingEngine: def process_billing_cycle(self, subscription_id): """Process billing for a subscription.""" subscription = self.get_subscription(subscription_id) # Check if billing is due if datetime.now() < subscription.current_period_end: return # Generate invoice invoice = self.generate_invoice(subscription) # Attempt payment payment_result = self.charge_customer( subscription.customer_id, invoice.total ) if payment_result.success: # Payment successful invoice.mark_paid() subscription.advance_billing_period() self.send_invoice(invoice) else: # Payment failed subscription.mark_past_due() self.start_dunning_process(subscription, invoice) def generate_invoice(self, subscription): """Generate invoice for billing period.""" invoice = Invoice( customer_id=subscription.customer_id, subscription_id=subscription.id, period_start=subscription.current_period_start, period_end=subscription.current_period_end ) # Add subscription line item invoice.add_line_item( description=subscription.plan.name, amount=subscription.plan.amount, quantity=subscription.quantity or 1 ) # Add usage-based charges if applicable if subscription.has_usage_billing: usage_charges = self.calculate_usage_charges(subscription) invoice.add_line_item( description="Usage charges", amount=usage_charges ) # Calculate tax tax = self.calculate_tax(invoice.subtotal, subscription.customer) invoice.tax = tax invoice.finalize() return invoice def charge_customer(self, customer_id, amount): """Charge customer using saved payment method.""" customer = self.get_customer(customer_id) try: # Charge using payment processor charge = stripe.Charge.create( customer=customer.stripe_id, amount=int(amount * 100), # Convert to cents currency='usd' ) return PaymentResult(success=True, transaction_id=charge.id) except stripe.error.CardError as e: return PaymentResult(success=False, error=str(e)) ``` ## Dunning Management ```python class DunningManager: """Manage failed payment recovery.""" def __init__(self): self.retry_schedule = [ {'days': 3, 'email_template': 'payment_failed_first'}, {'days': 7, 'email_template': 'payment_failed_reminder'}, {'days': 14, 'email_template': 'payment_failed_final'} ] def start_dunning_process(self, subscription, invoice): """Start dunning process for failed payment.""" dunning_attempt = DunningAttempt( subscription_id=subscription.id, invoice_id=invoice.id, attempt_number=1, next_retry=datetime.now() + timedelta(days=3) ) # Send initial failure notification self.send_dunning_email(subscription, 'payment_failed_first') # Schedule retries self.schedule_retries(dunning_attempt) def retry_payment(self, dunning_attempt): """Retry failed payment.""" subscription = self.get_subscription(dunning_attempt.subscription_id) invoice = self.get_invoice(dunning_attempt.invoice_id) # Attempt payment again result = self.charge_customer(subscription.customer_id, invoice.total) if result.success: # Payment succeeded invoice.mark_paid() subscription.status = SubscriptionStatus.ACTIVE self.send_dunning_email(subscription, 'payment_recovered') dunning_attempt.mark_resolved() else: # Still failing dunning_attempt.attempt_number += 1 if dunning_attempt.attempt_number < len(self.retry_schedule): # Schedule next retry next_retry_config = self.retry_schedule[dunning_attempt.attempt_number] dunning_attempt.next_retry = datetime.now() + timedelta(days=next_retry_config['days']) self.send_dunning_email(subscription, next_retry_config['email_template']) else: # Exhausted retries, cancel subscription subscription.cancel(at_period_end=False) self.send_dunning_email(subscription, 'subscription_canceled') def send_dunning_email(self, subscription, template): """Send dunning notification to customer.""" customer = self.get_customer(subscription.customer_id) email_content = self.render_template(template, { 'customer_name': customer.name, 'amount_due': subscription.plan.amount, 'update_payment_url': f"https://app.example.com/billing" }) send_email( to=customer.email, subject=email_content['subject'], body=email_content['body'] ) ``` ## Proration ```python class ProrationCalculator: """Calculate prorated charges for plan changes.""" @staticmethod def calculate_proration(old_plan, new_plan, period_start, period_end, change_date): """Calculate proration for plan change.""" # Days in current period total_days = (period_end - period_start).days # Days used on old plan days_used = (change_date - period_start).days # Days remaining on new plan days_remaining = (period_end - change_date).days # Calculate prorated amounts unused_amount = (old_plan.amount / total_days) * days_remaining new_plan_amount = (new_plan.amount / total_days) * days_remaining # Net charge/credit proration = new_plan_amount - unused_amount return { 'old_plan_credit': -unused_amount, 'new_plan_charge': new_plan_amount, 'net_proration': proration, 'days_used': days_used, 'days_remaining': days_remaining } @staticmethod def calculate_seat_proration(current_seats, new_seats, price_per_seat, period_start, period_end, change_date): """Calculate proration for seat changes.""" total_days = (period_end - period_start).days days_remaining = (period_end - change_date).days # Additional seats charge additional_seats = new_seats - current_seats prorated_amount = (additional_seats * price_per_seat / total_days) * days_remaining return { 'additional_seats': additional_seats, 'prorated_charge': max(0, prorated_amount), # No refund for removing seats mid-cycle 'effective_date': change_date } ``` ## Tax Calculation ```python class TaxCalculator: """Calculate sales tax, VAT, GST.""" def __init__(self): # Tax rates by region self.tax_rates = { 'US_CA': 0.0725, # California sales tax 'US_NY': 0.04, # New York sales tax 'GB': 0.20, # UK VAT 'DE': 0.19, # Germany VAT 'FR': 0.20, # France VAT 'AU': 0.10, # Australia GST } def calculate_tax(self, amount, customer): """Calculate applicable tax.""" # Determine tax jurisdiction jurisdiction = self.get_tax_jurisdiction(customer) if not jurisdiction: return 0 # Get tax rate tax_rate = self.tax_rates.get(jurisdiction, 0) # Calculate tax tax = amount * tax_rate return { 'tax_amount': tax, 'tax_rate': tax_rate, 'jurisdiction': jurisdiction, 'tax_type': self.get_tax_type(jurisdiction) } def get_tax_jurisdiction(self, customer): """Determine tax jurisdiction based on customer location.""" if customer.country == 'US': # US: Tax based on customer state return f"US_{customer.state}" elif customer.country in ['GB', 'DE', 'FR']: # EU: VAT return customer.country elif customer.country == 'AU': # Australia: GST return 'AU' else: return None def get_tax_type(self, jurisdiction): """Get type of tax for jurisdiction.""" if jurisdiction.startswith('US_'): return 'Sales Tax' elif jurisdiction in ['GB', 'DE', 'FR']: return 'VAT' elif jurisdiction == 'AU': return 'GST' return 'Tax' def validate_vat_number(self, vat_number, country): """Validate EU VAT number.""" # Use VIES API for validation # Returns True if valid, False otherwise pass ``` ## Invoice Generation ```python class Invoice: def __init__(self, customer_id, subscription_id=None): self.id = generate_invoice_number() self.customer_id = customer_id self.subscription_id = subscription_id self.status = 'draft' self.line_items = [] self.subtotal = 0 self.tax = 0 self.total = 0 self.created_at = datetime.now() def add_line_item(self, description, amount, quantity=1): """Add line item to invoice.""" line_item = { 'description': description, 'unit_amount': amount, 'quantity': quantity, 'total': amount * quantity } self.line_items.append(line_item) self.subtotal += line_item['total'] def finalize(self): """Finalize invoice and calculate total.""" self.total = self.subtotal + self.tax self.status = 'open' self.finalized_at = datetime.now() def mark_paid(self): """Mark invoice as paid.""" self.status = 'paid' self.paid_at = datetime.now() def to_pdf(self): """Generate PDF invoice.""" from reportlab.pdfgen import canvas # Generate PDF # Include: company info, customer info, line items, tax, total pass def to_html(self): """Generate HTML invoice.""" template = """ <!DOCTYPE html> <html> <head><title>Invoice #{invoice_number}</title></head> <body> <h1>Invoice #{invoice_number}</h1> <p>Date: {date}</p> <h2>Bill To:</h2> <p>{customer_name}<br>{customer_address}</p> <table> <tr><th>Description</th><th>Quantity</th><th>Amount</th></tr> {line_items} </table> <p>Subtotal: ${subtotal}</p> <p>Tax: ${tax}</p> <h3>Total: ${total}</h3> </body> </html> """ return template.format( invoice_number=self.id, date=self.created_at.strftime('%Y-%m-%d'), customer_name=self.customer.name, customer_address=self.customer.address, line_items=self.render_line_items(), subtotal=self.subtotal, tax=self.tax, total=self.total ) ``` ## Usage-Based Billing ```python class UsageBillingEngine: """Track and bill for usage.""" def track_usage(self, customer_id, metric, quantity): """Track usage event.""" UsageRecord.create( customer_id=customer_id, metric=metric, quantity=quantity, timestamp=datetime.now() ) def calculate_usage_charges(self, subscription, period_start, period_end): """Calculate charges for usage in billing period.""" usage_records = UsageRecord.get_for_period( subscription.customer_id, period_start, period_end ) total_usage = sum(record.quantity for record in usage_records) # Tiered pricing if subscription.plan.pricing_model == 'tiered': charge = self.calculate_tiered_pricing(total_usage, subscription.plan.tiers) # Per-unit pricing elif subscription.plan.pricing_model == 'per_unit': charge = total_usage * subscription.plan.unit_price # Volume pricing elif subscription.plan.pricing_model == 'volume': charge = self.calculate_volume_pricing(total_usage, subscription.plan.tiers) return charge def calculate_tiered_pricing(self, total_usage, tiers): """Calculate cost using tiered pricing.""" charge = 0 remaining = total_usage for tier in sorted(tiers, key=lambda x: x['up_to']): tier_usage = min(remaining, tier['up_to'] - tier['from']) charge += tier_usage * tier['unit_price'] remaining -= tier_usage if remaining <= 0: break return charge ``` ## Resources - **references/billing-cycles.md**: Billing cycle management - **references/dunning-management.md**: Failed payment recovery - **references/proration.md**: Prorated charge calculations - **references/tax-calculation.md**: Tax/VAT/GST handling - **references/invoice-lifecycle.md**: Invoice state management - **assets/billing-state-machine.yaml**: Billing workflow - **assets/invoice-template.html**: Invoice templates - **assets/dunning-policy.yaml**: Dunning configuration ## Best Practices 1. **Automate Everything**: Minimize manual intervention 2. **Clear Communication**: Notify customers of billing events 3. **Flexible Retry Logic**: Balance recovery with customer experience 4. **Accurate Proration**: Fair calculation for plan changes 5. **Tax Compliance**: Calculate correct tax for jurisdiction 6. **Audit Trail**: Log all billing events 7. **Graceful Degradation**: Handle edge cases without breaking ## Common Pitfalls - **Incorrect Proration**: Not accounting for partial periods - **Missing Tax**: Forgetting to add tax to invoices - **Aggressive Dunning**: Canceling too quickly - **No Notifications**: Not informing customers of failures - **Hardcoded Cycles**: Not supporting custom billing dates
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

python-observability

Python observability patterns including structured logging,

coding
⭐1
# Python Observability Instrument Python applications with structured logs, metrics, and traces. When something breaks in production, you need to answer "what, where, and why" without deploying new code. ## When to Use This Skill - Adding structured logging to applications - Implementing metrics collection with Prometheus - Setting up distributed tracing across services - Propagating correlation IDs through request chains - Debugging production issues - Building observability dashboards ## Core Concepts ### 1. Structured Logging Emit logs as JSON with consistent fields for production environments. Machine-readable logs enable powerful queries and alerts. For local development, consider human-readable formats. ### 2. The Four Golden Signals Track latency, traffic, errors, and saturation for every service boundary. ### 3. Correlation IDs Thread a unique ID through all logs and spans for a single request, enabling end-to-end tracing. ### 4. Bounded Cardinality Keep metric label values bounded. Unbounded labels (like user IDs) explode storage costs. ## Quick Start ```python import structlog structlog.configure( processors=[ structlog.processors.TimeStamper(fmt="iso"), structlog.processors.JSONRenderer(), ], ) logger = structlog.get_logger() logger.info("Request processed", user_id="123", duration_ms=45) ``` ## Fundamental Patterns ### Pattern 1: Structured Logging with Structlog Configure structlog for JSON output with consistent fields. ```python import logging import structlog def configure_logging(log_level: str = "INFO") -> None: """Configure structured logging for the application.""" structlog.configure( processors=[ structlog.contextvars.merge_contextvars, structlog.processors.add_log_level, structlog.processors.TimeStamper(fmt="iso"), structlog.processors.StackInfoRenderer(), structlog.processors.format_exc_info, structlog.processors.JSONRenderer(), ], wrapper_class=structlog.make_filtering_bound_logger( getattr(logging, log_level.upper()) ), context_class=dict, logger_factory=structlog.PrintLoggerFactory(), cache_logger_on_first_use=True, ) # Initialize at application startup configure_logging("INFO") logger = structlog.get_logger() ``` ### Pattern 2: Consistent Log Fields Every log entry should include standard fields for filtering and correlation. ```python import structlog from contextvars import ContextVar # Store correlation ID in context correlation_id: ContextVar[str] = ContextVar("correlation_id", default="") logger = structlog.get_logger() def process_request(request: Request) -> Response: """Process request with structured logging.""" logger.info( "Request received", correlation_id=correlation_id.get(), method=request.method, path=request.path, user_id=request.user_id, ) try: result = handle_request(request) logger.info( "Request completed", correlation_id=correlation_id.get(), status_code=200, duration_ms=elapsed, ) return result except Exception as e: logger.error( "Request failed", correlation_id=correlation_id.get(), error_type=type(e).__name__, error_message=str(e), ) raise ``` ### Pattern 3: Semantic Log Levels Use log levels consistently across the application. | Level | Purpose | Examples | |-------|---------|----------| | `DEBUG` | Development diagnostics | Variable values, internal state | | `INFO` | Request lifecycle, operations | Request start/end, job completion | | `WARNING` | Recoverable anomalies | Retry attempts, fallback used | | `ERROR` | Failures needing attention | Exceptions, service unavailable | ```python # DEBUG: Detailed internal information logger.debug("Cache lookup", key=cache_key, hit=cache_hit) # INFO: Normal operational events logger.info("Order created", order_id=order.id, total=order.total) # WARNING: Abnormal but handled situations logger.warning( "Rate limit approaching", current_rate=950, limit=1000, reset_seconds=30, ) # ERROR: Failures requiring investigation logger.error( "Payment processing failed", order_id=order.id, error=str(e), payment_provider="stripe", ) ``` Never log expected behavior at `ERROR`. A user entering a wrong password is `INFO`, not `ERROR`. ### Pattern 4: Correlation ID Propagation Generate a unique ID at ingress and thread it through all operations. ```python from contextvars import ContextVar import uuid import structlog correlation_id: ContextVar[str] = ContextVar("correlation_id", default="") def set_correlation_id(cid: str | None = None) -> str: """Set correlation ID for current context.""" cid = cid or str(uuid.uuid4()) correlation_id.set(cid) structlog.contextvars.bind_contextvars(correlation_id=cid) return cid # FastAPI middleware example from fastapi import Request async def correlation_middleware(request: Request, call_next): """Middleware to set and propagate correlation ID.""" # Use incoming header or generate new cid = request.headers.get("X-Correlation-ID") or str(uuid.uuid4()) set_correlation_id(cid) response = await call_next(request) response.headers["X-Correlation-ID"] = cid return response ``` Propagate to outbound requests: ```python import httpx async def call_downstream_service(endpoint: str, data: dict) -> dict: """Call downstream service with correlation ID.""" async with httpx.AsyncClient() as client: response = await client.post( endpoint, json=data, headers={"X-Correlation-ID": correlation_id.get()}, ) return response.json() ``` ## Advanced Patterns ### Pattern 5: The Four Golden Signals with Prometheus Track these metrics for every service boundary: ```python from prometheus_client import Counter, Histogram, Gauge # Latency: How long requests take REQUEST_LATENCY = Histogram( "http_request_duration_seconds", "Request latency in seconds", ["method", "endpoint", "status"], buckets=[0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10], ) # Traffic: Request rate REQUEST_COUNT = Counter( "http_requests_total", "Total HTTP requests", ["method", "endpoint", "status"], ) # Errors: Error rate ERROR_COUNT = Counter( "http_errors_total", "Total HTTP errors", ["method", "endpoint", "error_type"], ) # Saturation: Resource utilization DB_POOL_USAGE = Gauge( "db_connection_pool_used", "Number of database connections in use", ) ``` Instrument your endpoints: ```python import time from functools import wraps def track_request(func): """Decorator to track request metrics.""" @wraps(func) async def wrapper(request: Request, *args, **kwargs): method = request.method endpoint = request.url.path start = time.perf_counter() try: response = await func(request, *args, **kwargs) status = str(response.status_code) return response except Exception as e: status = "500" ERROR_COUNT.labels( method=method, endpoint=endpoint, error_type=type(e).__name__, ).inc() raise finally: duration = time.perf_counter() - start REQUEST_COUNT.labels(method=method, endpoint=endpoint, status=status).inc() REQUEST_LATENCY.labels(method=method, endpoint=endpoint, status=status).observe(duration) return wrapper ``` ### Pattern 6: Bounded Cardinality Avoid labels with unbounded values to prevent metric explosion. ```python # BAD: User ID has potentially millions of values REQUEST_COUNT.labels(method="GET", user_id=user.id) # Don't do this! # GOOD: Bounded values only REQUEST_COUNT.labels(method="GET", endpoint="/users", status="200") # If you need per-user metrics, use a different approach: # - Log the user_id and query logs # - Use a separate analytics system # - Bucket users by type/tier REQUEST_COUNT.labels( method="GET", endpoint="/users", user_tier="premium", # Bounded set of values ) ``` ### Pattern 7: Timed Operations with Context Manager Create a reusable timing context manager for operations. ```python from contextlib import contextmanager import time import structlog logger = structlog.get_logger() @contextmanager def timed_operation(name: str, **extra_fields): """Context manager for timing and logging operations.""" start = time.perf_counter() logger.debug("Operation started", operation=name, **extra_fields) try: yield except Exception as e: elapsed_ms = (time.perf_counter() - start) * 1000 logger.error( "Operation failed", operation=name, duration_ms=round(elapsed_ms, 2), error=str(e), **extra_fields, ) raise else: elapsed_ms = (time.perf_counter() - start) * 1000 logger.info( "Operation completed", operation=name, duration_ms=round(elapsed_ms, 2), **extra_fields, ) # Usage with timed_operation("fetch_user_orders", user_id=user.id): orders = await order_repository.get_by_user(user.id) ``` ### Pattern 8: OpenTelemetry Tracing Set up distributed tracing with OpenTelemetry. **Note:** OpenTelemetry is actively evolving. Check the [official Python documentation](https://opentelemetry.io/docs/languages/python/) for the latest API patterns and best practices. ```python from opentelemetry import trace from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import BatchSpanProcessor from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter def configure_tracing(service_name: str, otlp_endpoint: str) -> None: """Configure OpenTelemetry tracing.""" provider = TracerProvider() processor = BatchSpanProcessor(OTLPSpanExporter(endpoint=otlp_endpoint)) provider.add_span_processor(processor) trace.set_tracer_provider(provider) tracer = trace.get_tracer(__name__) async def process_order(order_id: str) -> Order: """Process order with tracing.""" with tracer.start_as_current_span("process_order") as span: span.set_attribute("order.id", order_id) with tracer.start_as_current_span("validate_order"): validate_order(order_id) with tracer.start_as_current_span("charge_payment"): charge_payment(order_id) with tracer.start_as_current_span("send_confirmation"): send_confirmation(order_id) return order ``` ## Best Practices Summary 1. **Use structured logging** - JSON logs with consistent fields 2. **Propagate correlation IDs** - Thread through all requests and logs 3. **Track the four golden signals** - Latency, traffic, errors, saturation 4. **Bound label cardinality** - Never use unbounded values as metric labels 5. **Log at appropriate levels** - Don't cry wolf with ERROR 6. **Include context** - User ID, request ID, operation name in logs 7. **Use context managers** - Consistent timing and error handling 8. **Separate concerns** - Observability code shouldn't pollute business logic 9. **Test your observability** - Verify logs and metrics in integration tests 10. **Set up alerts** - Metrics are useless without alerting
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered
πŸ€–system promptβ€’7 months ago

backtesting-frameworks

Build robust backtesting systems for trading strategies with proper

coding
⭐1
# Backtesting Frameworks Build robust, production-grade backtesting systems that avoid common pitfalls and produce reliable strategy performance estimates. ## When to Use This Skill - Developing trading strategy backtests - Building backtesting infrastructure - Validating strategy performance - Avoiding common backtesting biases - Implementing walk-forward analysis - Comparing strategy alternatives ## Core Concepts ### 1. Backtesting Biases | Bias | Description | Mitigation | | ---------------- | ------------------------- | ----------------------- | | **Look-ahead** | Using future information | Point-in-time data | | **Survivorship** | Only testing on survivors | Use delisted securities | | **Overfitting** | Curve-fitting to history | Out-of-sample testing | | **Selection** | Cherry-picking strategies | Pre-registration | | **Transaction** | Ignoring trading costs | Realistic cost models | ### 2. Proper Backtest Structure ``` Historical Data β”‚ β–Ό β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Training Set β”‚ β”‚ (Strategy Development & Optimization) β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β–Ό β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Validation Set β”‚ β”‚ (Parameter Selection, No Peeking) β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β–Ό β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚ Test Set β”‚ β”‚ (Final Performance Evaluation) β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ ``` ### 3. Walk-Forward Analysis ``` Window 1: [Train──────][Test] Window 2: [Train──────][Test] Window 3: [Train──────][Test] Window 4: [Train──────][Test] ─────▢ Time ``` ## Implementation Patterns ### Pattern 1: Event-Driven Backtester ```python from abc import ABC, abstractmethod from dataclasses import dataclass, field from datetime import datetime from decimal import Decimal from enum import Enum from typing import Dict, List, Optional import pandas as pd import numpy as np class OrderSide(Enum): BUY = "buy" SELL = "sell" class OrderType(Enum): MARKET = "market" LIMIT = "limit" STOP = "stop" @dataclass class Order: symbol: str side: OrderSide quantity: Decimal order_type: OrderType limit_price: Optional[Decimal] = None stop_price: Optional[Decimal] = None timestamp: Optional[datetime] = None @dataclass class Fill: order: Order fill_price: Decimal fill_quantity: Decimal commission: Decimal slippage: Decimal timestamp: datetime @dataclass class Position: symbol: str quantity: Decimal = Decimal("0") avg_cost: Decimal = Decimal("0") realized_pnl: Decimal = Decimal("0") def update(self, fill: Fill) -> None: if fill.order.side == OrderSide.BUY: new_quantity = self.quantity + fill.fill_quantity if new_quantity != 0: self.avg_cost = ( (self.quantity * self.avg_cost + fill.fill_quantity * fill.fill_price) / new_quantity ) self.quantity = new_quantity else: self.realized_pnl += fill.fill_quantity * (fill.fill_price - self.avg_cost) self.quantity -= fill.fill_quantity @dataclass class Portfolio: cash: Decimal positions: Dict[str, Position] = field(default_factory=dict) def get_position(self, symbol: str) -> Position: if symbol not in self.positions: self.positions[symbol] = Position(symbol=symbol) return self.positions[symbol] def process_fill(self, fill: Fill) -> None: position = self.get_position(fill.order.symbol) position.update(fill) if fill.order.side == OrderSide.BUY: self.cash -= fill.fill_price * fill.fill_quantity + fill.commission else: self.cash += fill.fill_price * fill.fill_quantity - fill.commission def get_equity(self, prices: Dict[str, Decimal]) -> Decimal: equity = self.cash for symbol, position in self.positions.items(): if position.quantity != 0 and symbol in prices: equity += position.quantity * prices[symbol] return equity class Strategy(ABC): @abstractmethod def on_bar(self, timestamp: datetime, data: pd.DataFrame) -> List[Order]: pass @abstractmethod def on_fill(self, fill: Fill) -> None: pass class ExecutionModel(ABC): @abstractmethod def execute(self, order: Order, bar: pd.Series) -> Optional[Fill]: pass class SimpleExecutionModel(ExecutionModel): def __init__(self, slippage_bps: float = 10, commission_per_share: float = 0.01): self.slippage_bps = slippage_bps self.commission_per_share = commission_per_share def execute(self, order: Order, bar: pd.Series) -> Optional[Fill]: if order.order_type == OrderType.MARKET: base_price = Decimal(str(bar["open"])) # Apply slippage slippage_mult = 1 + (self.slippage_bps / 10000) if order.side == OrderSide.BUY: fill_price = base_price * Decimal(str(slippage_mult)) else: fill_price = base_price / Decimal(str(slippage_mult)) commission = order.quantity * Decimal(str(self.commission_per_share)) slippage = abs(fill_price - base_price) * order.quantity return Fill( order=order, fill_price=fill_price, fill_quantity=order.quantity, commission=commission, slippage=slippage, timestamp=bar.name ) return None class Backtester: def __init__( self, strategy: Strategy, execution_model: ExecutionModel, initial_capital: Decimal = Decimal("100000") ): self.strategy = strategy self.execution_model = execution_model self.portfolio = Portfolio(cash=initial_capital) self.equity_curve: List[tuple] = [] self.trades: List[Fill] = [] def run(self, data: pd.DataFrame) -> pd.DataFrame: """Run backtest on OHLCV data with DatetimeIndex.""" pending_orders: List[Order] = [] for timestamp, bar in data.iterrows(): # Execute pending orders at today's prices for order in pending_orders: fill = self.execution_model.execute(order, bar) if fill: self.portfolio.process_fill(fill) self.strategy.on_fill(fill) self.trades.append(fill) pending_orders.clear() # Get current prices for equity calculation prices = {data.index.name or "default": Decimal(str(bar["close"]))} equity = self.portfolio.get_equity(prices) self.equity_curve.append((timestamp, float(equity))) # Generate new orders for next bar new_orders = self.strategy.on_bar(timestamp, data.loc[:timestamp]) pending_orders.extend(new_orders) return self._create_results() def _create_results(self) -> pd.DataFrame: equity_df = pd.DataFrame(self.equity_curve, columns=["timestamp", "equity"]) equity_df.set_index("timestamp", inplace=True) equity_df["returns"] = equity_df["equity"].pct_change() return equity_df ``` ### Pattern 2: Vectorized Backtester (Fast) ```python import pandas as pd import numpy as np from typing import Callable, Dict, Any class VectorizedBacktester: """Fast vectorized backtester for simple strategies.""" def __init__( self, initial_capital: float = 100000, commission: float = 0.001, # 0.1% slippage: float = 0.0005 # 0.05% ): self.initial_capital = initial_capital self.commission = commission self.slippage = slippage def run( self, prices: pd.DataFrame, signal_func: Callable[[pd.DataFrame], pd.Series] ) -> Dict[str, Any]: """ Run backtest with signal function. Args: prices: DataFrame with 'close' column signal_func: Function that returns position signals (-1, 0, 1) Returns: Dictionary with results """ # Generate signals (shifted to avoid look-ahead) signals = signal_func(prices).shift(1).fillna(0) # Calculate returns returns = prices["close"].pct_change() # Calculate strategy returns with costs position_changes = signals.diff().abs() trading_costs = position_changes * (self.commission + self.slippage) strategy_returns = signals * returns - trading_costs # Build equity curve equity = (1 + strategy_returns).cumprod() * self.initial_capital # Calculate metrics results = { "equity": equity, "returns": strategy_returns, "signals": signals, "metrics": self._calculate_metrics(strategy_returns, equity) } return results def _calculate_metrics( self, returns: pd.Series, equity: pd.Series ) -> Dict[str, float]: """Calculate performance metrics.""" total_return = (equity.iloc[-1] / self.initial_capital) - 1 annual_return = (1 + total_return) ** (252 / len(returns)) - 1 annual_vol = returns.std() * np.sqrt(252) sharpe = annual_return / annual_vol if annual_vol > 0 else 0 # Drawdown rolling_max = equity.cummax() drawdown = (equity - rolling_max) / rolling_max max_drawdown = drawdown.min() # Win rate winning_days = (returns > 0).sum() total_days = (returns != 0).sum() win_rate = winning_days / total_days if total_days > 0 else 0 return { "total_return": total_return, "annual_return": annual_return, "annual_volatility": annual_vol, "sharpe_ratio": sharpe, "max_drawdown": max_drawdown, "win_rate": win_rate, "num_trades": int((returns != 0).sum()) } # Example usage def momentum_signal(prices: pd.DataFrame, lookback: int = 20) -> pd.Series: """Simple momentum strategy: long when price > SMA, else flat.""" sma = prices["close"].rolling(lookback).mean() return (prices["close"] > sma).astype(int) # Run backtest # backtester = VectorizedBacktester() # results = backtester.run(price_data, lambda p: momentum_signal(p, 50)) ``` ### Pattern 3: Walk-Forward Optimization ```python from typing import Callable, Dict, List, Tuple, Any import pandas as pd import numpy as np from itertools import product class WalkForwardOptimizer: """Walk-forward analysis with anchored or rolling windows.""" def __init__( self, train_period: int, test_period: int, anchored: bool = False, n_splits: int = None ): """ Args: train_period: Number of bars in training window test_period: Number of bars in test window anchored: If True, training always starts from beginning n_splits: Number of train/test splits (auto-calculated if None) """ self.train_period = train_period self.test_period = test_period self.anchored = anchored self.n_splits = n_splits def generate_splits( self, data: pd.DataFrame ) -> List[Tuple[pd.DataFrame, pd.DataFrame]]: """Generate train/test splits.""" splits = [] n = len(data) if self.n_splits: step = (n - self.train_period) // self.n_splits else: step = self.test_period start = 0 while start + self.train_period + self.test_period <= n: if self.anchored: train_start = 0 else: train_start = start train_end = start + self.train_period test_end = min(train_end + self.test_period, n) train_data = data.iloc[train_start:train_end] test_data = data.iloc[train_end:test_end] splits.append((train_data, test_data)) start += step return splits def optimize( self, data: pd.DataFrame, strategy_func: Callable, param_grid: Dict[str, List], metric: str = "sharpe_ratio" ) -> Dict[str, Any]: """ Run walk-forward optimization. Args: data: Full dataset strategy_func: Function(data, **params) -> results dict param_grid: Parameter combinations to test metric: Metric to optimize Returns: Combined results from all test periods """ splits = self.generate_splits(data) all_results = [] optimal_params_history = [] for i, (train_data, test_data) in enumerate(splits): # Optimize on training data best_params, best_metric = self._grid_search( train_data, strategy_func, param_grid, metric ) optimal_params_history.append(best_params) # Test with optimal params test_results = strategy_func(test_data, **best_params) test_results["split"] = i test_results["params"] = best_params all_results.append(test_results) print(f"Split {i+1}/{len(splits)}: " f"Best {metric}={best_metric:.4f}, params={best_params}") return { "split_results": all_results, "param_history": optimal_params_history, "combined_equity": self._combine_equity_curves(all_results) } def _grid_search( self, data: pd.DataFrame, strategy_func: Callable, param_grid: Dict[str, List], metric: str ) -> Tuple[Dict, float]: """Grid search for best parameters.""" best_params = None best_metric = -np.inf # Generate all parameter combinations param_names = list(param_grid.keys()) param_values = list(param_grid.values()) for values in product(*param_values): params = dict(zip(param_names, values)) results = strategy_func(data, **params) if results["metrics"][metric] > best_metric: best_metric = results["metrics"][metric] best_params = params return best_params, best_metric def _combine_equity_curves( self, results: List[Dict] ) -> pd.Series: """Combine equity curves from all test periods.""" combined = pd.concat([r["equity"] for r in results]) return combined ``` ### Pattern 4: Monte Carlo Analysis ```python import numpy as np import pandas as pd from typing import Dict, List class MonteCarloAnalyzer: """Monte Carlo simulation for strategy robustness.""" def __init__(self, n_simulations: int = 1000, confidence: float = 0.95): self.n_simulations = n_simulations self.confidence = confidence def bootstrap_returns( self, returns: pd.Series, n_periods: int = None ) -> np.ndarray: """ Bootstrap simulation by resampling returns. Args: returns: Historical returns series n_periods: Length of each simulation (default: same as input) Returns: Array of shape (n_simulations, n_periods) """ if n_periods is None: n_periods = len(returns) simulations = np.zeros((self.n_simulations, n_periods)) for i in range(self.n_simulations): # Resample with replacement simulated_returns = np.random.choice( returns.values, size=n_periods, replace=True ) simulations[i] = simulated_returns return simulations def analyze_drawdowns( self, returns: pd.Series ) -> Dict[str, float]: """Analyze drawdown distribution via simulation.""" simulations = self.bootstrap_returns(returns) max_drawdowns = [] for sim_returns in simulations: equity = (1 + sim_returns).cumprod() rolling_max = np.maximum.accumulate(equity) drawdowns = (equity - rolling_max) / rolling_max max_drawdowns.append(drawdowns.min()) max_drawdowns = np.array(max_drawdowns) return { "expected_max_dd": np.mean(max_drawdowns), "median_max_dd": np.median(max_drawdowns), f"worst_{int(self.confidence*100)}pct": np.percentile( max_drawdowns, (1 - self.confidence) * 100 ), "worst_case": max_drawdowns.min() } def probability_of_loss( self, returns: pd.Series, holding_periods: List[int] = [21, 63, 126, 252] ) -> Dict[int, float]: """Calculate probability of loss over various holding periods.""" results = {} for period in holding_periods: if period > len(returns): continue simulations = self.bootstrap_returns(returns, period) total_returns = (1 + simulations).prod(axis=1) - 1 prob_loss = (total_returns < 0).mean() results[period] = prob_loss return results def confidence_interval( self, returns: pd.Series, periods: int = 252 ) -> Dict[str, float]: """Calculate confidence interval for future returns.""" simulations = self.bootstrap_returns(returns, periods) total_returns = (1 + simulations).prod(axis=1) - 1 lower = (1 - self.confidence) / 2 upper = 1 - lower return { "expected": total_returns.mean(), "lower_bound": np.percentile(total_returns, lower * 100), "upper_bound": np.percentile(total_returns, upper * 100), "std": total_returns.std() } ``` ## Performance Metrics ```python def calculate_metrics(returns: pd.Series, rf_rate: float = 0.02) -> Dict[str, float]: """Calculate comprehensive performance metrics.""" # Annualization factor (assuming daily returns) ann_factor = 252 # Basic metrics total_return = (1 + returns).prod() - 1 annual_return = (1 + total_return) ** (ann_factor / len(returns)) - 1 annual_vol = returns.std() * np.sqrt(ann_factor) # Risk-adjusted returns sharpe = (annual_return - rf_rate) / annual_vol if annual_vol > 0 else 0 # Sortino (downside deviation) downside_returns = returns[returns < 0] downside_vol = downside_returns.std() * np.sqrt(ann_factor) sortino = (annual_return - rf_rate) / downside_vol if downside_vol > 0 else 0 # Calmar ratio equity = (1 + returns).cumprod() rolling_max = equity.cummax() drawdowns = (equity - rolling_max) / rolling_max max_drawdown = drawdowns.min() calmar = annual_return / abs(max_drawdown) if max_drawdown != 0 else 0 # Win rate and profit factor wins = returns[returns > 0] losses = return
πŸ‘0
πŸ‘οΈ0
πŸ€– Auto-discovered