Skip to main content
EVOKORE// BROWSE
>

./browse/prompts

27 NODES
๐Ÿ“textโ€ข3 hours ago

Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering

Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering User centric prompting for analyzing inspiration pages or instructing on internal analysis cycles

architecture
โญ1
# Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering You are an **Expert Enterprise Architect, Product Director, and Lead Engineer**. Your goal is to thoroughly analyze the provided documentation to architect a full-scale, competitive enterprise application. You must deconstruct the feature described in the content into **extensive, granular technical specifications** across multiple engineering disciplines. **CRITICAL INSTRUCTIONS**: - **DO NOT SUMMARIZE**. Be exhaustive. - For every category below, aim to list **10+ specific items** if possible. - Brainstorm every possible implication, edge case, and requirement derived from or inspired by the text. - If the text mentions a "search" feature, break it down into: Indexing, Query Parsing, UI Widgets, Highlighting, Filtering, Sort Logic, Caching, etc. Please output the analysis in the following Markdown format: # 1. Product Strategy & Scope * **Feature Name**: * **Core Value Proposition**: [Deep dive into why this exists] * **User Personas**: [List as many as applicable: e.g. Admin, Power User, Viewer, Auditor, API Consumer...] * **User Stories**: [Extensive list of 10+ granular user stories e.g. "As a User, I want to..."] * **Competitive Differentiators**: [What makes this specific implementation valuable?] # 2. Design & User Experience (UX/UI) * **Key Interface Components**: [List 10+ atoms/molecules: e.g. Data Grid, Filter Chips, Modals, Tooltips, Empty States, Toasts, Dropdowns...] * **Interaction Patterns**: [List 10+ patterns: e.g. Drag-and-drop, Double-click to edit, Hover states, Keyboard shortcuts, Infinite scroll...] * **Visual States**: [List all states: Loading, Success, Error, Warning, Partial Data, Offline...] * **Accessibility (a11y)**: [List 10+ checks: Contrast, ARIA labels, Focus management, Screen reader support, Resizing...] # 3. Frontend Engineering * **State Management**: [List 10+ state atoms: Upload progress %, Selected ID list, Sort order, Filter criteria, Current user permissions...] * **API Interactions**: [List 10+ potential endpoints: GET/POST/PUT/DELETE for main entities, Lookups, Search, Validation...] * **Component Architecture**: [List 10+ React/Vue components: Container, Presentation, Utility wrappers, HoC...] * **Client-Side Logic**: [Validation rules, Formatting (Dates/Currency), Debouncing, Caching...] # 4. Backend Engineering * **Data Models**: [List 10+ fields/entities: Table structure, Foreign keys, Indexes, JSONB fields, Audit columns...] * **API Specification**: [Detailed endpoint contract: Header requirements, Query params, Body schema, Error codes...] * **Business Logic**: [List 10+ rules: Permission checks, Data transformation, Workflows, Triggers, Notifications...] * **Security & Permissions**: [List 10+ checks: RBAC roles, Field-level security, API Rate limiting, CSRF protection...] # 5. Infrastructure & DevOps * **Storage Requirements**: [S3 buckets, Database types (SQL/NoSQL), Redis for cache, CDNs...] * **Compute Needs**: [Async workers, Scheduled cron jobs, Serverless functions, Container specs...] * **Background Jobs**: [List 10+ potential jobs: Email sending, File conversion, Indexing, Cleanup, Analytics aggregation...] * **Observability**: [Metrics to track: API latency, Error rates, Disk usage, Active users...] # 6. Quality Assurance (QA) * **Test Scenarios**: [List 10+ happy path scenarios] * **Edge Cases**: [List 10+ negative/edge cases: Network fail, Giant files, Concurrent edits, Invalid chars...] * **Performance Metrics**: [Specific SLAs: <200ms API, <1s Page load, 99.9% Uptime...] * **Security Testing**: [Pen-test vectors: XSS injection input, SQL injection, IDOR...] # 7. Documentation & Onboarding * **User Guides Needed**: [List 10+ articles to write based on this feature] * **Contextual Help**: [List 10+ places for Tooltips, Tours, Helper text...] * **API Documentation**: [Swagger/OpenAPI requirements] # 8. Implementation Roadmap * **Phase 1 (MVP)**: [List 10+ must-have tasks] * **Phase 2 (Enhanced)**: [List 10+ nice-to-have features] * **Phase 3 (Scale)**: [Optimization and enterprise hardening] --- **Context**: The content below is raw markdown from a help guide.
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿ“textโ€ข3 hours ago

Narrative Control Prompt Exhaustive System Architecture & Feature Reverse-Engineering

Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering User concentric prompting for analyzing inspiration pages or instructing on internal analysis cycles

architecture
โญ1
# Narrative Control Prompt: Exhaustive System Architecture & Feature Reverse-Engineering You are an **Expert Enterprise Architect, Product Director, and Lead Engineer**. Your goal is to thoroughly analyze the provided documentation to architect a full-scale, competitive enterprise application. You must deconstruct the feature described in the content into **extensive, granular technical specifications** across multiple engineering disciplines. **CRITICAL INSTRUCTIONS**: - **DO NOT SUMMARIZE**. Be exhaustive. - For every category below, aim to list **10+ specific items** if possible. - Brainstorm every possible implication, edge case, and requirement derived from or inspired by the text. - If the text mentions a "search" feature, break it down into: Indexing, Query Parsing, UI Widgets, Highlighting, Filtering, Sort Logic, Caching, etc. Please output the analysis in the following Markdown format: # 1. Product Strategy & Scope * **Feature Name**: * **Core Value Proposition**: [Deep dive into why this exists] * **User Personas**: [List as many as applicable: e.g. Admin, Power User, Viewer, Auditor, API Consumer...] * **User Stories**: [Extensive list of 10+ granular user stories e.g. "As a User, I want to..."] * **Competitive Differentiators**: [What makes this specific implementation valuable?] # 2. Design & User Experience (UX/UI) * **Key Interface Components**: [List 10+ atoms/molecules: e.g. Data Grid, Filter Chips, Modals, Tooltips, Empty States, Toasts, Dropdowns...] * **Interaction Patterns**: [List 10+ patterns: e.g. Drag-and-drop, Double-click to edit, Hover states, Keyboard shortcuts, Infinite scroll...] * **Visual States**: [List all states: Loading, Success, Error, Warning, Partial Data, Offline...] * **Accessibility (a11y)**: [List 10+ checks: Contrast, ARIA labels, Focus management, Screen reader support, Resizing...] # 3. Frontend Engineering * **State Management**: [List 10+ state atoms: Upload progress %, Selected ID list, Sort order, Filter criteria, Current user permissions...] * **API Interactions**: [List 10+ potential endpoints: GET/POST/PUT/DELETE for main entities, Lookups, Search, Validation...] * **Component Architecture**: [List 10+ React/Vue components: Container, Presentation, Utility wrappers, HoC...] * **Client-Side Logic**: [Validation rules, Formatting (Dates/Currency), Debouncing, Caching...] # 4. Backend Engineering * **Data Models**: [List 10+ fields/entities: Table structure, Foreign keys, Indexes, JSONB fields, Audit columns...] * **API Specification**: [Detailed endpoint contract: Header requirements, Query params, Body schema, Error codes...] * **Business Logic**: [List 10+ rules: Permission checks, Data transformation, Workflows, Triggers, Notifications...] * **Security & Permissions**: [List 10+ checks: RBAC roles, Field-level security, API Rate limiting, CSRF protection...] # 5. Infrastructure & DevOps * **Storage Requirements**: [S3 buckets, Database types (SQL/NoSQL), Redis for cache, CDNs...] * **Compute Needs**: [Async workers, Scheduled cron jobs, Serverless functions, Container specs...] * **Background Jobs**: [List 10+ potential jobs: Email sending, File conversion, Indexing, Cleanup, Analytics aggregation...] * **Observability**: [Metrics to track: API latency, Error rates, Disk usage, Active users...] # 6. Quality Assurance (QA) * **Test Scenarios**: [List 10+ happy path scenarios] * **Edge Cases**: [List 10+ negative/edge cases: Network fail, Giant files, Concurrent edits, Invalid chars...] * **Performance Metrics**: [Specific SLAs: <200ms API, <1s Page load, 99.9% Uptime...] * **Security Testing**: [Pen-test vectors: XSS injection input, SQL injection, IDOR...] # 7. Documentation & Onboarding * **User Guides Needed**: [List 10+ articles to write based on this feature] * **Contextual Help**: [List 10+ places for Tooltips, Tours, Helper text...] * **API Documentation**: [Swagger/OpenAPI requirements] # 8. Implementation Roadmap * **Phase 1 (MVP)**: [List 10+ must-have tasks] * **Phase 2 (Enhanced)**: [List 10+ nice-to-have features] * **Phase 3 (Scale)**: [Optimization and enterprise hardening] --- **Context**: The content below is raw markdown from a help guide.
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿ“textโ€ข6 months ago

Evaluation and Regression Gate

Create evaluation runs, submit results, compute metrics, save baselines, and triage regressions before shipping changes.

analysis
โญ1
# Evaluation and Regression Gate Imported from curated first-party documentation sources. ## What this covers Use this workflow to formalize benchmarking, regression detection, and baseline management around agent or product changes. ## Use this when - Validating changes against prior baselines - Making regressions visible before rollout - Capturing repeatable evaluation evidence ## Expected outcomes - Evaluation runs produce comparable metrics - Baselines are stored after successful validation - Regression triage becomes a workflow instead of an afterthought ## Source synthesis - AGENT33/docs/functionality-and-workflows.md (https://github.com/mattmre/AGENT33/blob/main/docs/functionality-and-workflows.md) - AGENT33/docs/use-cases.md (https://github.com/mattmre/AGENT33/blob/main/docs/use-cases.md) ## Dedupe notes Combines AGENT33 evaluation lifecycle definitions with evaluation and regression use-case framing. ## Source excerpts ### AGENT33/docs/functionality-and-workflows.md ### 4.3 Evaluation Lifecycle Flow: 1. Create run (`/v1/evaluations/runs`) 2. Submit task results (`/runs/{id}/results`) 3. Compute metrics + gate report 4. Save baseline (`/runs/{id}/baseline`) 5. Triage/resolve regressions ### AGENT33/docs/use-cases.md ## 4. Evaluation and Regression Gates Goal: - Quantify quality and block regressions across PR/merge/release gates. Use these modules: - `api/routes/evaluations.py` - `evaluation/service.py` - `evaluation/gates.py` - `evaluation/regression.py` Typical flow: 1. Create run for gate type (`G-PR`, `G-MRG`, `G-REL`, `G-MON`). 2. Submit task results and quality metadata. 3. Compute metrics and gate verdict. 4. Save baseline for future comparison. 5. Triage/resolve regression records. Best fit: - Teams with golden-task style quality gates.
๐Ÿ‘0
๐Ÿ‘๏ธ0
docs
๐Ÿค–system promptโ€ข7 months ago

context-driven-development

Creates and maintains project context artifacts (product.md,

coding
โญ1
# Context-Driven Development Guide for implementing and maintaining context as a managed artifact alongside code, enabling consistent AI interactions and team alignment through structured project documentation. ## When to Use This Skill - Setting up new projects with Conductor - Understanding the relationship between context artifacts - Maintaining consistency across AI-assisted development sessions - Onboarding team members to an existing Conductor project - Deciding when to update context documents - Managing greenfield vs brownfield project contexts ## Core Philosophy Context-Driven Development treats project context as a first-class artifact managed alongside code. Instead of relying on ad-hoc prompts or scattered documentation, establish a persistent, structured foundation that informs all AI interactions. Key principles: 1. **Context precedes code**: Define what you're building and how before implementation 2. **Living documentation**: Context artifacts evolve with the project 3. **Single source of truth**: One canonical location for each type of information 4. **AI alignment**: Consistent context produces consistent AI behavior ## The Workflow Follow the **Context โ†’ Spec & Plan โ†’ Implement** workflow: 1. **Context Phase**: Establish or verify project context artifacts exist and are current 2. **Specification Phase**: Define requirements and acceptance criteria for work units 3. **Planning Phase**: Break specifications into phased, actionable tasks 4. **Implementation Phase**: Execute tasks following established workflow patterns ## Artifact Relationships ### product.md - Defines WHAT and WHY Purpose: Captures product vision, goals, target users, and business context. Contents: - Product name and one-line description - Problem statement and solution approach - Target user personas - Core features and capabilities - Success metrics and KPIs - Product roadmap (high-level) Update when: - Product vision or goals change - New major features are planned - Target audience shifts - Business priorities evolve ### product-guidelines.md - Defines HOW to Communicate Purpose: Establishes brand voice, messaging standards, and communication patterns. Contents: - Brand voice and tone guidelines - Terminology and glossary - Error message conventions - User-facing copy standards - Documentation style Update when: - Brand guidelines change - New terminology is introduced - Communication patterns need refinement ### tech-stack.md - Defines WITH WHAT Purpose: Documents technology choices, dependencies, and architectural decisions. Contents: - Primary languages and frameworks - Key dependencies with versions - Infrastructure and deployment targets - Development tools and environment - Testing frameworks - Code quality tools Update when: - Adding new dependencies - Upgrading major versions - Changing infrastructure - Adopting new tools or patterns ### workflow.md - Defines HOW to Work Purpose: Establishes development practices, quality gates, and team workflows. Contents: - Development methodology (TDD, etc.) - Git workflow and commit conventions - Code review requirements - Testing requirements and coverage targets - Quality assurance gates - Deployment procedures Update when: - Team practices evolve - Quality standards change - New workflow patterns are adopted ### tracks.md - Tracks WHAT'S HAPPENING Purpose: Registry of all work units with status and metadata. Contents: - Active tracks with current status - Completed tracks with completion dates - Track metadata (type, priority, assignee) - Links to individual track directories Update when: - New tracks are created - Track status changes - Tracks are completed or archived See [references/artifact-templates.md](references/artifact-templates.md) for copy-paste starter templates. ## Context Maintenance Principles ### Keep Artifacts Synchronized Ensure changes in one artifact reflect in related documents: - New feature in product.md โ†’ Update tech-stack.md if new dependencies needed - Completed track โ†’ Update product.md to reflect new capabilities - Workflow change โ†’ Update all affected track plans ### Update tech-stack.md When Adding Dependencies Before adding any new dependency: 1. Check if existing dependencies solve the need 2. Document the rationale for new dependencies 3. Add version constraints 4. Note any configuration requirements ### Update product.md When Features Complete After completing a feature track: 1. Move feature from "planned" to "implemented" in product.md 2. Update any affected success metrics 3. Document any scope changes from original plan ### Verify Context Before Implementation Before starting any track: 1. Read all context artifacts 2. Flag any outdated information 3. Propose updates before proceeding 4. Confirm context accuracy with stakeholders ## Greenfield vs Brownfield Handling ### Greenfield Projects (New) For new projects: 1. Run `/conductor:setup` to create all artifacts interactively 2. Answer questions about product vision, tech preferences, and workflow 3. Generate initial style guides for chosen languages 4. Create empty tracks registry Characteristics: - Full control over context structure - Define standards before code exists - Establish patterns early ### Brownfield Projects (Existing) For existing codebases: 1. Run `/conductor:setup` with existing codebase detection 2. System analyzes existing code, configs, and documentation 3. Pre-populate artifacts based on discovered patterns 4. Review and refine generated context Characteristics: - Extract implicit context from existing code - Reconcile existing patterns with desired patterns - Document technical debt and modernization plans - Preserve working patterns while establishing standards ## Benefits ### Team Alignment - New team members onboard faster with explicit context - Consistent terminology and conventions across the team - Shared understanding of product goals and technical decisions ### AI Consistency - AI assistants produce aligned outputs across sessions - Reduced need to re-explain context in each interaction - Predictable behavior based on documented standards ### Institutional Memory - Decisions and rationale are preserved - Context survives team changes - Historical context informs future decisions ### Quality Assurance - Standards are explicit and verifiable - Deviations from context are detectable - Quality gates are documented and enforceable ## Directory Structure ``` conductor/ โ”œโ”€โ”€ index.md # Navigation hub linking all artifacts โ”œโ”€โ”€ product.md # Product vision and goals โ”œโ”€โ”€ product-guidelines.md # Communication standards โ”œโ”€โ”€ tech-stack.md # Technology preferences โ”œโ”€โ”€ workflow.md # Development practices โ”œโ”€โ”€ tracks.md # Work unit registry โ”œโ”€โ”€ setup_state.json # Resumable setup state โ”œโ”€โ”€ code_styleguides/ # Language-specific conventions โ”‚ โ”œโ”€โ”€ python.md โ”‚ โ”œโ”€โ”€ typescript.md โ”‚ โ””โ”€โ”€ ... โ””โ”€โ”€ tracks/ โ””โ”€โ”€ <track-id>/ โ”œโ”€โ”€ spec.md โ”œโ”€โ”€ plan.md โ”œโ”€โ”€ metadata.json โ””โ”€โ”€ index.md ``` ## Context Lifecycle 1. **Creation**: Initial setup via `/conductor:setup` 2. **Validation**: Verify before each track 3. **Evolution**: Update as project grows 4. **Synchronization**: Keep artifacts aligned 5. **Archival**: Document historical decisions ## Context Validation Checklist Before starting implementation on any track, validate context: ### Product Context - [ ] product.md reflects current product vision - [ ] Target users are accurately described - [ ] Feature list is up to date - [ ] Success metrics are defined ### Technical Context - [ ] tech-stack.md lists all current dependencies - [ ] Version numbers are accurate - [ ] Infrastructure targets are correct - [ ] Development tools are documented ### Workflow Context - [ ] workflow.md describes current practices - [ ] Quality gates are defined - [ ] Coverage targets are specified - [ ] Commit conventions are documented ### Track Context - [ ] tracks.md shows all active work - [ ] No stale or abandoned tracks - [ ] Dependencies between tracks are noted ## Common Anti-Patterns Avoid these context management mistakes: ### Stale Context Problem: Context documents become outdated and misleading. Solution: Update context as part of each track's completion process. ### Context Sprawl Problem: Information scattered across multiple locations. Solution: Use the defined artifact structure; resist creating new document types. ### Implicit Context Problem: Relying on knowledge not captured in artifacts. Solution: If you reference something repeatedly, add it to the appropriate artifact. ### Context Hoarding Problem: One person maintains context without team input. Solution: Review context artifacts in pull requests; make updates collaborative. ### Over-Specification Problem: Context becomes so detailed it's impossible to maintain. Solution: Keep artifacts focused on decisions that affect AI behavior and team alignment. ## Integration with Development Tools ### IDE Integration Configure your IDE to display context files prominently: - Pin conductor/product.md for quick reference - Add tech-stack.md to project notes - Create snippets for common patterns from style guides ### Git Hooks Consider pre-commit hooks that: - Warn when dependencies change without tech-stack.md update - Remind to update product.md when feature branches merge - Validate context artifact syntax ### CI/CD Integration Include context validation in pipelines: - Check tech-stack.md matches actual dependencies - Verify links in context documents resolve - Ensure tracks.md status matches git branch state ## Session Continuity Conductor supports multi-session development through context persistence: ### Starting a New Session 1. Read index.md to orient yourself 2. Check tracks.md for active work 3. Review relevant track's plan.md for current task 4. Verify context artifacts are current ### Ending a Session 1. Update plan.md with current progress 2. Note any blockers or decisions made 3. Commit in-progress work with clear status 4. Update tracks.md if status changed ### Handling Interruptions If interrupted mid-task: 1. Mark task as `[~]` with note about stopping point 2. Commit work-in-progress to feature branch 3. Document any uncommitted decisions in plan.md ## Best Practices 1. **Read context first**: Always read relevant artifacts before starting work 2. **Small updates**: Make incremental context changes, not massive rewrites 3. **Link decisions**: Reference context when making implementation choices 4. **Version context**: Commit context changes alongside code changes 5. **Review context**: Include context artifact reviews in code reviews 6. **Validate regularly**: Run context validation checklist before major work 7. **Communicate changes**: Notify team when context artifacts change significantly 8. **Preserve history**: Use git to track context evolution over time 9. **Question staleness**: If context feels wrong, investigate and update 10. **Keep it actionable**: Every context item should inform a decision or behavior
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

track-management

Use this skill when creating, managing, or working with Conductor

coding
โญ1
# Track Management Guide for creating, managing, and completing Conductor tracks - the logical work units that organize features, bugs, and refactors through specification, planning, and implementation phases. ## When to Use This Skill - Creating new feature, bug, or refactor tracks - Writing or reviewing spec.md files - Creating or updating plan.md files - Managing track lifecycle from creation to completion - Understanding track status markers and conventions - Working with the tracks.md registry - Interpreting or updating track metadata ## Track Concept A track is a logical work unit that encapsulates a complete piece of work. Each track has: - A unique identifier - A specification defining requirements - A phased plan breaking work into tasks - Metadata tracking status and progress Tracks provide semantic organization for work, enabling: - Clear scope boundaries - Progress tracking - Git-aware operations (revert by track) - Team coordination ## Track Types ### feature New functionality or capabilities. Use for: - New user-facing features - New API endpoints - New integrations - Significant enhancements ### bug Defect fixes. Use for: - Incorrect behavior - Error conditions - Performance regressions - Security vulnerabilities ### chore Maintenance and housekeeping. Use for: - Dependency updates - Configuration changes - Documentation updates - Cleanup tasks ### refactor Code improvement without behavior change. Use for: - Code restructuring - Pattern adoption - Technical debt reduction - Performance optimization (same behavior, better performance) ## Track ID Format Track IDs follow the pattern: `{shortname}_{YYYYMMDD}` - **shortname**: 2-4 word kebab-case description (e.g., `user-auth`, `api-rate-limit`) - **YYYYMMDD**: Creation date in ISO format Examples: - `user-auth_20250115` - `fix-login-error_20250115` - `upgrade-deps_20250115` - `refactor-api-client_20250115` ## Track Lifecycle ### 1. Creation (newTrack) **Define Requirements** 1. Gather requirements through interactive Q&A 2. Identify acceptance criteria 3. Determine scope boundaries 4. Identify dependencies **Generate Specification** 1. Create `spec.md` with structured requirements 2. Document functional and non-functional requirements 3. Define acceptance criteria 4. List dependencies and constraints **Generate Plan** 1. Create `plan.md` with phased task breakdown 2. Organize tasks into logical phases 3. Add verification tasks after phases 4. Estimate effort and complexity **Register Track** 1. Add entry to `tracks.md` registry 2. Create track directory structure 3. Generate `metadata.json` 4. Create track `index.md` ### 2. Implementation **Execute Tasks** 1. Select next pending task from plan 2. Mark task as in-progress 3. Implement following workflow (TDD) 4. Mark task complete with commit SHA **Update Status** 1. Update task markers in plan.md 2. Record commit SHAs for traceability 3. Update phase progress 4. Update track status in tracks.md **Verify Progress** 1. Complete verification tasks 2. Wait for checkpoint approval 3. Record checkpoint commits ### 3. Completion **Sync Documentation** 1. Update product.md if features added 2. Update tech-stack.md if dependencies changed 3. Verify all acceptance criteria met **Archive or Delete** 1. Mark track as completed in tracks.md 2. Record completion date 3. Archive or retain track directory ## Specification (spec.md) Structure ```markdown # {Track Title} ## Overview Brief description of what this track accomplishes and why. ## Functional Requirements ### FR-1: {Requirement Name} Description of the functional requirement. - Acceptance: How to verify this requirement is met ### FR-2: {Requirement Name} ... ## Non-Functional Requirements ### NFR-1: {Requirement Name} Description of the non-functional requirement (performance, security, etc.) - Target: Specific measurable target - Verification: How to test ## Acceptance Criteria - [ ] Criterion 1: Specific, testable condition - [ ] Criterion 2: Specific, testable condition - [ ] Criterion 3: Specific, testable condition ## Scope ### In Scope - Explicitly included items - Features to implement - Components to modify ### Out of Scope - Explicitly excluded items - Future considerations - Related but separate work ## Dependencies ### Internal - Other tracks or components this depends on - Required context artifacts ### External - Third-party services or APIs - External dependencies ## Risks and Mitigations | Risk | Impact | Mitigation | | ---------------- | --------------- | ------------------- | | Risk description | High/Medium/Low | Mitigation strategy | ## Open Questions - [ ] Question that needs resolution - [x] Resolved question - Answer ``` ## Plan (plan.md) Structure ```markdown # Implementation Plan: {Track Title} Track ID: `{track-id}` Created: YYYY-MM-DD Status: pending | in-progress | completed ## Overview Brief description of implementation approach. ## Phase 1: {Phase Name} ### Tasks - [ ] **Task 1.1**: Task description - Sub-task or detail - Sub-task or detail - [ ] **Task 1.2**: Task description - [ ] **Task 1.3**: Task description ### Verification - [ ] **Verify 1.1**: Verification step for phase ## Phase 2: {Phase Name} ### Tasks - [ ] **Task 2.1**: Task description - [ ] **Task 2.2**: Task description ### Verification - [ ] **Verify 2.1**: Verification step for phase ## Phase 3: Finalization ### Tasks - [ ] **Task 3.1**: Update documentation - [ ] **Task 3.2**: Final integration test ### Verification - [ ] **Verify 3.1**: All acceptance criteria met ## Checkpoints | Phase | Checkpoint SHA | Date | Status | | ------- | -------------- | ---- | ------- | | Phase 1 | | | pending | | Phase 2 | | | pending | | Phase 3 | | | pending | ``` ## Status Marker Conventions Use consistent markers in plan.md: | Marker | Meaning | Usage | | ------ | ----------- | --------------------------- | | `[ ]` | Pending | Task not started | | `[~]` | In Progress | Currently being worked | | `[x]` | Complete | Task finished (include SHA) | | `[-]` | Skipped | Intentionally not done | | `[!]` | Blocked | Waiting on dependency | Example: ```markdown - [x] **Task 1.1**: Set up database schema `abc1234` - [~] **Task 1.2**: Implement user model - [ ] **Task 1.3**: Add validation logic - [!] **Task 1.4**: Integrate auth service (blocked: waiting for API key) - [-] **Task 1.5**: Legacy migration (skipped: not needed) ``` ## Track Registry (tracks.md) Format ```markdown # Track Registry ## Active Tracks | Track ID | Type | Status | Phase | Started | Assignee | | ------------------------------------------------ | ------- | ----------- | ----- | ---------- | ---------- | | [user-auth_20250115](tracks/user-auth_20250115/) | feature | in-progress | 2/3 | 2025-01-15 | @developer | | [fix-login_20250114](tracks/fix-login_20250114/) | bug | pending | 0/2 | 2025-01-14 | - | ## Completed Tracks | Track ID | Type | Completed | Duration | | ---------------------------------------------- | ----- | ---------- | -------- | | [setup-ci_20250110](tracks/setup-ci_20250110/) | chore | 2025-01-12 | 2 days | ## Archived Tracks | Track ID | Reason | Archived | | ---------------------------------------------------- | ---------- | ---------- | | [old-feature_20241201](tracks/old-feature_20241201/) | Superseded | 2025-01-05 | ``` ## Metadata (metadata.json) Fields ```json { "id": "user-auth_20250115", "title": "User Authentication System", "type": "feature", "status": "in-progress", "priority": "high", "created": "2025-01-15T10:30:00Z", "updated": "2025-01-15T14:45:00Z", "started": "2025-01-15T11:00:00Z", "completed": null, "assignee": "@developer", "phases": { "total": 3, "current": 2, "completed": 1 }, "tasks": { "total": 12, "completed": 5, "in_progress": 1, "pending": 6 }, "checkpoints": [ { "phase": 1, "sha": "abc1234", "date": "2025-01-15T13:00:00Z" } ], "dependencies": [], "tags": ["auth", "security"] } ``` ## Track Operations ### Creating a Track 1. Run `/conductor:new-track` 2. Answer interactive questions 3. Review generated spec.md 4. Review generated plan.md 5. Confirm track creation ### Starting Implementation 1. Read spec.md and plan.md 2. Verify context artifacts are current 3. Mark first task as `[~]` 4. Begin TDD workflow ### Completing a Phase 1. Ensure all phase tasks are `[x]` 2. Complete verification tasks 3. Wait for checkpoint approval 4. Record checkpoint SHA 5. Proceed to next phase ### Completing a Track 1. Verify all phases complete 2. Verify all acceptance criteria met 3. Update product.md if needed 4. Mark track completed in tracks.md 5. Update metadata.json ### Reverting a Track 1. Run `/conductor:revert` 2. Select track to revert 3. Choose granularity (track/phase/task) 4. Confirm revert operation 5. Update status markers ## Handling Track Dependencies ### Identifying Dependencies During track creation, identify: - **Hard dependencies**: Must complete before this track can start - **Soft dependencies**: Can proceed in parallel but may affect integration - **External dependencies**: Third-party services, APIs, or team decisions ### Documenting Dependencies In spec.md, list dependencies with: - Dependency type (hard/soft/external) - Current status (available/pending/blocked) - Resolution path (what needs to happen) ### Managing Blocked Tracks When a track is blocked: 1. Mark blocked tasks with `[!]` and reason 2. Update tracks.md status 3. Document blocker in metadata.json 4. Consider creating dependency track if needed ## Track Sizing Guidelines ### Right-Sized Tracks Aim for tracks that: - Complete in 1-5 days of work - Have 2-4 phases - Contain 8-20 tasks total - Deliver a coherent, testable unit ### Too Large Signs a track is too large: - More than 5 phases - More than 25 tasks - Multiple unrelated features - Estimated duration > 1 week Solution: Split into multiple tracks with clear boundaries. ### Too Small Signs a track is too small: - Single phase with 1-2 tasks - No meaningful verification needed - Could be a sub-task of another track - Less than a few hours of work Solution: Combine with related work or handle as part of existing track. ## Specification Quality Checklist Before finalizing spec.md, verify: ### Requirements Quality - [ ] Each requirement has clear acceptance criteria - [ ] Requirements are testable - [ ] Requirements are independent (can verify separately) - [ ] No ambiguous language ("should be fast" โ†’ "response < 200ms") ### Scope Clarity - [ ] In-scope items are specific - [ ] Out-of-scope items prevent scope creep - [ ] Boundaries are clear to implementer ### Dependencies Identified - [ ] All internal dependencies listed - [ ] External dependencies have owners/contacts - [ ] Dependency status is current ### Risks Addressed - [ ] Major risks identified - [ ] Impact assessment realistic - [ ] Mitigations are actionable ## Plan Quality Checklist Before starting implementation, verify plan.md: ### Task Quality - [ ] Tasks are atomic (one logical action) - [ ] Tasks are independently verifiable - [ ] Task descriptions are clear - [ ] Sub-tasks provide helpful detail ### Phase Organization - [ ] Phases group related tasks - [ ] Each phase delivers something testable - [ ] Verification tasks after each phase - [ ] Phases build on each other logically ### Completeness - [ ] All spec requirements have corresponding tasks - [ ] Documentation tasks included - [ ] Testing tasks included - [ ] Integration tasks included ## Common Track Patterns ### Feature Track Pattern ``` Phase 1: Foundation - Data models - Database migrations - Basic API structure Phase 2: Core Logic - Business logic implementation - Input validation - Error handling Phase 3: Integration - UI integration - API documentation - End-to-end tests ``` ### Bug Fix Track Pattern ``` Phase 1: Reproduction - Write failing test capturing bug - Document reproduction steps Phase 2: Fix - Implement fix - Verify test passes - Check for regressions Phase 3: Verification - Manual verification - Update documentation if needed ``` ### Refactor Track Pattern ``` Phase 1: Preparation - Add characterization tests - Document current behavior Phase 2: Refactoring - Apply changes incrementally - Maintain green tests throughout Phase 3: Cleanup - Remove dead code - Update documentation ``` ## Best Practices 1. **One track, one concern**: Keep tracks focused on a single logical change 2. **Small phases**: Break work into phases of 3-5 tasks maximum 3. **Verification after phases**: Always include verification tasks 4. **Update markers immediately**: Mark task status as you work 5. **Record SHAs**: Always note commit SHAs for completed tasks 6. **Review specs before planning**: Ensure spec is complete before creating plan 7. **Link dependencies**: Explicitly note track dependencies 8. **Archive, don't delete**: Preserve completed tracks for reference 9. **Size appropriately**: Keep tracks between 1-5 days of work 10. **Clear acceptance criteria**: Every requirement must be testable
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

workflow-patterns

Use this skill when implementing tasks according to Conductor's TDD

coding
โญ1
# Workflow Patterns Guide for implementing tasks using Conductor's TDD workflow, managing phase checkpoints, handling git commits, and executing the verification protocol that ensures quality throughout implementation. ## When to Use This Skill - Implementing tasks from a track's plan.md - Following TDD red-green-refactor cycle - Completing phase checkpoints - Managing git commits and notes - Understanding quality assurance gates - Handling verification protocols - Recording progress in plan files ## TDD Task Lifecycle Follow these 11 steps for each task: ### Step 1: Select Next Task Read plan.md and identify the next pending `[ ]` task. Select tasks in order within the current phase. Do not skip ahead to later phases. ### Step 2: Mark as In Progress Update plan.md to mark the task as `[~]`: ```markdown - [~] **Task 2.1**: Implement user validation ``` Commit this status change separately from implementation. ### Step 3: RED - Write Failing Tests Write tests that define the expected behavior before writing implementation: - Create test file if needed - Write test cases covering happy path - Write test cases covering edge cases - Write test cases covering error conditions - Run tests - they should FAIL Example: ```python def test_validate_user_email_valid(): user = User(email="test@example.com") assert user.validate_email() is True def test_validate_user_email_invalid(): user = User(email="invalid") assert user.validate_email() is False ``` ### Step 4: GREEN - Implement Minimum Code Write the minimum code necessary to make tests pass: - Focus on making tests green, not perfection - Avoid premature optimization - Keep implementation simple - Run tests - they should PASS ### Step 5: REFACTOR - Improve Clarity With green tests, improve the code: - Extract common patterns - Improve naming - Remove duplication - Simplify logic - Run tests after each change - they should remain GREEN ### Step 6: Verify Coverage Check test coverage meets the 80% target: ```bash pytest --cov=module --cov-report=term-missing ``` If coverage is below 80%: - Identify uncovered lines - Add tests for missing paths - Re-run coverage check ### Step 7: Document Deviations If implementation deviated from plan or introduced new dependencies: - Update tech-stack.md with new dependencies - Note deviations in plan.md task comments - Update spec.md if requirements changed ### Step 8: Commit Implementation Create a focused commit for the task: ```bash git add -A git commit -m "feat(user): implement email validation - Add validate_email method to User class - Handle empty and malformed emails - Add comprehensive test coverage Task: 2.1 Track: user-auth_20250115" ``` Commit message format: - Type: feat, fix, refactor, test, docs, chore - Scope: affected module or component - Summary: imperative, present tense - Body: bullet points of changes - Footer: task and track references ### Step 9: Attach Git Notes Add rich task summary as git note: ```bash git notes add -m "Task 2.1: Implement user validation Summary: - Added email validation using regex pattern - Handles edge cases: empty, no @, no domain - Coverage: 94% on validation module Files changed: - src/models/user.py (modified) - tests/test_user.py (modified) Decisions: - Used simple regex over email-validator library - Reason: No external dependency for basic validation" ``` ### Step 10: Update Plan with SHA Update plan.md to mark task complete with commit SHA: ```markdown - [x] **Task 2.1**: Implement user validation `abc1234` ``` ### Step 11: Commit Plan Update Commit the plan status update: ```bash git add conductor/tracks/*/plan.md git commit -m "docs: update plan - task 2.1 complete Track: user-auth_20250115" ``` ## Phase Completion Protocol When all tasks in a phase are complete, execute the verification protocol: ### Identify Changed Files List all files modified since the last checkpoint: ```bash git diff --name-only <last-checkpoint-sha>..HEAD ``` ### Ensure Test Coverage For each modified file: 1. Identify corresponding test file 2. Verify tests exist for new/changed code 3. Run coverage for modified modules 4. Add tests if coverage < 80% ### Run Full Test Suite Execute complete test suite: ```bash pytest -v --tb=short ``` All tests must pass before proceeding. ### Generate Manual Verification Steps Create checklist of manual verifications: ```markdown ## Phase 1 Verification Checklist - [ ] User can register with valid email - [ ] Invalid email shows appropriate error - [ ] Database stores user correctly - [ ] API returns expected response codes ``` ### WAIT for User Approval Present verification checklist to user: ``` Phase 1 complete. Please verify: 1. [ ] Test suite passes (automated) 2. [ ] Coverage meets target (automated) 3. [ ] Manual verification items (requires human) Respond with 'approved' to continue, or note issues. ``` Do NOT proceed without explicit approval. ### Create Checkpoint Commit After approval, create checkpoint commit: ```bash git add -A git commit -m "checkpoint: phase 1 complete - user-auth_20250115 Verified: - All tests passing - Coverage: 87% - Manual verification approved Phase 1 tasks: - [x] Task 1.1: Setup database schema - [x] Task 1.2: Implement user model - [x] Task 1.3: Add validation logic" ``` ### Record Checkpoint SHA Update plan.md checkpoints table: ```markdown ## Checkpoints | Phase | Checkpoint SHA | Date | Status | | ------- | -------------- | ---------- | -------- | | Phase 1 | def5678 | 2025-01-15 | verified | | Phase 2 | | | pending | ``` ## Quality Assurance Gates Before marking any task complete, verify these gates: ### Passing Tests - All existing tests pass - New tests pass - No test regressions ### Coverage >= 80% - New code has 80%+ coverage - Overall project coverage maintained - Critical paths fully covered ### Style Compliance - Code follows style guides - Linting passes - Formatting correct ### Documentation - Public APIs documented - Complex logic explained - README updated if needed ### Type Safety - Type hints present (if applicable) - Type checker passes - No type: ignore without reason ### No Linting Errors - Zero linter errors - Warnings addressed or justified - Static analysis clean ### Mobile Compatibility If applicable: - Responsive design verified - Touch interactions work - Performance acceptable ### Security Audit - No secrets in code - Input validation present - Authentication/authorization correct - Dependencies vulnerability-free ## Git Integration ### Commit Message Format ``` <type>(<scope>): <subject> <body> <footer> ``` Types: - `feat`: New feature - `fix`: Bug fix - `refactor`: Code change without feature/fix - `test`: Adding tests - `docs`: Documentation - `chore`: Maintenance ### Git Notes for Rich Summaries Attach detailed notes to commits: ```bash git notes add -m "<detailed summary>" ``` View notes: ```bash git log --show-notes ``` Benefits: - Preserves context without cluttering commit message - Enables semantic queries across commits - Supports track-based operations ### SHA Recording in plan.md Always record the commit SHA when completing tasks: ```markdown - [x] **Task 1.1**: Setup schema `abc1234` - [x] **Task 1.2**: Add model `def5678` ``` This enables: - Traceability from plan to code - Semantic revert operations - Progress auditing ## Verification Checkpoints ### Why Checkpoints Matter Checkpoints create restore points for semantic reversion: - Revert to end of any phase - Maintain logical code state - Enable safe experimentation ### When to Create Checkpoints Create checkpoint after: - All phase tasks complete - All phase verifications pass - User approval received ### Checkpoint Commit Content Include in checkpoint commit: - All uncommitted changes - Updated plan.md - Updated metadata.json - Any documentation updates ### How to Use Checkpoints For reverting: ```bash # Revert to end of Phase 1 git revert --no-commit <phase-2-commits>... git commit -m "revert: rollback to phase 1 checkpoint" ``` For review: ```bash # See what changed in Phase 2 git diff <phase-1-sha>..<phase-2-sha> ``` ## Handling Deviations During implementation, deviations from the plan may occur. Handle them systematically: ### Types of Deviations **Scope Addition** Discovered requirement not in original spec. - Document in spec.md as new requirement - Add tasks to plan.md - Note addition in task comments **Scope Reduction** Feature deemed unnecessary during implementation. - Mark tasks as `[-]` (skipped) with reason - Update spec.md scope section - Document decision rationale **Technical Deviation** Different implementation approach than planned. - Note deviation in task completion comment - Update tech-stack.md if dependencies changed - Document why original approach was unsuitable **Requirement Change** Understanding of requirement changes during work. - Update spec.md with corrected requirement - Adjust plan.md tasks if needed - Re-verify acceptance criteria ### Deviation Documentation Format When completing a task with deviation: ```markdown - [x] **Task 2.1**: Implement validation `abc1234` - DEVIATION: Used library instead of custom code - Reason: Better edge case handling - Impact: Added email-validator to dependencies ``` ## Error Recovery ### Failed Tests After GREEN If tests fail after reaching GREEN: 1. Do NOT proceed to REFACTOR 2. Identify which test started failing 3. Check if refactoring broke something 4. Revert to last known GREEN state 5. Re-approach the implementation ### Checkpoint Rejection If user rejects a checkpoint: 1. Note rejection reason in plan.md 2. Create tasks to address issues 3. Complete remediation tasks 4. Request checkpoint approval again ### Blocked by Dependency If task cannot proceed: 1. Mark task as `[!]` with blocker description 2. Check if other tasks can proceed 3. Document expected resolution timeline 4. Consider creating dependency resolution track ## TDD Variations by Task Type ### Data Model Tasks ``` RED: Write test for model creation and validation GREEN: Implement model class with fields REFACTOR: Add computed properties, improve types ``` ### API Endpoint Tasks ``` RED: Write test for request/response contract GREEN: Implement endpoint handler REFACTOR: Extract validation, improve error handling ``` ### Integration Tasks ``` RED: Write test for component interaction GREEN: Wire components together REFACTOR: Improve error propagation, add logging ``` ### Refactoring Tasks ``` RED: Add characterization tests for current behavior GREEN: Apply refactoring (tests should stay green) REFACTOR: Clean up any introduced complexity ``` ## Working with Existing Tests When modifying code with existing tests: ### Extend, Don't Replace - Keep existing tests passing - Add new tests for new behavior - Update tests only when requirements change ### Test Migration When refactoring changes test structure: 1. Run existing tests (should pass) 2. Add new tests for refactored code 3. Migrate test cases to new structure 4. Remove old tests only after new tests pass ### Regression Prevention After any change: 1. Run full test suite 2. Check for unexpected failures 3. Investigate any new failures 4. Fix regressions before proceeding ## Checkpoint Verification Details ### Automated Verification Run before requesting approval: ```bash # Test suite pytest -v --tb=short # Coverage pytest --cov=src --cov-report=term-missing # Linting ruff check src/ tests/ # Type checking (if applicable) mypy src/ ``` ### Manual Verification Guidance For manual items, provide specific instructions: ```markdown ## Manual Verification Steps ### User Registration 1. Navigate to /register 2. Enter valid email: test@example.com 3. Enter password meeting requirements 4. Click Submit 5. Verify success message appears 6. Verify user appears in database ### Error Handling 1. Enter invalid email: "notanemail" 2. Verify error message shows 3. Verify form retains other entered data ``` ## Performance Considerations ### Test Suite Performance Keep test suite fast: - Use fixtures to avoid redundant setup - Mock slow external calls - Run subset during development, full suite at checkpoints ### Commit Performance Keep commits atomic: - One logical change per commit - Complete thought, not work-in-progress - Tests should pass after every commit ## Best Practices 1. **Never skip RED**: Always write failing tests first 2. **Small commits**: One logical change per commit 3. **Immediate updates**: Update plan.md right after task completion 4. **Wait for approval**: Never skip checkpoint verification 5. **Rich git notes**: Include context that helps future understanding 6. **Coverage discipline**: Don't accept coverage below target 7. **Quality gates**: Check all gates before marking complete 8. **Sequential phases**: Complete phases in order 9. **Document deviations**: Note any changes from original plan 10. **Clean state**: Each commit should leave code in working state 11. **Fast feedback**: Run relevant tests frequently during development 12. **Clear blockers**: Address blockers promptly, don't work around them
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

airflow-dag-patterns

Build production Apache Airflow DAGs with best practices for

data
โญ1
# Apache Airflow DAG Patterns Production-ready patterns for Apache Airflow including DAG design, operators, sensors, testing, and deployment strategies. ## When to Use This Skill - Creating data pipeline orchestration with Airflow - Designing DAG structures and dependencies - Implementing custom operators and sensors - Testing Airflow DAGs locally - Setting up Airflow in production - Debugging failed DAG runs ## Core Concepts ### 1. DAG Design Principles | Principle | Description | | --------------- | ----------------------------------- | | **Idempotent** | Running twice produces same result | | **Atomic** | Tasks succeed or fail completely | | **Incremental** | Process only new/changed data | | **Observable** | Logs, metrics, alerts at every step | ### 2. Task Dependencies ```python # Linear task1 >> task2 >> task3 # Fan-out task1 >> [task2, task3, task4] # Fan-in [task1, task2, task3] >> task4 # Complex task1 >> task2 >> task4 task1 >> task3 >> task4 ``` ## Quick Start ```python # dags/example_dag.py from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.empty import EmptyOperator default_args = { 'owner': 'data-team', 'depends_on_past': False, 'email_on_failure': True, 'email_on_retry': False, 'retries': 3, 'retry_delay': timedelta(minutes=5), 'retry_exponential_backoff': True, 'max_retry_delay': timedelta(hours=1), } with DAG( dag_id='example_etl', default_args=default_args, description='Example ETL pipeline', schedule='0 6 * * *', # Daily at 6 AM start_date=datetime(2024, 1, 1), catchup=False, tags=['etl', 'example'], max_active_runs=1, ) as dag: start = EmptyOperator(task_id='start') def extract_data(**context): execution_date = context['ds'] # Extract logic here return {'records': 1000} extract = PythonOperator( task_id='extract', python_callable=extract_data, ) end = EmptyOperator(task_id='end') start >> extract >> end ``` ## Patterns ### Pattern 1: TaskFlow API (Airflow 2.0+) ```python # dags/taskflow_example.py from datetime import datetime from airflow.decorators import dag, task from airflow.models import Variable @dag( dag_id='taskflow_etl', schedule='@daily', start_date=datetime(2024, 1, 1), catchup=False, tags=['etl', 'taskflow'], ) def taskflow_etl(): """ETL pipeline using TaskFlow API""" @task() def extract(source: str) -> dict: """Extract data from source""" import pandas as pd df = pd.read_csv(f's3://bucket/{source}/{{ ds }}.csv') return {'data': df.to_dict(), 'rows': len(df)} @task() def transform(extracted: dict) -> dict: """Transform extracted data""" import pandas as pd df = pd.DataFrame(extracted['data']) df['processed_at'] = datetime.now() df = df.dropna() return {'data': df.to_dict(), 'rows': len(df)} @task() def load(transformed: dict, target: str): """Load data to target""" import pandas as pd df = pd.DataFrame(transformed['data']) df.to_parquet(f's3://bucket/{target}/{{ ds }}.parquet') return transformed['rows'] @task() def notify(rows_loaded: int): """Send notification""" print(f'Loaded {rows_loaded} rows') # Define dependencies with XCom passing extracted = extract(source='raw_data') transformed = transform(extracted) loaded = load(transformed, target='processed_data') notify(loaded) # Instantiate the DAG taskflow_etl() ``` ### Pattern 2: Dynamic DAG Generation ```python # dags/dynamic_dag_factory.py from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.models import Variable import json # Configuration for multiple similar pipelines PIPELINE_CONFIGS = [ {'name': 'customers', 'schedule': '@daily', 'source': 's3://raw/customers'}, {'name': 'orders', 'schedule': '@hourly', 'source': 's3://raw/orders'}, {'name': 'products', 'schedule': '@weekly', 'source': 's3://raw/products'}, ] def create_dag(config: dict) -> DAG: """Factory function to create DAGs from config""" dag_id = f"etl_{config['name']}" default_args = { 'owner': 'data-team', 'retries': 3, 'retry_delay': timedelta(minutes=5), } dag = DAG( dag_id=dag_id, default_args=default_args, schedule=config['schedule'], start_date=datetime(2024, 1, 1), catchup=False, tags=['etl', 'dynamic', config['name']], ) with dag: def extract_fn(source, **context): print(f"Extracting from {source} for {context['ds']}") def transform_fn(**context): print(f"Transforming data for {context['ds']}") def load_fn(table_name, **context): print(f"Loading to {table_name} for {context['ds']}") extract = PythonOperator( task_id='extract', python_callable=extract_fn, op_kwargs={'source': config['source']}, ) transform = PythonOperator( task_id='transform', python_callable=transform_fn, ) load = PythonOperator( task_id='load', python_callable=load_fn, op_kwargs={'table_name': config['name']}, ) extract >> transform >> load return dag # Generate DAGs for config in PIPELINE_CONFIGS: globals()[f"dag_{config['name']}"] = create_dag(config) ``` ### Pattern 3: Branching and Conditional Logic ```python # dags/branching_example.py from airflow.decorators import dag, task from airflow.operators.python import BranchPythonOperator from airflow.operators.empty import EmptyOperator from airflow.utils.trigger_rule import TriggerRule @dag( dag_id='branching_pipeline', schedule='@daily', start_date=datetime(2024, 1, 1), catchup=False, ) def branching_pipeline(): @task() def check_data_quality() -> dict: """Check data quality and return metrics""" quality_score = 0.95 # Simulated return {'score': quality_score, 'rows': 10000} def choose_branch(**context) -> str: """Determine which branch to execute""" ti = context['ti'] metrics = ti.xcom_pull(task_ids='check_data_quality') if metrics['score'] >= 0.9: return 'high_quality_path' elif metrics['score'] >= 0.7: return 'medium_quality_path' else: return 'low_quality_path' quality_check = check_data_quality() branch = BranchPythonOperator( task_id='branch', python_callable=choose_branch, ) high_quality = EmptyOperator(task_id='high_quality_path') medium_quality = EmptyOperator(task_id='medium_quality_path') low_quality = EmptyOperator(task_id='low_quality_path') # Join point - runs after any branch completes join = EmptyOperator( task_id='join', trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, ) quality_check >> branch >> [high_quality, medium_quality, low_quality] >> join branching_pipeline() ``` ### Pattern 4: Sensors and External Dependencies ```python # dags/sensor_patterns.py from datetime import datetime, timedelta from airflow import DAG from airflow.sensors.filesystem import FileSensor from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor from airflow.sensors.external_task import ExternalTaskSensor from airflow.operators.python import PythonOperator with DAG( dag_id='sensor_example', schedule='@daily', start_date=datetime(2024, 1, 1), catchup=False, ) as dag: # Wait for file on S3 wait_for_file = S3KeySensor( task_id='wait_for_s3_file', bucket_name='data-lake', bucket_key='raw/{{ ds }}/data.parquet', aws_conn_id='aws_default', timeout=60 * 60 * 2, # 2 hours poke_interval=60 * 5, # Check every 5 minutes mode='reschedule', # Free up worker slot while waiting ) # Wait for another DAG to complete wait_for_upstream = ExternalTaskSensor( task_id='wait_for_upstream_dag', external_dag_id='upstream_etl', external_task_id='final_task', execution_date_fn=lambda dt: dt, # Same execution date timeout=60 * 60 * 3, mode='reschedule', ) # Custom sensor using @task.sensor decorator @task.sensor(poke_interval=60, timeout=3600, mode='reschedule') def wait_for_api() -> PokeReturnValue: """Custom sensor for API availability""" import requests response = requests.get('https://api.example.com/health') is_done = response.status_code == 200 return PokeReturnValue(is_done=is_done, xcom_value=response.json()) api_ready = wait_for_api() def process_data(**context): api_result = context['ti'].xcom_pull(task_ids='wait_for_api') print(f"API returned: {api_result}") process = PythonOperator( task_id='process', python_callable=process_data, ) [wait_for_file, wait_for_upstream, api_ready] >> process ``` ### Pattern 5: Error Handling and Alerts ```python # dags/error_handling.py from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.utils.trigger_rule import TriggerRule from airflow.models import Variable def task_failure_callback(context): """Callback on task failure""" task_instance = context['task_instance'] exception = context.get('exception') # Send to Slack/PagerDuty/etc message = f""" Task Failed! DAG: {task_instance.dag_id} Task: {task_instance.task_id} Execution Date: {context['ds']} Error: {exception} Log URL: {task_instance.log_url} """ # send_slack_alert(message) print(message) def dag_failure_callback(context): """Callback on DAG failure""" # Aggregate failures, send summary pass with DAG( dag_id='error_handling_example', schedule='@daily', start_date=datetime(2024, 1, 1), catchup=False, on_failure_callback=dag_failure_callback, default_args={ 'on_failure_callback': task_failure_callback, 'retries': 3, 'retry_delay': timedelta(minutes=5), }, ) as dag: def might_fail(**context): import random if random.random() < 0.3: raise ValueError("Random failure!") return "Success" risky_task = PythonOperator( task_id='risky_task', python_callable=might_fail, ) def cleanup(**context): """Cleanup runs regardless of upstream failures""" print("Cleaning up...") cleanup_task = PythonOperator( task_id='cleanup', python_callable=cleanup, trigger_rule=TriggerRule.ALL_DONE, # Run even if upstream fails ) def notify_success(**context): """Only runs if all upstream succeeded""" print("All tasks succeeded!") success_notification = PythonOperator( task_id='notify_success', python_callable=notify_success, trigger_rule=TriggerRule.ALL_SUCCESS, ) risky_task >> [cleanup_task, success_notification] ``` ### Pattern 6: Testing DAGs ```python # tests/test_dags.py import pytest from datetime import datetime from airflow.models import DagBag @pytest.fixture def dagbag(): return DagBag(dag_folder='dags/', include_examples=False) def test_dag_loaded(dagbag): """Test that all DAGs load without errors""" assert len(dagbag.import_errors) == 0, f"DAG import errors: {dagbag.import_errors}" def test_dag_structure(dagbag): """Test specific DAG structure""" dag = dagbag.get_dag('example_etl') assert dag is not None assert len(dag.tasks) == 3 assert dag.schedule_interval == '0 6 * * *' def test_task_dependencies(dagbag): """Test task dependencies are correct""" dag = dagbag.get_dag('example_etl') extract_task = dag.get_task('extract') assert 'start' in [t.task_id for t in extract_task.upstream_list] assert 'end' in [t.task_id for t in extract_task.downstream_list] def test_dag_integrity(dagbag): """Test DAG has no cycles and is valid""" for dag_id, dag in dagbag.dags.items(): assert dag.test_cycle() is None, f"Cycle detected in {dag_id}" # Test individual task logic def test_extract_function(): """Unit test for extract function""" from dags.example_dag import extract_data result = extract_data(ds='2024-01-01') assert 'records' in result assert isinstance(result['records'], int) ``` ## Project Structure ``` airflow/ โ”œโ”€โ”€ dags/ โ”‚ โ”œโ”€โ”€ __init__.py โ”‚ โ”œโ”€โ”€ common/ โ”‚ โ”‚ โ”œโ”€โ”€ __init__.py โ”‚ โ”‚ โ”œโ”€โ”€ operators.py # Custom operators โ”‚ โ”‚ โ”œโ”€โ”€ sensors.py # Custom sensors โ”‚ โ”‚ โ””โ”€โ”€ callbacks.py # Alert callbacks โ”‚ โ”œโ”€โ”€ etl/ โ”‚ โ”‚ โ”œโ”€โ”€ customers.py โ”‚ โ”‚ โ””โ”€โ”€ orders.py โ”‚ โ””โ”€โ”€ ml/ โ”‚ โ””โ”€โ”€ training.py โ”œโ”€โ”€ plugins/ โ”‚ โ””โ”€โ”€ custom_plugin.py โ”œโ”€โ”€ tests/ โ”‚ โ”œโ”€โ”€ __init__.py โ”‚ โ”œโ”€โ”€ test_dags.py โ”‚ โ””โ”€โ”€ test_operators.py โ”œโ”€โ”€ docker-compose.yml โ””โ”€โ”€ requirements.txt ``` ## Best Practices ### Do's - **Use TaskFlow API** - Cleaner code, automatic XCom - **Set timeouts** - Prevent zombie tasks - **Use `mode='reschedule'`** - For sensors, free up workers - **Test DAGs** - Unit tests and integration tests - **Idempotent tasks** - Safe to retry ### Don'ts - **Don't use `depends_on_past=True`** - Creates bottlenecks - **Don't hardcode dates** - Use `{{ ds }}` macros - **Don't use global state** - Tasks should be stateless - **Don't skip catchup blindly** - Understand implications - **Don't put heavy logic in DAG file** - Import from modules ## Resources - [Airflow Documentation](https://airflow.apache.org/docs/) - [Astronomer Guides](https://docs.astronomer.io/learn) - [TaskFlow API](https://airflow.apache.org/docs/apache-airflow/stable/tutorial/taskflow.html)
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

data-quality-frameworks

Implement data quality validation with Great Expectations, dbt

data
โญ1
# Data Quality Frameworks Production patterns for implementing data quality with Great Expectations, dbt tests, and data contracts to ensure reliable data pipelines. ## When to Use This Skill - Implementing data quality checks in pipelines - Setting up Great Expectations validation - Building comprehensive dbt test suites - Establishing data contracts between teams - Monitoring data quality metrics - Automating data validation in CI/CD ## Core Concepts ### 1. Data Quality Dimensions | Dimension | Description | Example Check | | ---------------- | ------------------------ | -------------------------------------------------- | | **Completeness** | No missing values | `expect_column_values_to_not_be_null` | | **Uniqueness** | No duplicates | `expect_column_values_to_be_unique` | | **Validity** | Values in expected range | `expect_column_values_to_be_in_set` | | **Accuracy** | Data matches reality | Cross-reference validation | | **Consistency** | No contradictions | `expect_column_pair_values_A_to_be_greater_than_B` | | **Timeliness** | Data is recent | `expect_column_max_to_be_between` | ### 2. Testing Pyramid for Data ``` /\ / \ Integration Tests (cross-table) /โ”€โ”€โ”€โ”€\ / \ Unit Tests (single column) /โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€\ / \ Schema Tests (structure) /โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€\ ``` ## Quick Start ### Great Expectations Setup ```bash # Install pip install great_expectations # Initialize project great_expectations init # Create datasource great_expectations datasource new ``` ```python # great_expectations/checkpoints/daily_validation.yml import great_expectations as gx # Create context context = gx.get_context() # Create expectation suite suite = context.add_expectation_suite("orders_suite") # Add expectations suite.add_expectation( gx.expectations.ExpectColumnValuesToNotBeNull(column="order_id") ) suite.add_expectation( gx.expectations.ExpectColumnValuesToBeUnique(column="order_id") ) # Validate results = context.run_checkpoint(checkpoint_name="daily_orders") ``` ## Patterns ### Pattern 1: Great Expectations Suite ```python # expectations/orders_suite.py import great_expectations as gx from great_expectations.core import ExpectationSuite from great_expectations.core.expectation_configuration import ExpectationConfiguration def build_orders_suite() -> ExpectationSuite: """Build comprehensive orders expectation suite""" suite = ExpectationSuite(expectation_suite_name="orders_suite") # Schema expectations suite.add_expectation(ExpectationConfiguration( expectation_type="expect_table_columns_to_match_set", kwargs={ "column_set": ["order_id", "customer_id", "amount", "status", "created_at"], "exact_match": False # Allow additional columns } )) # Primary key suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_not_be_null", kwargs={"column": "order_id"} )) suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_be_unique", kwargs={"column": "order_id"} )) # Foreign key suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_not_be_null", kwargs={"column": "customer_id"} )) # Categorical values suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_be_in_set", kwargs={ "column": "status", "value_set": ["pending", "processing", "shipped", "delivered", "cancelled"] } )) # Numeric ranges suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_be_between", kwargs={ "column": "amount", "min_value": 0, "max_value": 100000, "strict_min": True # amount > 0 } )) # Date validity suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_values_to_be_dateutil_parseable", kwargs={"column": "created_at"} )) # Freshness - data should be recent suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_max_to_be_between", kwargs={ "column": "created_at", "min_value": {"$PARAMETER": "now - timedelta(days=1)"}, "max_value": {"$PARAMETER": "now"} } )) # Row count sanity suite.add_expectation(ExpectationConfiguration( expectation_type="expect_table_row_count_to_be_between", kwargs={ "min_value": 1000, # Expect at least 1000 rows "max_value": 10000000 } )) # Statistical expectations suite.add_expectation(ExpectationConfiguration( expectation_type="expect_column_mean_to_be_between", kwargs={ "column": "amount", "min_value": 50, "max_value": 500 } )) return suite ``` ### Pattern 2: Great Expectations Checkpoint ```yaml # great_expectations/checkpoints/orders_checkpoint.yml name: orders_checkpoint config_version: 1.0 class_name: Checkpoint run_name_template: "%Y%m%d-%H%M%S-orders-validation" validations: - batch_request: datasource_name: warehouse data_connector_name: default_inferred_data_connector_name data_asset_name: orders data_connector_query: index: -1 # Latest batch expectation_suite_name: orders_suite action_list: - name: store_validation_result action: class_name: StoreValidationResultAction - name: store_evaluation_parameters action: class_name: StoreEvaluationParametersAction - name: update_data_docs action: class_name: UpdateDataDocsAction # Slack notification on failure - name: send_slack_notification action: class_name: SlackNotificationAction slack_webhook: ${SLACK_WEBHOOK} notify_on: failure renderer: module_name: great_expectations.render.renderer.slack_renderer class_name: SlackRenderer ``` ```python # Run checkpoint import great_expectations as gx context = gx.get_context() result = context.run_checkpoint(checkpoint_name="orders_checkpoint") if not result.success: failed_expectations = [ r for r in result.run_results.values() if not r.success ] raise ValueError(f"Data quality check failed: {failed_expectations}") ``` ### Pattern 3: dbt Data Tests ```yaml # models/marts/core/_core__models.yml version: 2 models: - name: fct_orders description: Order fact table tests: # Table-level tests - dbt_utils.recency: datepart: day field: created_at interval: 1 - dbt_utils.at_least_one - dbt_utils.expression_is_true: expression: "total_amount >= 0" columns: - name: order_id description: Primary key tests: - unique - not_null - name: customer_id description: Foreign key to dim_customers tests: - not_null - relationships: to: ref('dim_customers') field: customer_id - name: order_status tests: - accepted_values: values: ["pending", "processing", "shipped", "delivered", "cancelled"] - name: total_amount tests: - not_null - dbt_utils.expression_is_true: expression: ">= 0" - name: created_at tests: - not_null - dbt_utils.expression_is_true: expression: "<= current_timestamp" - name: dim_customers columns: - name: customer_id tests: - unique - not_null - name: email tests: - unique - not_null # Custom regex test - dbt_utils.expression_is_true: expression: "email ~ '^[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\\.[A-Za-z]{2,}$'" ``` ### Pattern 4: Custom dbt Tests ```sql -- tests/generic/test_row_count_in_range.sql {% test row_count_in_range(model, min_count, max_count) %} with row_count as ( select count(*) as cnt from {{ model }} ) select cnt from row_count where cnt < {{ min_count }} or cnt > {{ max_count }} {% endtest %} -- Usage in schema.yml: -- tests: -- - row_count_in_range: -- min_count: 1000 -- max_count: 10000000 ``` ```sql -- tests/generic/test_sequential_values.sql {% test sequential_values(model, column_name, interval=1) %} with lagged as ( select {{ column_name }}, lag({{ column_name }}) over (order by {{ column_name }}) as prev_value from {{ model }} ) select * from lagged where {{ column_name }} - prev_value != {{ interval }} and prev_value is not null {% endtest %} ``` ```sql -- tests/singular/assert_orders_customers_match.sql -- Singular test: specific business rule with orders_customers as ( select distinct customer_id from {{ ref('fct_orders') }} ), dim_customers as ( select customer_id from {{ ref('dim_customers') }} ), orphaned_orders as ( select o.customer_id from orders_customers o left join dim_customers c using (customer_id) where c.customer_id is null ) select * from orphaned_orders -- Test passes if this returns 0 rows ``` ### Pattern 5: Data Contracts ```yaml # contracts/orders_contract.yaml apiVersion: datacontract.com/v1.0.0 kind: DataContract metadata: name: orders version: 1.0.0 owner: data-platform-team contact: data-team@company.com info: title: Orders Data Contract description: Contract for order event data from the ecommerce platform purpose: Analytics, reporting, and ML features servers: production: type: snowflake account: company.us-east-1 database: ANALYTICS schema: CORE terms: usage: Internal analytics only limitations: PII must not be exposed in downstream marts billing: Charged per query TB scanned schema: type: object properties: order_id: type: string format: uuid description: Unique order identifier required: true unique: true pii: false customer_id: type: string format: uuid description: Customer identifier required: true pii: true piiClassification: indirect total_amount: type: number minimum: 0 maximum: 100000 description: Order total in USD created_at: type: string format: date-time description: Order creation timestamp required: true status: type: string enum: [pending, processing, shipped, delivered, cancelled] description: Current order status quality: type: SodaCL specification: checks for orders: - row_count > 0 - missing_count(order_id) = 0 - duplicate_count(order_id) = 0 - invalid_count(status) = 0: valid values: [pending, processing, shipped, delivered, cancelled] - freshness(created_at) < 24h sla: availability: 99.9% freshness: 1 hour latency: 5 minutes ``` ### Pattern 6: Automated Quality Pipeline ```python # quality_pipeline.py from dataclasses import dataclass from typing import List, Dict, Any import great_expectations as gx from datetime import datetime @dataclass class QualityResult: table: str passed: bool total_expectations: int failed_expectations: int details: List[Dict[str, Any]] timestamp: datetime class DataQualityPipeline: """Orchestrate data quality checks across tables""" def __init__(self, context: gx.DataContext): self.context = context self.results: List[QualityResult] = [] def validate_table(self, table: str, suite: str) -> QualityResult: """Validate a single table against expectation suite""" checkpoint_config = { "name": f"{table}_validation", "config_version": 1.0, "class_name": "Checkpoint", "validations": [{ "batch_request": { "datasource_name": "warehouse", "data_asset_name": table, }, "expectation_suite_name": suite, }], } result = self.context.run_checkpoint(**checkpoint_config) # Parse results validation_result = list(result.run_results.values())[0] results = validation_result.results failed = [r for r in results if not r.success] return QualityResult( table=table, passed=result.success, total_expectations=len(results), failed_expectations=len(failed), details=[{ "expectation": r.expectation_config.expectation_type, "success": r.success, "observed_value": r.result.get("observed_value"), } for r in results], timestamp=datetime.now() ) def run_all(self, tables: Dict[str, str]) -> Dict[str, QualityResult]: """Run validation for all tables""" results = {} for table, suite in tables.items(): print(f"Validating {table}...") results[table] = self.validate_table(table, suite) return results def generate_report(self, results: Dict[str, QualityResult]) -> str: """Generate quality report""" report = ["# Data Quality Report", f"Generated: {datetime.now()}", ""] total_passed = sum(1 for r in results.values() if r.passed) total_tables = len(results) report.append(f"## Summary: {total_passed}/{total_tables} tables passed") report.append("") for table, result in results.items(): status = "โœ…" if result.passed else "โŒ" report.append(f"### {status} {table}") report.append(f"- Expectations: {result.total_expectations}") report.append(f"- Failed: {result.failed_expectations}") if not result.passed: report.append("- Failed checks:") for detail in result.details: if not detail["success"]: report.append(f" - {detail['expectation']}: {detail['observed_value']}") report.append("") return "\n".join(report) # Usage context = gx.get_context() pipeline = DataQualityPipeline(context) tables_to_validate = { "orders": "orders_suite", "customers": "customers_suite", "products": "products_suite", } results = pipeline.run_all(tables_to_validate) report = pipeline.generate_report(results) # Fail pipeline if any table failed if not all(r.passed for r in results.values()): print(report) raise ValueError("Data quality checks failed!") ``` ## Best Practices ### Do's - **Test early** - Validate source data before transformations - **Test incrementally** - Add tests as you find issues - **Document expectations** - Clear descriptions for each test - **Alert on failures** - Integrate with monitoring - **Version contracts** - Track schema changes ### Don'ts - **Don't test everything** - Focus on critical columns - **Don't ignore warnings** - They often precede failures - **Don't skip freshness** - Stale data is bad data - **Don't hardcode thresholds** - Use dynamic baselines - **Don't test in isolation** - Test relationships too ## Resources - [Great Expectations Documentation](https://docs.greatexpectations.io/) - [dbt Testing Documentation](https://docs.getdbt.com/docs/build/tests) - [Data Contract Specification](https://datacontract.com/) - [Soda Core](https://docs.soda.io/soda-core/overview.html)
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

dbt-transformation-patterns

Master dbt (data build tool) for analytics engineering with model

data
โญ1
# dbt Transformation Patterns Production-ready patterns for dbt (data build tool) including model organization, testing strategies, documentation, and incremental processing. ## When to Use This Skill - Building data transformation pipelines with dbt - Organizing models into staging, intermediate, and marts layers - Implementing data quality tests - Creating incremental models for large datasets - Documenting data models and lineage - Setting up dbt project structure ## Core Concepts ### 1. Model Layers (Medallion Architecture) ``` sources/ Raw data definitions โ†“ staging/ 1:1 with source, light cleaning โ†“ intermediate/ Business logic, joins, aggregations โ†“ marts/ Final analytics tables ``` ### 2. Naming Conventions | Layer | Prefix | Example | | ------------ | -------------- | ----------------------------- | | Staging | `stg_` | `stg_stripe__payments` | | Intermediate | `int_` | `int_payments_pivoted` | | Marts | `dim_`, `fct_` | `dim_customers`, `fct_orders` | ## Quick Start ```yaml # dbt_project.yml name: "analytics" version: "1.0.0" profile: "analytics" model-paths: ["models"] analysis-paths: ["analyses"] test-paths: ["tests"] seed-paths: ["seeds"] macro-paths: ["macros"] vars: start_date: "2020-01-01" models: analytics: staging: +materialized: view +schema: staging intermediate: +materialized: ephemeral marts: +materialized: table +schema: analytics ``` ``` # Project structure models/ โ”œโ”€โ”€ staging/ โ”‚ โ”œโ”€โ”€ stripe/ โ”‚ โ”‚ โ”œโ”€โ”€ _stripe__sources.yml โ”‚ โ”‚ โ”œโ”€โ”€ _stripe__models.yml โ”‚ โ”‚ โ”œโ”€โ”€ stg_stripe__customers.sql โ”‚ โ”‚ โ””โ”€โ”€ stg_stripe__payments.sql โ”‚ โ””โ”€โ”€ shopify/ โ”‚ โ”œโ”€โ”€ _shopify__sources.yml โ”‚ โ””โ”€โ”€ stg_shopify__orders.sql โ”œโ”€โ”€ intermediate/ โ”‚ โ””โ”€โ”€ finance/ โ”‚ โ””โ”€โ”€ int_payments_pivoted.sql โ””โ”€โ”€ marts/ โ”œโ”€โ”€ core/ โ”‚ โ”œโ”€โ”€ _core__models.yml โ”‚ โ”œโ”€โ”€ dim_customers.sql โ”‚ โ””โ”€โ”€ fct_orders.sql โ””โ”€โ”€ finance/ โ””โ”€โ”€ fct_revenue.sql ``` ## Patterns ### Pattern 1: Source Definitions ```yaml # models/staging/stripe/_stripe__sources.yml version: 2 sources: - name: stripe description: Raw Stripe data loaded via Fivetran database: raw schema: stripe loader: fivetran loaded_at_field: _fivetran_synced freshness: warn_after: { count: 12, period: hour } error_after: { count: 24, period: hour } tables: - name: customers description: Stripe customer records columns: - name: id description: Primary key tests: - unique - not_null - name: email description: Customer email - name: created description: Account creation timestamp - name: payments description: Stripe payment transactions columns: - name: id tests: - unique - not_null - name: customer_id tests: - not_null - relationships: to: source('stripe', 'customers') field: id ``` ### Pattern 2: Staging Models ```sql -- models/staging/stripe/stg_stripe__customers.sql with source as ( select * from {{ source('stripe', 'customers') }} ), renamed as ( select -- ids id as customer_id, -- strings lower(email) as email, name as customer_name, -- timestamps created as created_at, -- metadata _fivetran_synced as _loaded_at from source ) select * from renamed ``` ```sql -- models/staging/stripe/stg_stripe__payments.sql {{ config( materialized='incremental', unique_key='payment_id', on_schema_change='append_new_columns' ) }} with source as ( select * from {{ source('stripe', 'payments') }} {% if is_incremental() %} where _fivetran_synced > (select max(_loaded_at) from {{ this }}) {% endif %} ), renamed as ( select -- ids id as payment_id, customer_id, invoice_id, -- amounts (convert cents to dollars) amount / 100.0 as amount, amount_refunded / 100.0 as amount_refunded, -- status status as payment_status, -- timestamps created as created_at, -- metadata _fivetran_synced as _loaded_at from source ) select * from renamed ``` ### Pattern 3: Intermediate Models ```sql -- models/intermediate/finance/int_payments_pivoted_to_customer.sql with payments as ( select * from {{ ref('stg_stripe__payments') }} ), customers as ( select * from {{ ref('stg_stripe__customers') }} ), payment_summary as ( select customer_id, count(*) as total_payments, count(case when payment_status = 'succeeded' then 1 end) as successful_payments, sum(case when payment_status = 'succeeded' then amount else 0 end) as total_amount_paid, min(created_at) as first_payment_at, max(created_at) as last_payment_at from payments group by customer_id ) select customers.customer_id, customers.email, customers.created_at as customer_created_at, coalesce(payment_summary.total_payments, 0) as total_payments, coalesce(payment_summary.successful_payments, 0) as successful_payments, coalesce(payment_summary.total_amount_paid, 0) as lifetime_value, payment_summary.first_payment_at, payment_summary.last_payment_at from customers left join payment_summary using (customer_id) ``` ### Pattern 4: Mart Models (Dimensions and Facts) ```sql -- models/marts/core/dim_customers.sql {{ config( materialized='table', unique_key='customer_id' ) }} with customers as ( select * from {{ ref('int_payments_pivoted_to_customer') }} ), orders as ( select * from {{ ref('stg_shopify__orders') }} ), order_summary as ( select customer_id, count(*) as total_orders, sum(total_price) as total_order_value, min(created_at) as first_order_at, max(created_at) as last_order_at from orders group by customer_id ), final as ( select -- surrogate key {{ dbt_utils.generate_surrogate_key(['customers.customer_id']) }} as customer_key, -- natural key customers.customer_id, -- attributes customers.email, customers.customer_created_at, -- payment metrics customers.total_payments, customers.successful_payments, customers.lifetime_value, customers.first_payment_at, customers.last_payment_at, -- order metrics coalesce(order_summary.total_orders, 0) as total_orders, coalesce(order_summary.total_order_value, 0) as total_order_value, order_summary.first_order_at, order_summary.last_order_at, -- calculated fields case when customers.lifetime_value >= 1000 then 'high' when customers.lifetime_value >= 100 then 'medium' else 'low' end as customer_tier, -- timestamps current_timestamp as _loaded_at from customers left join order_summary using (customer_id) ) select * from final ``` ```sql -- models/marts/core/fct_orders.sql {{ config( materialized='incremental', unique_key='order_id', incremental_strategy='merge' ) }} with orders as ( select * from {{ ref('stg_shopify__orders') }} {% if is_incremental() %} where updated_at > (select max(updated_at) from {{ this }}) {% endif %} ), customers as ( select * from {{ ref('dim_customers') }} ), final as ( select -- keys orders.order_id, customers.customer_key, orders.customer_id, -- dimensions orders.order_status, orders.fulfillment_status, orders.payment_status, -- measures orders.subtotal, orders.tax, orders.shipping, orders.total_price, orders.total_discount, orders.item_count, -- timestamps orders.created_at, orders.updated_at, orders.fulfilled_at, -- metadata current_timestamp as _loaded_at from orders left join customers on orders.customer_id = customers.customer_id ) select * from final ``` ### Pattern 5: Testing and Documentation ```yaml # models/marts/core/_core__models.yml version: 2 models: - name: dim_customers description: Customer dimension with payment and order metrics columns: - name: customer_key description: Surrogate key for the customer dimension tests: - unique - not_null - name: customer_id description: Natural key from source system tests: - unique - not_null - name: email description: Customer email address tests: - not_null - name: customer_tier description: Customer value tier based on lifetime value tests: - accepted_values: values: ["high", "medium", "low"] - name: lifetime_value description: Total amount paid by customer tests: - dbt_utils.expression_is_true: expression: ">= 0" - name: fct_orders description: Order fact table with all order transactions tests: - dbt_utils.recency: datepart: day field: created_at interval: 1 columns: - name: order_id tests: - unique - not_null - name: customer_key tests: - not_null - relationships: to: ref('dim_customers') field: customer_key ``` ### Pattern 6: Macros and DRY Code ```sql -- macros/cents_to_dollars.sql {% macro cents_to_dollars(column_name, precision=2) %} round({{ column_name }} / 100.0, {{ precision }}) {% endmacro %} -- macros/generate_schema_name.sql {% macro generate_schema_name(custom_schema_name, node) %} {%- set default_schema = target.schema -%} {%- if custom_schema_name is none -%} {{ default_schema }} {%- else -%} {{ default_schema }}_{{ custom_schema_name }} {%- endif -%} {% endmacro %} -- macros/limit_data_in_dev.sql {% macro limit_data_in_dev(column_name, days=3) %} {% if target.name == 'dev' %} where {{ column_name }} >= dateadd(day, -{{ days }}, current_date) {% endif %} {% endmacro %} -- Usage in model select * from {{ ref('stg_orders') }} {{ limit_data_in_dev('created_at') }} ``` ### Pattern 7: Incremental Strategies ```sql -- Delete+Insert (default for most warehouses) {{ config( materialized='incremental', unique_key='id', incremental_strategy='delete+insert' ) }} -- Merge (best for late-arriving data) {{ config( materialized='incremental', unique_key='id', incremental_strategy='merge', merge_update_columns=['status', 'amount', 'updated_at'] ) }} -- Insert Overwrite (partition-based) {{ config( materialized='incremental', incremental_strategy='insert_overwrite', partition_by={ "field": "created_date", "data_type": "date", "granularity": "day" } ) }} select *, date(created_at) as created_date from {{ ref('stg_events') }} {% if is_incremental() %} where created_date >= dateadd(day, -3, current_date) {% endif %} ``` ## dbt Commands ```bash # Development dbt run # Run all models dbt run --select staging # Run staging models only dbt run --select +fct_orders # Run fct_orders and its upstream dbt run --select fct_orders+ # Run fct_orders and its downstream dbt run --full-refresh # Rebuild incremental models # Testing dbt test # Run all tests dbt test --select stg_stripe # Test specific models dbt build # Run + test in DAG order # Documentation dbt docs generate # Generate docs dbt docs serve # Serve docs locally # Debugging dbt compile # Compile SQL without running dbt debug # Test connection dbt ls --select tag:critical # List models by tag ``` ## Best Practices ### Do's - **Use staging layer** - Clean data once, use everywhere - **Test aggressively** - Not null, unique, relationships - **Document everything** - Column descriptions, model descriptions - **Use incremental** - For tables > 1M rows - **Version control** - dbt project in Git ### Don'ts - **Don't skip staging** - Raw โ†’ mart is tech debt - **Don't hardcode dates** - Use `{{ var('start_date') }}` - **Don't repeat logic** - Extract to macros - **Don't test in prod** - Use dev target - **Don't ignore freshness** - Monitor source data ## Resources - [dbt Documentation](https://docs.getdbt.com/) - [dbt Best Practices](https://docs.getdbt.com/guides/best-practices) - [dbt-utils Package](https://hub.getdbt.com/dbt-labs/dbt_utils/latest/) - [dbt Discourse](https://discourse.getdbt.com/)
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

dotnet-backend-patterns

Master C#/.NET backend development patterns for building robust

coding
โญ1
# .NET Backend Development Patterns Master C#/.NET patterns for building production-grade APIs, MCP servers, and enterprise backends with modern best practices (2024/2025). ## When to Use This Skill - Developing new .NET Web APIs or MCP servers - Reviewing C# code for quality and performance - Designing service architectures with dependency injection - Implementing caching strategies with Redis - Writing unit and integration tests - Optimizing database access with EF Core or Dapper - Configuring applications with IOptions pattern - Handling errors and implementing resilience patterns ## Core Concepts ### 1. Project Structure (Clean Architecture) ``` src/ โ”œโ”€โ”€ Domain/ # Core business logic (no dependencies) โ”‚ โ”œโ”€โ”€ Entities/ โ”‚ โ”œโ”€โ”€ Interfaces/ โ”‚ โ”œโ”€โ”€ Exceptions/ โ”‚ โ””โ”€โ”€ ValueObjects/ โ”œโ”€โ”€ Application/ # Use cases, DTOs, validation โ”‚ โ”œโ”€โ”€ Services/ โ”‚ โ”œโ”€โ”€ DTOs/ โ”‚ โ”œโ”€โ”€ Validators/ โ”‚ โ””โ”€โ”€ Interfaces/ โ”œโ”€โ”€ Infrastructure/ # External implementations โ”‚ โ”œโ”€โ”€ Data/ # EF Core, Dapper repositories โ”‚ โ”œโ”€โ”€ Caching/ # Redis, Memory cache โ”‚ โ”œโ”€โ”€ External/ # HTTP clients, third-party APIs โ”‚ โ””โ”€โ”€ DependencyInjection/ # Service registration โ””โ”€โ”€ Api/ # Entry point โ”œโ”€โ”€ Controllers/ # Or MinimalAPI endpoints โ”œโ”€โ”€ Middleware/ โ”œโ”€โ”€ Filters/ โ””โ”€โ”€ Program.cs ``` ### 2. Dependency Injection Patterns ```csharp // Service registration by lifetime public static class ServiceCollectionExtensions { public static IServiceCollection AddApplicationServices( this IServiceCollection services, IConfiguration configuration) { // Scoped: One instance per HTTP request services.AddScoped<IProductService, ProductService>(); services.AddScoped<IOrderService, OrderService>(); // Singleton: One instance for app lifetime services.AddSingleton<ICacheService, RedisCacheService>(); services.AddSingleton<IConnectionMultiplexer>(_ => ConnectionMultiplexer.Connect(configuration["Redis:Connection"]!)); // Transient: New instance every time services.AddTransient<IValidator<CreateOrderRequest>, CreateOrderValidator>(); // Options pattern for configuration services.Configure<CatalogOptions>(configuration.GetSection("Catalog")); services.Configure<RedisOptions>(configuration.GetSection("Redis")); // Factory pattern for conditional creation services.AddScoped<IPriceCalculator>(sp => { var options = sp.GetRequiredService<IOptions<PricingOptions>>().Value; return options.UseNewEngine ? sp.GetRequiredService<NewPriceCalculator>() : sp.GetRequiredService<LegacyPriceCalculator>(); }); // Keyed services (.NET 8+) services.AddKeyedScoped<IPaymentProcessor, StripeProcessor>("stripe"); services.AddKeyedScoped<IPaymentProcessor, PayPalProcessor>("paypal"); return services; } } // Usage with keyed services public class CheckoutService { public CheckoutService( [FromKeyedServices("stripe")] IPaymentProcessor stripeProcessor) { _processor = stripeProcessor; } } ``` ### 3. Async/Await Patterns ```csharp // โœ… CORRECT: Async all the way down public async Task<Product> GetProductAsync(string id, CancellationToken ct = default) { return await _repository.GetByIdAsync(id, ct); } // โœ… CORRECT: Parallel execution with WhenAll public async Task<(Stock, Price)> GetStockAndPriceAsync( string productId, CancellationToken ct = default) { var stockTask = _stockService.GetAsync(productId, ct); var priceTask = _priceService.GetAsync(productId, ct); await Task.WhenAll(stockTask, priceTask); return (await stockTask, await priceTask); } // โœ… CORRECT: ConfigureAwait in libraries public async Task<T> LibraryMethodAsync<T>(CancellationToken ct = default) { var result = await _httpClient.GetAsync(url, ct).ConfigureAwait(false); return await result.Content.ReadFromJsonAsync<T>(ct).ConfigureAwait(false); } // โœ… CORRECT: ValueTask for hot paths with caching public ValueTask<Product?> GetCachedProductAsync(string id) { if (_cache.TryGetValue(id, out Product? product)) return ValueTask.FromResult(product); return new ValueTask<Product?>(GetFromDatabaseAsync(id)); } // โŒ WRONG: Blocking on async (deadlock risk) var result = GetProductAsync(id).Result; // NEVER do this var result2 = GetProductAsync(id).GetAwaiter().GetResult(); // Also bad // โŒ WRONG: async void (except event handlers) public async void ProcessOrder() { } // Exceptions are lost // โŒ WRONG: Unnecessary Task.Run for already async code await Task.Run(async () => await GetDataAsync()); // Wastes thread ``` ### 4. Configuration with IOptions ```csharp // Configuration classes public class CatalogOptions { public const string SectionName = "Catalog"; public int DefaultPageSize { get; set; } = 50; public int MaxPageSize { get; set; } = 200; public TimeSpan CacheDuration { get; set; } = TimeSpan.FromMinutes(15); public bool EnableEnrichment { get; set; } = true; } public class RedisOptions { public const string SectionName = "Redis"; public string Connection { get; set; } = "localhost:6379"; public string KeyPrefix { get; set; } = "mcp:"; public int Database { get; set; } = 0; } // appsettings.json { "Catalog": { "DefaultPageSize": 50, "MaxPageSize": 200, "CacheDuration": "00:15:00", "EnableEnrichment": true }, "Redis": { "Connection": "localhost:6379", "KeyPrefix": "mcp:", "Database": 0 } } // Registration services.Configure<CatalogOptions>(configuration.GetSection(CatalogOptions.SectionName)); services.Configure<RedisOptions>(configuration.GetSection(RedisOptions.SectionName)); // Usage with IOptions (singleton, read once at startup) public class CatalogService { private readonly CatalogOptions _options; public CatalogService(IOptions<CatalogOptions> options) { _options = options.Value; } } // Usage with IOptionsSnapshot (scoped, re-reads on each request) public class DynamicService { private readonly CatalogOptions _options; public DynamicService(IOptionsSnapshot<CatalogOptions> options) { _options = options.Value; // Fresh value per request } } // Usage with IOptionsMonitor (singleton, notified on changes) public class MonitoredService { private CatalogOptions _options; public MonitoredService(IOptionsMonitor<CatalogOptions> monitor) { _options = monitor.CurrentValue; monitor.OnChange(newOptions => _options = newOptions); } } ``` ### 5. Result Pattern (Avoiding Exceptions for Flow Control) ```csharp // Generic Result type public class Result<T> { public bool IsSuccess { get; } public T? Value { get; } public string? Error { get; } public string? ErrorCode { get; } private Result(bool isSuccess, T? value, string? error, string? errorCode) { IsSuccess = isSuccess; Value = value; Error = error; ErrorCode = errorCode; } public static Result<T> Success(T value) => new(true, value, null, null); public static Result<T> Failure(string error, string? code = null) => new(false, default, error, code); public Result<TNew> Map<TNew>(Func<T, TNew> mapper) => IsSuccess ? Result<TNew>.Success(mapper(Value!)) : Result<TNew>.Failure(Error!, ErrorCode); public async Task<Result<TNew>> MapAsync<TNew>(Func<T, Task<TNew>> mapper) => IsSuccess ? Result<TNew>.Success(await mapper(Value!)) : Result<TNew>.Failure(Error!, ErrorCode); } // Usage in service public async Task<Result<Order>> CreateOrderAsync(CreateOrderRequest request, CancellationToken ct) { // Validation var validation = await _validator.ValidateAsync(request, ct); if (!validation.IsValid) return Result<Order>.Failure( validation.Errors.First().ErrorMessage, "VALIDATION_ERROR"); // Business rule check var stock = await _stockService.CheckAsync(request.ProductId, request.Quantity, ct); if (!stock.IsAvailable) return Result<Order>.Failure( $"Insufficient stock: {stock.Available} available, {request.Quantity} requested", "INSUFFICIENT_STOCK"); // Create order var order = await _repository.CreateAsync(request.ToEntity(), ct); return Result<Order>.Success(order); } // Usage in controller/endpoint app.MapPost("/orders", async ( CreateOrderRequest request, IOrderService orderService, CancellationToken ct) => { var result = await orderService.CreateOrderAsync(request, ct); return result.IsSuccess ? Results.Created($"/orders/{result.Value!.Id}", result.Value) : Results.BadRequest(new { error = result.Error, code = result.ErrorCode }); }); ``` ## Data Access Patterns ### Entity Framework Core ```csharp // DbContext configuration public class AppDbContext : DbContext { public DbSet<Product> Products => Set<Product>(); public DbSet<Order> Orders => Set<Order>(); protected override void OnModelCreating(ModelBuilder modelBuilder) { // Apply all configurations from assembly modelBuilder.ApplyConfigurationsFromAssembly(typeof(AppDbContext).Assembly); // Global query filters modelBuilder.Entity<Product>().HasQueryFilter(p => !p.IsDeleted); } } // Entity configuration public class ProductConfiguration : IEntityTypeConfiguration<Product> { public void Configure(EntityTypeBuilder<Product> builder) { builder.ToTable("Products"); builder.HasKey(p => p.Id); builder.Property(p => p.Id).HasMaxLength(40); builder.Property(p => p.Name).HasMaxLength(200).IsRequired(); builder.Property(p => p.Price).HasPrecision(18, 2); builder.HasIndex(p => p.Sku).IsUnique(); builder.HasIndex(p => new { p.CategoryId, p.Name }); builder.HasMany(p => p.OrderItems) .WithOne(oi => oi.Product) .HasForeignKey(oi => oi.ProductId); } } // Repository with EF Core public class ProductRepository : IProductRepository { private readonly AppDbContext _context; public async Task<Product?> GetByIdAsync(string id, CancellationToken ct = default) { return await _context.Products .AsNoTracking() .FirstOrDefaultAsync(p => p.Id == id, ct); } public async Task<IReadOnlyList<Product>> SearchAsync( ProductSearchCriteria criteria, CancellationToken ct = default) { var query = _context.Products.AsNoTracking(); if (!string.IsNullOrWhiteSpace(criteria.SearchTerm)) query = query.Where(p => EF.Functions.Like(p.Name, $"%{criteria.SearchTerm}%")); if (criteria.CategoryId.HasValue) query = query.Where(p => p.CategoryId == criteria.CategoryId); if (criteria.MinPrice.HasValue) query = query.Where(p => p.Price >= criteria.MinPrice); if (criteria.MaxPrice.HasValue) query = query.Where(p => p.Price <= criteria.MaxPrice); return await query .OrderBy(p => p.Name) .Skip((criteria.Page - 1) * criteria.PageSize) .Take(criteria.PageSize) .ToListAsync(ct); } } ``` ### Dapper for Performance ```csharp public class DapperProductRepository : IProductRepository { private readonly IDbConnection _connection; public async Task<Product?> GetByIdAsync(string id, CancellationToken ct = default) { const string sql = """ SELECT Id, Name, Sku, Price, CategoryId, Stock, CreatedAt FROM Products WHERE Id = @Id AND IsDeleted = 0 """; return await _connection.QueryFirstOrDefaultAsync<Product>( new CommandDefinition(sql, new { Id = id }, cancellationToken: ct)); } public async Task<IReadOnlyList<Product>> SearchAsync( ProductSearchCriteria criteria, CancellationToken ct = default) { var sql = new StringBuilder(""" SELECT Id, Name, Sku, Price, CategoryId, Stock, CreatedAt FROM Products WHERE IsDeleted = 0 """); var parameters = new DynamicParameters(); if (!string.IsNullOrWhiteSpace(criteria.SearchTerm)) { sql.Append(" AND Name LIKE @SearchTerm"); parameters.Add("SearchTerm", $"%{criteria.SearchTerm}%"); } if (criteria.CategoryId.HasValue) { sql.Append(" AND CategoryId = @CategoryId"); parameters.Add("CategoryId", criteria.CategoryId); } if (criteria.MinPrice.HasValue) { sql.Append(" AND Price >= @MinPrice"); parameters.Add("MinPrice", criteria.MinPrice); } if (criteria.MaxPrice.HasValue) { sql.Append(" AND Price <= @MaxPrice"); parameters.Add("MaxPrice", criteria.MaxPrice); } sql.Append(" ORDER BY Name OFFSET @Offset ROWS FETCH NEXT @PageSize ROWS ONLY"); parameters.Add("Offset", (criteria.Page - 1) * criteria.PageSize); parameters.Add("PageSize", criteria.PageSize); var results = await _connection.QueryAsync<Product>( new CommandDefinition(sql.ToString(), parameters, cancellationToken: ct)); return results.ToList(); } // Multi-mapping for related data public async Task<Order?> GetOrderWithItemsAsync(int orderId, CancellationToken ct = default) { const string sql = """ SELECT o.*, oi.*, p.* FROM Orders o LEFT JOIN OrderItems oi ON o.Id = oi.OrderId LEFT JOIN Products p ON oi.ProductId = p.Id WHERE o.Id = @OrderId """; var orderDictionary = new Dictionary<int, Order>(); await _connection.QueryAsync<Order, OrderItem, Product, Order>( new CommandDefinition(sql, new { OrderId = orderId }, cancellationToken: ct), (order, item, product) => { if (!orderDictionary.TryGetValue(order.Id, out var existingOrder)) { existingOrder = order; existingOrder.Items = new List<OrderItem>(); orderDictionary.Add(order.Id, existingOrder); } if (item != null) { item.Product = product; existingOrder.Items.Add(item); } return existingOrder; }, splitOn: "Id,Id"); return orderDictionary.Values.FirstOrDefault(); } } ``` ## Caching Patterns ### Multi-Level Cache with Redis ```csharp public class CachedProductService : IProductService { private readonly IProductRepository _repository; private readonly IMemoryCache _memoryCache; private readonly IDistributedCache _distributedCache; private readonly ILogger<CachedProductService> _logger; private static readonly TimeSpan MemoryCacheDuration = TimeSpan.FromMinutes(1); private static readonly TimeSpan DistributedCacheDuration = TimeSpan.FromMinutes(15); public async Task<Product?> GetByIdAsync(string id, CancellationToken ct = default) { var cacheKey = $"product:{id}"; // L1: Memory cache (in-process, fastest) if (_memoryCache.TryGetValue(cacheKey, out Product? cached)) { _logger.LogDebug("L1 cache hit for {CacheKey}", cacheKey); return cached; } // L2: Distributed cache (Redis) var distributed = await _distributedCache.GetStringAsync(cacheKey, ct); if (distributed != null) { _logger.LogDebug("L2 cache hit for {CacheKey}", cacheKey); var product = JsonSerializer.Deserialize<Product>(distributed); // Populate L1 _memoryCache.Set(cacheKey, product, MemoryCacheDuration); return product; } // L3: Database _logger.LogDebug("Cache miss for {CacheKey}, fetching from database", cacheKey); var fromDb = await _repository.GetByIdAsync(id, ct); if (fromDb != null) { var serialized = JsonSerializer.Serialize(fromDb); // Populate both caches await _distributedCache.SetStringAsync( cacheKey, serialized, new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = DistributedCacheDuration }, ct); _memoryCache.Set(cacheKey, fromDb, MemoryCacheDuration); } return fromDb; } public async Task InvalidateAsync(string id, CancellationToken ct = default) { var cacheKey = $"product:{id}"; _memoryCache.Remove(cacheKey); await _distributedCache.RemoveAsync(cacheKey, ct); _logger.LogInformation("Invalidated cache for {CacheKey}", cacheKey); } } // Stale-while-revalidate pattern public class StaleWhileRevalidateCache<T> { private readonly IDistributedCache _cache; private readonly TimeSpan _freshDuration; private readonly TimeSpan _staleDuration; public async Task<T?> GetOrCreateAsync( string key, Func<CancellationToken, Task<T>> factory, CancellationToken ct = default) { var cached = await _cache.GetStringAsync(key, ct); if (cached != null) { var entry = JsonSerializer.Deserialize<CacheEntry<T>>(cached)!; if (entry.IsStale && !entry.IsExpired) { // Return stale data immediately, refresh in background _ = Task.Run(async () => { var fresh = await factory(CancellationToken.None); await SetAsync(key, fresh, CancellationToken.None); }); } if (!entry.IsExpired) return entry.Value; } // Cache miss or expired var value = await factory(ct); await SetAsync(key, value, ct); return value; } private record CacheEntry<TValue>(TValue Value, DateTime CreatedAt) { public bool IsStale => DateTime.UtcNow - CreatedAt > _freshDuration; public bool IsExpired => DateTime.UtcNow - CreatedAt > _staleDuration; } } ``` ## Testing Patterns ### Unit Tests with xUnit and Moq ```csharp public class OrderServiceTests { private readonly Mock<IOrderRepository> _mockRepository; private readonly Mock<IStockService> _mockStockService; private readonly Mock<IValidator<CreateOrderRequest>> _mockValidator; private readonly OrderService _sut; // System Under Test public OrderServiceTests() { _mockRepository = new Mock<IOrderRepository>(); _mockStockService = new Mock<IStockService>(); _mockValidator = new Mock<IValidator<CreateOrderRequest>>(); // Default: validation passes _mockValidator .Setup(v => v.ValidateAsync(It.IsAny<CreateOrderRequest>(), It.IsAny<
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

on-call-handoff-patterns

Master on-call shift handoffs with context transfer, escalation

coding
โญ1
# On-Call Handoff Patterns Effective patterns for on-call shift transitions, ensuring continuity, context transfer, and reliable incident response across shifts. ## When to Use This Skill - Transitioning on-call responsibilities - Writing shift handoff summaries - Documenting ongoing investigations - Establishing on-call rotation procedures - Improving handoff quality - Onboarding new on-call engineers ## Core Concepts ### 1. Handoff Components | Component | Purpose | | -------------------------- | ----------------------- | | **Active Incidents** | What's currently broken | | **Ongoing Investigations** | Issues being debugged | | **Recent Changes** | Deployments, configs | | **Known Issues** | Workarounds in place | | **Upcoming Events** | Maintenance, releases | ### 2. Handoff Timing ``` Recommended: 30 min overlap between shifts Outgoing: โ”œโ”€โ”€ 15 min: Write handoff document โ””โ”€โ”€ 15 min: Sync call with incoming Incoming: โ”œโ”€โ”€ 15 min: Review handoff document โ”œโ”€โ”€ 15 min: Sync call with outgoing โ””โ”€โ”€ 5 min: Verify alerting setup ``` ## Templates ### Template 1: Shift Handoff Document ````markdown # On-Call Handoff: Platform Team **Outgoing**: @alice (2024-01-15 to 2024-01-22) **Incoming**: @bob (2024-01-22 to 2024-01-29) **Handoff Time**: 2024-01-22 09:00 UTC --- ## ๐Ÿ”ด Active Incidents ### None currently active No active incidents at handoff time. --- ## ๐ŸŸก Ongoing Investigations ### 1. Intermittent API Timeouts (ENG-1234) **Status**: Investigating **Started**: 2024-01-20 **Impact**: ~0.1% of requests timing out **Context**: - Timeouts correlate with database backup window (02:00-03:00 UTC) - Suspect backup process causing lock contention - Added extra logging in PR #567 (deployed 01/21) **Next Steps**: - [ ] Review new logs after tonight's backup - [ ] Consider moving backup window if confirmed **Resources**: - Dashboard: [API Latency](https://grafana/d/api-latency) - Thread: #platform-eng (01/20, 14:32) --- ### 2. Memory Growth in Auth Service (ENG-1235) **Status**: Monitoring **Started**: 2024-01-18 **Impact**: None yet (proactive) **Context**: - Memory usage growing ~5% per day - No memory leak found in profiling - Suspect connection pool not releasing properly **Next Steps**: - [ ] Review heap dump from 01/21 - [ ] Consider restart if usage > 80% **Resources**: - Dashboard: [Auth Service Memory](https://grafana/d/auth-memory) - Analysis doc: [Memory Investigation](https://docs/eng-1235) --- ## ๐ŸŸข Resolved This Shift ### Payment Service Outage (2024-01-19) - **Duration**: 23 minutes - **Root Cause**: Database connection exhaustion - **Resolution**: Rolled back v2.3.4, increased pool size - **Postmortem**: [POSTMORTEM-89](https://docs/postmortem-89) - **Follow-up tickets**: ENG-1230, ENG-1231 --- ## ๐Ÿ“‹ Recent Changes ### Deployments | Service | Version | Time | Notes | | ------------ | ------- | ----------- | -------------------------- | | api-gateway | v3.2.1 | 01/21 14:00 | Bug fix for header parsing | | user-service | v2.8.0 | 01/20 10:00 | New profile features | | auth-service | v4.1.2 | 01/19 16:00 | Security patch | ### Configuration Changes - 01/21: Increased API rate limit from 1000 to 1500 RPS - 01/20: Updated database connection pool max from 50 to 75 ### Infrastructure - 01/20: Added 2 nodes to Kubernetes cluster - 01/19: Upgraded Redis from 6.2 to 7.0 --- ## โš ๏ธ Known Issues & Workarounds ### 1. Slow Dashboard Loading **Issue**: Grafana dashboards slow on Monday mornings **Workaround**: Wait 5 min after 08:00 UTC for cache warm-up **Ticket**: OPS-456 (P3) ### 2. Flaky Integration Test **Issue**: `test_payment_flow` fails intermittently in CI **Workaround**: Re-run failed job (usually passes on retry) **Ticket**: ENG-1200 (P2) --- ## ๐Ÿ“… Upcoming Events | Date | Event | Impact | Contact | | ----------- | -------------------- | ------------------- | ------------- | | 01/23 02:00 | Database maintenance | 5 min read-only | @dba-team | | 01/24 14:00 | Major release v5.0 | Monitor closely | @release-team | | 01/25 | Marketing campaign | 2x traffic expected | @platform | --- ## ๐Ÿ“ž Escalation Reminders | Issue Type | First Escalation | Second Escalation | | --------------- | -------------------- | ----------------- | | Payment issues | @payments-oncall | @payments-manager | | Auth issues | @auth-oncall | @security-team | | Database issues | @dba-team | @infra-manager | | Unknown/severe | @engineering-manager | @vp-engineering | --- ## ๐Ÿ”ง Quick Reference ### Common Commands ```bash # Check service health kubectl get pods -A | grep -v Running # Recent deployments kubectl get events --sort-by='.lastTimestamp' | tail -20 # Database connections psql -c "SELECT count(*) FROM pg_stat_activity;" # Clear cache (emergency only) redis-cli FLUSHDB ``` ```` ### Important Links - [Runbooks](https://wiki/runbooks) - [Service Catalog](https://wiki/services) - [Incident Slack](https://slack.com/incidents) - [PagerDuty](https://pagerduty.com/schedules) --- ## Handoff Checklist ### Outgoing Engineer - [x] Document active incidents - [x] Document ongoing investigations - [x] List recent changes - [x] Note known issues - [x] Add upcoming events - [x] Sync with incoming engineer ### Incoming Engineer - [ ] Read this document - [ ] Join sync call - [ ] Verify PagerDuty is routing to you - [ ] Verify Slack notifications working - [ ] Check VPN/access working - [ ] Review critical dashboards ```` ### Template 2: Quick Handoff (Async) ```markdown # Quick Handoff: @alice โ†’ @bob ## TL;DR - No active incidents - 1 investigation ongoing (API timeouts, see ENG-1234) - Major release tomorrow (01/24) - be ready for issues ## Watch List 1. API latency around 02:00-03:00 UTC (backup window) 2. Auth service memory (restart if > 80%) ## Recent - Deployed api-gateway v3.2.1 yesterday (stable) - Increased rate limits to 1500 RPS ## Coming Up - 01/23 02:00 - DB maintenance (5 min read-only) - 01/24 14:00 - v5.0 release ## Questions? I'll be available on Slack until 17:00 today. ```` ### Template 3: Incident Handoff (Mid-Incident) ```markdown # INCIDENT HANDOFF: Payment Service Degradation **Incident Start**: 2024-01-22 08:15 UTC **Current Status**: Mitigating **Severity**: SEV2 --- ## Current State - Error rate: 15% (down from 40%) - Mitigation in progress: scaling up pods - ETA to resolution: ~30 min ## What We Know 1. Root cause: Memory pressure on payment-service pods 2. Triggered by: Unusual traffic spike (3x normal) 3. Contributing: Inefficient query in checkout flow ## What We've Done - Scaled payment-service from 5 โ†’ 15 pods - Enabled rate limiting on checkout endpoint - Disabled non-critical features ## What Needs to Happen 1. Monitor error rate - should reach <1% in ~15 min 2. If not improving, escalate to @payments-manager 3. Once stable, begin root cause investigation ## Key People - Incident Commander: @alice (handing off) - Comms Lead: @charlie - Technical Lead: @bob (incoming) ## Communication - Status page: Updated at 08:45 - Customer support: Notified - Exec team: Aware ## Resources - Incident channel: #inc-20240122-payment - Dashboard: [Payment Service](https://grafana/d/payments) - Runbook: [Payment Degradation](https://wiki/runbooks/payments) --- **Incoming on-call (@bob) - Please confirm you have:** - [ ] Joined #inc-20240122-payment - [ ] Access to dashboards - [ ] Understand current state - [ ] Know escalation path ``` ## Handoff Sync Meeting ### Agenda (15 minutes) ```markdown ## Handoff Sync: @alice โ†’ @bob 1. **Active Issues** (5 min) - Walk through any ongoing incidents - Discuss investigation status - Transfer context and theories 2. **Recent Changes** (3 min) - Deployments to watch - Config changes - Known regressions 3. **Upcoming Events** (3 min) - Maintenance windows - Expected traffic changes - Releases planned 4. **Questions** (4 min) - Clarify anything unclear - Confirm access and alerting - Exchange contact info ``` ## On-Call Best Practices ### Before Your Shift ```markdown ## Pre-Shift Checklist ### Access Verification - [ ] VPN working - [ ] kubectl access to all clusters - [ ] Database read access - [ ] Log aggregator access (Splunk/Datadog) - [ ] PagerDuty app installed and logged in ### Alerting Setup - [ ] PagerDuty schedule shows you as primary - [ ] Phone notifications enabled - [ ] Slack notifications for incident channels - [ ] Test alert received and acknowledged ### Knowledge Refresh - [ ] Review recent incidents (past 2 weeks) - [ ] Check service changelog - [ ] Skim critical runbooks - [ ] Know escalation contacts ### Environment Ready - [ ] Laptop charged and accessible - [ ] Phone charged - [ ] Quiet space available for calls - [ ] Secondary contact identified (if traveling) ``` ### During Your Shift ```markdown ## Daily On-Call Routine ### Morning (start of day) - [ ] Check overnight alerts - [ ] Review dashboards for anomalies - [ ] Check for any P0/P1 tickets created - [ ] Skim incident channels for context ### Throughout Day - [ ] Respond to alerts within SLA - [ ] Document investigation progress - [ ] Update team on significant issues - [ ] Triage incoming pages ### End of Day - [ ] Hand off any active issues - [ ] Update investigation docs - [ ] Note anything for next shift ``` ### After Your Shift ```markdown ## Post-Shift Checklist - [ ] Complete handoff document - [ ] Sync with incoming on-call - [ ] Verify PagerDuty routing changed - [ ] Close/update investigation tickets - [ ] File postmortems for any incidents - [ ] Take time off if shift was stressful ``` ## Escalation Guidelines ### When to Escalate ```markdown ## Escalation Triggers ### Immediate Escalation - SEV1 incident declared - Data breach suspected - Unable to diagnose within 30 min - Customer or legal escalation received ### Consider Escalation - Issue spans multiple teams - Requires expertise you don't have - Business impact exceeds threshold - You're uncertain about next steps ### How to Escalate 1. Page the appropriate escalation path 2. Provide brief context in Slack 3. Stay engaged until escalation acknowledges 4. Hand off cleanly, don't just disappear ``` ## Best Practices ### Do's - **Document everything** - Future you will thank you - **Escalate early** - Better safe than sorry - **Take breaks** - Alert fatigue is real - **Keep handoffs synchronous** - Async loses context - **Test your setup** - Before incidents, not during ### Don'ts - **Don't skip handoffs** - Context loss causes incidents - **Don't hero** - Escalate when needed - **Don't ignore alerts** - Even if they seem minor - **Don't work sick** - Swap shifts instead - **Don't disappear** - Stay reachable during shift ## Resources - [Google SRE - Being On-Call](https://sre.google/sre-book/being-on-call/) - [PagerDuty On-Call Guide](https://www.pagerduty.com/resources/learn/on-call-management/) - [Increment On-Call Issue](https://increment.com/on-call/)
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

embedding-strategies

Select and optimize embedding models for semantic search and RAG

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

hybrid-search-implementation

Combine vector and keyword search for improved retrieval. Use when

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

team-composition-patterns

Design optimal agent team compositions with sizing heuristics,

coding
โญ1
# Team Composition Patterns Best practices for composing multi-agent teams, selecting team sizes, choosing agent types, and configuring display modes for Claude Code's Agent Teams feature. ## When to Use This Skill - Deciding how many teammates to spawn for a task - Choosing between preset team configurations - Selecting the right agent type (subagent_type) for each role - Configuring teammate display modes (tmux, iTerm2, in-process) - Building custom team compositions for non-standard workflows ## Team Sizing Heuristics | Complexity | Team Size | When to Use | | ------------ | --------- | ----------------------------------------------------------- | | Simple | 1-2 | Single-dimension review, isolated bug, small feature | | Moderate | 2-3 | Multi-file changes, 2-3 concerns, medium features | | Complex | 3-4 | Cross-cutting concerns, large features, deep debugging | | Very Complex | 4-5 | Full-stack features, comprehensive reviews, systemic issues | **Rule of thumb**: Start with the smallest team that covers all required dimensions. Adding teammates increases coordination overhead. ## Preset Team Compositions ### Review Team - **Size**: 3 reviewers - **Agents**: 3x `team-reviewer` - **Default dimensions**: security, performance, architecture - **Use when**: Code changes need multi-dimensional quality assessment ### Debug Team - **Size**: 3 investigators - **Agents**: 3x `team-debugger` - **Default hypotheses**: 3 competing hypotheses - **Use when**: Bug has multiple plausible root causes ### Feature Team - **Size**: 3 (1 lead + 2 implementers) - **Agents**: 1x `team-lead` + 2x `team-implementer` - **Use when**: Feature can be decomposed into parallel work streams ### Fullstack Team - **Size**: 4 (1 lead + 3 implementers) - **Agents**: 1x `team-lead` + 1x frontend `team-implementer` + 1x backend `team-implementer` + 1x test `team-implementer` - **Use when**: Feature spans frontend, backend, and test layers ### Research Team - **Size**: 3 researchers - **Agents**: 3x `general-purpose` - **Default areas**: Each assigned a different research question, module, or topic - **Capabilities**: Codebase search (Grep, Glob, Read), web search (WebSearch, WebFetch) - **Use when**: Need to understand a codebase, research libraries, compare approaches, or gather information from code and web sources in parallel ### Security Team - **Size**: 4 reviewers - **Agents**: 4x `team-reviewer` - **Default dimensions**: OWASP/vulnerabilities, auth/access control, dependencies/supply chain, secrets/configuration - **Use when**: Comprehensive security audit covering multiple attack surfaces ### Migration Team - **Size**: 4 (1 lead + 2 implementers + 1 reviewer) - **Agents**: 1x `team-lead` + 2x `team-implementer` + 1x `team-reviewer` - **Use when**: Large codebase migration (framework upgrade, language port, API version bump) requiring parallel work with correctness verification ## Agent Type Selection When spawning teammates with the Task tool, choose `subagent_type` based on what tools the teammate needs: | Agent Type | Tools Available | Use For | | ------------------------------ | ----------------------------------------- | ---------------------------------------------------------- | | `general-purpose` | All tools (Read, Write, Edit, Bash, etc.) | Implementation, debugging, any task requiring file changes | | `Explore` | Read-only tools (Read, Grep, Glob) | Research, code exploration, analysis | | `Plan` | Read-only tools | Architecture planning, task decomposition | | `agent-teams:team-reviewer` | All tools | Code review with structured findings | | `agent-teams:team-debugger` | All tools | Hypothesis-driven investigation | | `agent-teams:team-implementer` | All tools | Building features within file ownership boundaries | | `agent-teams:team-lead` | All tools | Team orchestration and coordination | **Key distinction**: Read-only agents (Explore, Plan) cannot modify files. Never assign implementation tasks to read-only agents. ## Display Mode Configuration Configure in `~/.claude/settings.json`: ```json { "teammateMode": "tmux" } ``` | Mode | Behavior | Best For | | -------------- | ------------------------------ | ------------------------------------------------- | | `"tmux"` | Each teammate in a tmux pane | Development workflows, monitoring multiple agents | | `"iterm2"` | Each teammate in an iTerm2 tab | macOS users who prefer iTerm2 | | `"in-process"` | All teammates in same process | Simple tasks, CI/CD environments | ## Custom Team Guidelines When building custom teams: 1. **Every team needs a coordinator** โ€” Either designate a `team-lead` or have the user coordinate directly 2. **Match roles to agent types** โ€” Use specialized agents (reviewer, debugger, implementer) when available 3. **Avoid duplicate roles** โ€” Two agents doing the same thing wastes resources 4. **Define boundaries upfront** โ€” Each teammate needs clear ownership of files or responsibilities 5. **Keep it small** โ€” 2-4 teammates is the sweet spot; 5+ requires significant coordination overhead
๐Ÿ‘0
๐Ÿ‘๏ธ0
๐Ÿค– Auto-discovered
๐Ÿค–system promptโ€ข7 months ago

prompt-engineering-patterns

Master advanced prompt engineering techniques to maximize LLM

coding
โญ1
# Prompt Engineering Patterns Master advanced prompt engineering techniques to maximize LLM performance, reliability, and controllability. ## When to Use This Skill - Designing complex prompts for production LLM applications - Optimizing prompt performance and consistency - Implementing structured reasoning patterns (chain-of-thought, tree-of-thought) - Building few-shot learning systems with dynamic example selection - Creating reusable prompt templates with variable interpolation - Debugging and refining prompts that produce inconsistent outputs - Implementing system prompts for specialized AI assistants - Using structured outputs (JSON mode) for reliable parsing ## Core Capabilities ### 1. Few-Shot Learning - Example selection strategies (semantic similarity, diversity sampling) - Balancing example count with context window constraints - Constructing effective demonstrations with input-output pairs - Dynamic example retrieval from knowledge bases - Handling edge cases through strategic example selection ### 2. Chain-of-Thought Prompting - Step-by-step reasoning elicitation - Zero-shot CoT with "Let's think step by step" - Few-shot CoT with reasoning traces - Self-consistency techniques (sampling multiple reasoning paths) - Verification and validation steps ### 3. Structured Outputs - JSON mode for reliable parsing - Pydantic schema enforcement - Type-safe response handling - Error handling for malformed outputs ### 4. Prompt Optimization - Iterative refinement workflows - A/B testing prompt variations - Measuring prompt performance metrics (accuracy, consistency, latency) - Reducing token usage while maintaining quality - Handling edge cases and failure modes ### 5. Template Systems - Variable interpolation and formatting - Conditional prompt sections - Multi-turn conversation templates - Role-based prompt composition - Modular prompt components ### 6. System Prompt Design - Setting model behavior and constraints - Defining output formats and structure - Establishing role and expertise - Safety guidelines and content policies - Context setting and background information ## Quick Start ```python from langchain_anthropic import ChatAnthropic from langchain_core.prompts import ChatPromptTemplate from pydantic import BaseModel, Field # Define structured output schema class SQLQuery(BaseModel): query: str = Field(description="The SQL query") explanation: str = Field(description="Brief explanation of what the query does") tables_used: list[str] = Field(description="List of tables referenced") # Initialize model with structured output llm = ChatAnthropic(model="claude-sonnet-4-6") structured_llm = llm.with_structured_output(SQLQuery) # Create prompt template prompt = ChatPromptTemplate.from_messages([ ("system", """You are an expert SQL developer. Generate efficient, secure SQL queries. Always use parameterized queries to prevent SQL injection. Explain your reasoning briefly."""), ("user", "Convert this to SQL: {query}") ]) # Create chain chain = prompt | structured_llm # Use result = await chain.ainvoke({ "query": "Find all users who registered in the last 30 days" }) print(result.query) print(result.explanation) ``` ## Key Patterns ### Pattern 1: Structured Output with Pydantic ```python from anthropic import Anthropic from pydantic import BaseModel, Field from typing import Literal import json class SentimentAnalysis(BaseModel): sentiment: Literal["positive", "negative", "neutral"] confidence: float = Field(ge=0, le=1) key_phrases: list[str] reasoning: str async def analyze_sentiment(text: str) -> SentimentAnalysis: """Analyze sentiment with structured output.""" client = Anthropic() message = client.messages.create( model="claude-sonnet-4-6", max_tokens=500, messages=[{ "role": "user", "content": f"""Analyze the sentiment of this text. Text: {text} Respond with JSON matching this schema: {{ "sentiment": "positive" | "negative" | "neutral", "confidence": 0.0-1.0, "key_phrases": ["phrase1", "phrase2"], "reasoning": "brief explanation" }}""" }] ) return SentimentAnalysis(**json.loads(message.content[0].text)) ``` ### Pattern 2: Chain-of-Thought with Self-Verification ```python from langchain_core.prompts import ChatPromptTemplate cot_prompt = ChatPromptTemplate.from_template(""" Solve this problem step by step. Problem: {problem} Instructions: 1. Break down the problem into clear steps 2. Work through each step showing your reasoning 3. State your final answer 4. Verify your answer by checking it against the original problem Format your response as: ## Steps [Your step-by-step reasoning] ## Answer [Your final answer] ## Verification [Check that your answer is correct] """) ``` ### Pattern 3: Few-Shot with Dynamic Example Selection ```python from langchain_voyageai import VoyageAIEmbeddings from langchain_core.example_selectors import SemanticSimilarityExampleSelector from langchain_chroma import Chroma # Create example selector with semantic similarity example_selector = SemanticSimilarityExampleSelector.from_examples( examples=[ {"input": "How do I reset my password?", "output": "Go to Settings > Security > Reset Password"}, {"input": "Where can I see my order history?", "output": "Navigate to Account > Orders"}, {"input": "How do I contact support?", "output": "Click Help > Contact Us or email support@example.com"}, ], embeddings=VoyageAIEmbeddings(model="voyage-3-large"), vectorstore_cls=Chroma, k=2 # Select 2 most similar examples ) async def get_few_shot_prompt(query: str) -> str: """Build prompt with dynamically selected examples.""" examples = await example_selector.aselect_examples({"input": query}) examples_text = "\n".join( f"User: {ex['input']}\nAssistant: {ex['output']}" for ex in examples ) return f"""You are a helpful customer support assistant. Here are some example interactions: {examples_text} Now respond to this query: User: {query} Assistant:""" ``` ### Pattern 4: Progressive Disclosure Start with simple prompts, add complexity only when needed: ```python PROMPT_LEVELS = { # Level 1: Direct instruction "simple": "Summarize this article: {text}", # Level 2: Add constraints "constrained": """Summarize this article in 3 bullet points, focusing on: - Key findings - Main conclusions - Practical implications Article: {text}""", # Level 3: Add reasoning "reasoning": """Read this article carefully. 1. First, identify the main topic and thesis 2. Then, extract the key supporting points 3. Finally, summarize in 3 bullet points Article: {text} Summary:""", # Level 4: Add examples "few_shot": """Read articles and provide concise summaries. Example: Article: "New research shows that regular exercise can reduce anxiety by up to 40%..." Summary: โ€ข Regular exercise reduces anxiety by up to 40% โ€ข 30 minutes of moderate activity 3x/week is sufficient โ€ข Benefits appear within 2 weeks of starting Now summarize this article: Article: {text} Summary:""" } ``` ### Pattern 5: Error Recovery and Fallback ```python from pydantic import BaseModel, ValidationError import json class ResponseWithConfidence(BaseModel): answer: str confidence: float sources: list[str] alternative_interpretations: list[str] = [] ERROR_RECOVERY_PROMPT = """ Answer the question based on the context provided. Context: {context} Question: {question} Instructions: 1. If you can answer confidently (>0.8), provide a direct answer 2. If you're somewhat confident (0.5-0.8), provide your best answer with caveats 3. If you're uncertain (<0.5), explain what information is missing 4. Always provide alternative interpretations if the question is ambiguous Respond in JSON: {{ "answer": "your answer or 'I cannot determine this from the context'", "confidence": 0.0-1.0, "sources": ["relevant context excerpts"], "alternative_interpretations": ["if question is ambiguous"] }} """ async def answer_with_fallback( context: str, question: str, llm ) -> ResponseWithConfidence: """Answer with error recovery and fallback.""" prompt = ERROR_RECOVERY_PROMPT.format(context=context, question=question) try: response = await llm.ainvoke(prompt) return ResponseWithConfidence(**json.loads(response.content)) except (json.JSONDecodeError, ValidationError) as e: # Fallback: try to extract answer without structure simple_prompt = f"Based on: {context}\n\nAnswer: {question}" simple_response = await llm.ainvoke(simple_prompt) return ResponseWithConfidence( answer=simple_response.content, confidence=0.5, sources=["fallback extraction"], alternative_interpretations=[] ) ``` ### Pattern 6: Role-Based System Prompts ```python SYSTEM_PROMPTS = { "analyst": """You are a senior data analyst with expertise in SQL, Python, and business intelligence. Your responsibilities: - Write efficient, well-documented queries - Explain your analysis methodology - Highlight key insights and recommendations - Flag any data quality concerns Communication style: - Be precise and technical when discussing methodology - Translate technical findings into business impact - Use clear visualizations when helpful""", "assistant": """You are a helpful AI assistant focused on accuracy and clarity. Core principles: - Always cite sources when making factual claims - Acknowledge uncertainty rather than guessing - Ask clarifying questions when the request is ambiguous - Provide step-by-step explanations for complex topics Constraints: - Do not provide medical, legal, or financial advice - Redirect harmful requests appropriately - Protect user privacy""", "code_reviewer": """You are a senior software engineer conducting code reviews. Review criteria: - Correctness: Does the code work as intended? - Security: Are there any vulnerabilities? - Performance: Are there efficiency concerns? - Maintainability: Is the code readable and well-structured? - Best practices: Does it follow language idioms? Output format: 1. Summary assessment (approve/request changes) 2. Critical issues (must fix) 3. Suggestions (nice to have) 4. Positive feedback (what's done well)""" } ``` ## Integration Patterns ### With RAG Systems ```python RAG_PROMPT = """You are a knowledgeable assistant that answers questions based on provided context. Context (retrieved from knowledge base): {context} Instructions: 1. Answer ONLY based on the provided context 2. If the context doesn't contain the answer, say "I don't have information about that in my knowledge base" 3. Cite specific passages using [1], [2] notation 4. If the question is ambiguous, ask for clarification Question: {question} Answer:""" ``` ### With Validation and Verification ```python VALIDATED_PROMPT = """Complete the following task: Task: {task} After generating your response, verify it meets ALL these criteria: โœ“ Directly addresses the original request โœ“ Contains no factual errors โœ“ Is appropriately detailed (not too brief, not too verbose) โœ“ Uses proper formatting โœ“ Is safe and appropriate If verification fails on any criterion, revise before responding. Response:""" ``` ## Performance Optimization ### Token Efficiency ```python # Before: Verbose prompt (150+ tokens) verbose_prompt = """ I would like you to please take the following text and provide me with a comprehensive summary of the main points. The summary should capture the key ideas and important details while being concise and easy to understand. """ # After: Concise prompt (30 tokens) concise_prompt = """Summarize the key points concisely: {text} Summary:""" ``` ### Caching Common Prefixes ```python from anthropic import Anthropic client = Anthropic() # Use prompt caching for repeated system prompts response = client.messages.create( model="claude-sonnet-4-6", max_tokens=1000, system=[ { "type": "text", "text": LONG_SYSTEM_PROMPT, "cache_control": {"type": "ephemeral"} } ], messages=[{"role": "user", "content": user_query}] ) ``` ## Best Practices 1. **Be Specific**: Vague prompts produce inconsistent results 2. **Show, Don't Tell**: Examples are more effective than descriptions 3. **Use Structured Outputs**: Enforce schemas with Pydantic for reliability 4. **Test Extensively**: Evaluate on diverse, representative inputs 5. **Iterate Rapidly**: Small changes can have large impacts 6. **Monitor Performance**: Track metrics in production 7. **Version Control**: Treat prompts as code with proper versioning 8. **Document Intent**: Explain why prompts are structured as they are ## Common Pitfalls - **Over-engineering**: Starting with complex prompts before trying simple ones - **Example pollution**: Using examples that don't match the target task - **Context overflow**: Exceeding token limits with excessive examples - **Ambiguous instructions**: Leaving room for multiple interpretations - **Ignoring edge cases**: Not testing on unusual or boundary inputs - **No error handling**: Assuming outputs will always be well-formed - **Hardcoded values**: Not parameterizing prompts for reuse ## Success Metrics Track these KPIs for your prompts: - **Accuracy**: Correctness of outputs - **Consistency**: Reproducibility across similar inputs - **Latency**: Response time (P50, P95, P99) - **Token Usage**: Average tokens per request - **Success Rate**: Percentage of valid, parseable outputs - **User Satisfaction**: Ratings and feedback ## Resources - [Anthropic Prompt Engineering Guide](https://docs.anthropic.com/en/docs/build-with-claude/prompt-engineering) - [Claude Prompt Caching](https://docs.anthropic.com/en/docs/build-with-claude/prompt-caching) - [OpenAI Prompt Engineering](https://platform.openai.com/docs/guides/prompt-engineering) - [LangChain Prompts](https://python.langchain.com/docs/concepts/prompts/)
๐Ÿ‘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

ml-pipeline-workflow

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

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

sast-configuration

Configure Static Application Security Testing (SAST) tools for

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