Senior Data Engineer
Core Capabilities
- Batch Pipeline Orchestration - Design and implement production-ready ETL/ELT pipelines with Airflow, intelligent dependency resolution, retry logic, and comprehensive monitoring
- Real-Time Streaming - Build event-driven streaming pipelines with Kafka, Flink, Kinesis, and Spark Streaming with exactly-once semantics and sub-second latency
- Data Quality Management - Comprehensive batch and streaming data quality validation covering completeness, accuracy, consistency, timeliness, and validity
- Streaming Quality Monitoring - Track consumer lag, data freshness, schema drift, throughput, and dead letter queue rates for streaming pipelines
- Performance Optimization - Analyze and optimize pipeline performance with query optimization, Spark tuning, and cost analysis recommendations
Key Workflows
Workflow 1: Build ETL Pipeline
Time: 2-4 hours
Steps:
- Design pipeline architecture using Lambda, Kappa, or Medallion pattern
- Configure YAML pipeline definition with sources, transformations, targets
- Generate Airflow DAG with
pipeline_orchestrator.py
- Define data quality validation rules
- Deploy and configure monitoring/alerting
Expected Output: Production-ready ETL pipeline with 99%+ success rate, automated quality checks, and comprehensive monitoring
Workflow 2: Build Real-Time Streaming Pipeline
Time: 3-5 days
Steps:
- Select streaming architecture (Kappa vs Lambda) based on requirements
- Configure streaming pipeline YAML (sources, processing, sinks, quality)
- Generate Kafka configurations with
kafka_config_generator.py
- Generate Flink/Spark job scaffolding with
stream_processor.py
- Deploy and monitor with
streaming_quality_validator.py
Expected Output: Streaming pipeline processing 10K+ events/sec with P99 latency < 1s, exactly-once delivery, and real-time quality monitoring
World-class data engineering for production-grade data systems, scalable pipelines, and enterprise data platforms.
Overview
This skill provides comprehensive expertise in data engineering fundamentals through advanced production patterns. From designing medallion architectures to implementing real-time streaming pipelines, it covers the full spectrum of modern data engineering including ETL/ELT design, data quality frameworks, pipeline orchestration, and DataOps practices.
What This Skill Provides:
- Production-ready pipeline templates (Airflow, Spark, dbt)
- Comprehensive data quality validation framework
- Performance optimization and cost analysis tools
- Data architecture patterns (Lambda, Kappa, Medallion)
- Complete DataOps CI/CD workflows
Best For:
- Building scalable data pipelines for enterprise systems
- Implementing data quality and governance frameworks
- Optimizing ETL performance and cloud costs
- Designing modern data architectures (lake, warehouse, lakehouse)
- Production ML/AI data infrastructure
Quick Start
Pipeline Orchestration
# Generate Airflow DAG from configuration
python scripts/pipeline_orchestrator.py --config pipeline_config.yaml --output dags/
# Validate pipeline configuration
python scripts/pipeline_orchestrator.py --config pipeline_config.yaml --validate
# Use incremental load template
python scripts/pipeline_orchestrator.py --template incremental --output dags/
Data Quality Validation
# Validate CSV file with quality checks
python scripts/data_quality_validator.py --input data/sales.csv --output report.html
# Validate database table with custom rules
python scripts/data_quality_validator.py \
--connection postgresql://user:pass@host/db \
--table sales_transactions \
--rules rules/sales_validation.yaml \
--threshold 0.95
Performance Optimization
# Analyze pipeline performance and get recommendations
python scripts/etl_performance_optimizer.py \
--airflow-db postgresql://host/airflow \
--dag-id sales_etl_pipeline \
--days 30 \
--optimize
# Analyze Spark job performance
python scripts/etl_performance_optimizer.py \
--spark-history-server http://spark-history:18080 \
--app-id app-20250115-001
Real-Time Streaming
# Validate streaming pipeline configuration
python scripts/stream_processor.py --config streaming_config.yaml --validate
# Generate Kafka topic and client configurations
python scripts/kafka_config_generator.py \
--topic user-events \
--partitions 12 \
--replication 3 \
--output kafka/topics/
# Generate exactly-once producer configuration
python scripts/kafka_config_generator.py \
--producer \
--profile exactly-once \
--output kafka/producer.properties
# Generate Flink job scaffolding
python scripts/stream_processor.py \
--config streaming_config.yaml \
--mode flink \
--generate \
--output flink-jobs/
# Monitor streaming quality
python scripts/streaming_quality_validator.py \
--lag --consumer-group events-processor --threshold 10000 \
--freshness --topic processed-events --max-latency-ms 5000 \
--output streaming-health-report.html
Core Workflows
1. Building Production Data Pipelines
Steps:
- Design Architecture: Choose pattern (Lambda, Kappa, Medallion) based on requirements
- Configure Pipeline: Create YAML configuration with sources, transformations, targets
- Generate DAG:
python scripts/pipeline_orchestrator.py --config config.yaml
- Add Quality Checks: Define validation rules for data quality
- Deploy & Monitor: Deploy to Airflow, configure alerts, track metrics
Pipeline Patterns: See frameworks.md for Lambda Architecture, Kappa Architecture, Medallion Architecture (Bronze/Silver/Gold), and Microservices Data patterns.
Templates: See templates.md for complete Airflow DAG templates, Spark job templates, dbt models, and Docker configurations.
2. Data Quality Management
Steps:
- Define Rules: Create validation rules covering completeness, accuracy, consistency
- Run Validation:
python scripts/data_quality_validator.py --rules rules.yaml
- Review Results: Analyze quality scores and failed checks
- Integrate CI/CD: Add validation to pipeline deployment process
- Monitor Trends: Track quality scores over time
Quality Framework: See frameworks.md for complete Data Quality Framework covering all dimensions (completeness, accuracy, consistency, timeliness, validity).
Validation Templates: See templates.md for validation configuration examples and Python API usage.
3. Data Modeling & Transformation
Steps:
- Choose Modeling Approach: Dimensional (Kimball), Data Vault 2.0, or One Big Table
- Design Schema: Define fact tables, dimensions, and relationships
- Implement with dbt: Create staging, intermediate, and mart models
- Handle SCD: Implement slowly changing dimension logic (Type 1/2/3)
- Test & Deploy: Run dbt tests, generate documentation, deploy
Modeling Patterns: See frameworks.md for Dimensional Modeling (Kimball), Data Vault 2.0, One Big Table (OBT), and SCD implementations.
dbt Templates: See templates.md for complete dbt model templates including staging, intermediate, fact tables, and SCD Type 2 logic.
4. Performance Optimization
Steps:
- Profile Pipeline: Run performance analyzer on recent pipeline executions
- Identify Bottlenecks: Review execution time breakdown and slow tasks
- Apply Optimizations: Implement recommendations (partitioning, indexing, batching)
- Tune Spark Jobs: Optimize memory, parallelism, and shuffle settings
- Measure Impact: Compare before/after metrics, track cost savings
Optimization Strategies: See frameworks.md for performance best practices including partitioning strategies, query optimization, and Spark tuning.
Analysis Tools: See tools.md for complete documentation on etl_performance_optimizer.py with query analysis and Spark tuning.
5. Building Real-Time Streaming Pipelines
Steps:
- Architecture Selection: Choose Kappa (streaming-only) or Lambda (batch + streaming) architecture
- Configure Pipeline: Create YAML config with sources, processing engine, sinks, quality thresholds
- Generate Kafka Configs:
python scripts/kafka_config_generator.py --topic events --partitions 12
- Generate Job Scaffolding:
python scripts/stream_processor.py --mode flink --generate
- Deploy Infrastructure: Use Docker Compose for local dev, Kubernetes for production
- Monitor Quality:
python scripts/streaming_quality_validator.py --lag --freshness --throughput
Streaming Patterns: See frameworks.md for stateful processing, stream joins, windowing, exactly-once semantics, and CDC patterns.
Templates: See templates.md for Flink DataStream jobs, Kafka Streams applications, PyFlink templates, and Docker Compose configurations.
Python Tools
pipeline_orchestrator.py
Automated Airflow DAG generation with intelligent dependency resolution and monitoring.
Key Features:
- Generate production-ready DAGs from YAML configuration
- Automatic task dependency resolution
- Built-in retry logic and error handling
- Multi-source support (PostgreSQL, S3, BigQuery, Snowflake)
- Integrated quality checks and alerting
Usage:
# Basic DAG generation
python scripts/pipeline_orchestrator.py --config pipeline_config.yaml --output dags/
# With validation
python scripts/pipeline_orchestrator.py --config config.yaml --validate
# From template
python scripts/pipeline_orchestrator.py --template incremental --output dags/
Complete Documentation: See tools.md for full configuration options, templates, and integration examples.
data_quality_validator.py
Comprehensive data quality validation framework with automated checks and reporting.
Capabilities:
- Multi-dimensional validation (completeness, accuracy, consistency, timeliness, validity)
- Great Expectations integration
- Custom business rule validation
- HTML/PDF report generation
- Anomaly detection
- Historical trend tracking
Usage:
# Validate with custom rules
python scripts/data_quality_validator.py \
--input data/sales.csv \
--rules rules/sales_validation.yaml \
--output report.html
# Database table validation
python scripts/data_quality_validator.py \
--connection postgresql://host/db \
--table sales_transactions \
--threshold 0.95
Complete Documentation: See tools.md for rule configuration, API usage, and integration patterns.
etl_performance_optimizer.py
Pipeline performance analysis with actionable optimization recommendations.
Capabilities:
- Airflow DAG execution profiling
- Bottleneck detection and analysis
- SQL query optimization suggestions
- Spark job tuning recommendations
- Cost analysis and optimization
- Historical performance trending
Usage:
# Analyze Airflow DAG
python scripts/etl_performance_optimizer.py \
--airflow-db postgresql://host/airflow \
--dag-id sales_etl_pipeline \
--days 30 \
--optimize
# Spark job analysis
python scripts/etl_performance_optimizer.py \
--spark-history-server http://spark-history:18080 \
--app-id app-20250115-001
Complete Documentation: See tools.md for profiling options, optimization strategies, and cost analysis.
stream_processor.py
Streaming pipeline configuration generator and validator for Kafka, Flink, and Kinesis.
Capabilities:
- Multi-platform support (Kafka, Flink, Kinesis, Spark Streaming)
- Configuration validation with best practice checks
- Flink/Spark job scaffolding generation
- Kafka topic configuration generation
- Docker Compose for local streaming stacks
- Exactly-once semantics configuration
Usage:
# Validate configuration
python scripts/stream_processor.py --config streaming_config.yaml --validate
# Generate Kafka configurations
python scripts/stream_processor.py --config streaming_config.yaml --mode kafka --generate
# Generate Flink job scaffolding
python scripts/stream_processor.py --config streaming_config.yaml --mode flink --generate --output flink-jobs/
# Generate Docker Compose for local development
python scripts/stream_processor.py --config streaming_config.yaml --mode docker --generate
Complete Documentation: See tools.md for configuration format, validation checks, and generated outputs.
streaming_quality_validator.py
Real-time streaming data quality monitoring with comprehensive health scoring.
Capabilities:
- Consumer lag monitoring with thresholds
- Data freshness validation (P50/P95/P99 latency)
- Schema drift detection
- Throughput analysis (events/sec, bytes/sec)
- Dead letter queue rate monitoring
- Overall quality scoring with recommendations
- Prometheus metrics export
Usage:
# Monitor consumer lag
python scripts/streaming_quality_validator.py \
--lag --consumer-group events-processor --threshold 10000
# Monitor data freshness
python scripts/streaming_quality_validator.py \
--freshness --topic processed-events --max-latency-ms 5000
# Full quality validation
python scripts/streaming_quality_validator.py \
--lag --freshness --throughput --dlq \
--output streaming-health-report.html
Complete Documentation: See tools.md for all monitoring dimensions and integration patterns.
kafka_config_generator.py
Production-grade Kafka configuration generator with performance and security profiles.
Capabilities:
- Topic configuration (partitions, replication, retention, compaction)
- Producer profiles (high-throughput, exactly-once, low-latency, ordered)
- Consumer profiles (exactly-once, high-throughput, batch)
- Kafka Streams configuration with state store tuning
- Security configuration (SASL-PLAIN, SASL-SCRAM, mTLS)
- Kafka Connect source/sink configurations
- Multiple output formats (properties, YAML, JSON)
Usage:
# Generate topic configuration
python scripts/kafka_config_generator.py \
--topic user-events --partitions 12 --replication 3 --retention-hours 168
# Generate exactly-once producer
python scripts/kafka_config_generator.py \
--producer --profile exactly-once --transactional-id producer-001
# Generate Kafka Streams config
python scripts/kafka_config_generator.py \
--streams --application-id events-processor --exactly-once
Complete Documentation: See tools.md for all profiles, security options, and Connect configurations.
Reference Documentation
Frameworks (frameworks.md)
Comprehensive data engineering frameworks and patterns:
- Architecture Patterns: Lambda, Kappa, Medallion, Microservices data architecture
- Data Modeling: Dimensional (Kimball), Data Vault 2.0, One Big Table
- ETL/ELT Patterns: Full load, incremental load, CDC, SCD, idempotent pipelines
- Data Quality: Complete framework covering all quality dimensions
- DataOps: CI/CD for data pipelines, testing strategies, monitoring
- Orchestration: Airflow DAG patterns, backfill strategies
- Real-Time Streaming: Stateful processing, stream joins, windowing strategies, exactly-once semantics, event time processing, watermarks, backpressure, Apache Flink patterns, AWS Kinesis patterns, CDC for streaming
- Governance: Data catalog, lineage tracking, access control
Templates (templates.md)
Production-ready code templates and examples:
- Airflow DAGs: Complete ETL DAG, incremental load, dynamic task generation
- Spark Jobs: Batch processing, streaming, optimized configurations
- dbt Models: Staging, intermediate, fact tables, dimensions with SCD Type 2
- SQL Patterns: Incremental merge (upsert), deduplication, date spine, window functions
- Python Pipelines: Data quality validation class, retry decorators, error handling
- Real-Time Streaming: Apache Flink DataStream jobs (Java), Kafka Streams applications, PyFlink jobs, AWS Kinesis consumers, Docker Compose for streaming stack
- Kafka Configs: Producer/consumer properties templates, topic configurations, security configurations
- Docker: Dockerfiles for data pipelines, Docker Compose for local development including streaming stack (Kafka, Flink, Schema Registry)
- Configuration: dbt project config, Spark configuration, Airflow variables, streaming pipeline YAML
- Testing: pytest fixtures, integration tests, data quality tests
Tools (tools.md)
Python automation tool documentation:
- pipeline_orchestrator.py: Complete usage guide, configuration format, DAG templates
- data_quality_validator.py: Validation rules, dimension checks, Great Expectations integration
- etl_performance_optimizer.py: Performance analysis, query optimization, Spark tuning
- stream_processor.py: Streaming pipeline configuration, validation, job scaffolding generation
- streaming_quality_validator.py: Consumer lag, data freshness, schema drift, throughput monitoring
- kafka_config_generator.py: Topic, producer, consumer, Kafka Streams, and Connect configurations
- Integration Patterns: Airflow, dbt, CI/CD, monitoring systems, Prometheus
- Best Practices: Configuration management, error handling, performance, monitoring, streaming quality
Tech Stack
Core Technologies:
- Languages: Python 3.8+, SQL, Scala (Spark), Java (Flink)
- Orchestration: Apache Airflow, Prefect, Dagster
- Batch Processing: Apache Spark, dbt, Pandas
- Stream Processing: Apache Kafka, Apache Flink, Kafka Streams, Spark Structured Streaming, AWS Kinesis
- Storage: PostgreSQL, BigQuery, Snowflake, Redshift, S3, GCS
- Schema Management: Confluent Schema Registry, AWS Glue Schema Registry
- Containerization: Docker, Kubernetes
- Monitoring: Datadog, Prometheus, Grafana, Kafka UI
Data Platforms:
- Cloud Data Warehouses: Snowflake, BigQuery, Redshift
- Data Lakes: Delta Lake, Apache Iceberg, Apache Hudi
- Streaming Platforms: Apache Kafka, AWS Kinesis, Google Pub/Sub, Azure Event Hubs
- Stream Processing Engines: Apache Flink, Kafka Streams, Spark Structured Streaming
- Workflow: Airflow, Prefect, Dagster
Integration Points
This skill integrates with:
- Orchestration: Airflow, Prefect, Dagster for workflow management
- Transformation: dbt for SQL transformations and testing
- Quality: Great Expectations for data validation
- Monitoring: Datadog, Prometheus for pipeline monitoring
- BI Tools: Looker, Tableau, Power BI for analytics
- ML Platforms: MLflow, Kubeflow for ML pipeline integration
- Version Control: Git for pipeline code and configuration
See tools.md for detailed integration patterns and examples.
Best Practices
Pipeline Design:
- Idempotent operations for safe reruns
- Incremental processing where possible
- Clear data lineage and documentation
- Comprehensive error handling
- Automated recovery mechanisms
Data Quality:
- Define quality rules early
- Validate at every pipeline stage
- Automate quality monitoring
- Track quality trends over time
- Block bad data from downstream
Performance:
- Partition large tables by date/region
- Use columnar formats (Parquet, ORC)
- Leverage predicate pushdown
- Optimize for your query patterns
- Monitor and tune regularly
Operations:
- Version control everything
- Automate testing and deployment
- Implement comprehensive monitoring
- Document runbooks for incidents
- Regular performance reviews
Performance Targets
Batch Pipeline Execution:
- P50 latency: < 5 minutes (hourly pipelines)
- P95 latency: < 15 minutes
- Success rate: > 99%
- Data freshness: < 1 hour behind source
Streaming Pipeline Execution:
- Throughput: 10K+ events/second sustained
- End-to-end latency: P99 < 1 second
- Consumer lag: < 10K records behind
- Exactly-once delivery: Zero duplicates or losses
Data Quality (Batch):
- Quality score: > 95%
- Completeness: > 99%
- Timeliness: < 2 hours data lag
- Zero critical failures
Streaming Quality:
- Data freshness: P95 < 5 minutes from event generation
- Late data rate: < 5% outside watermark window
- Dead letter queue rate: < 1%
- Schema compatibility: 100% backward/forward compatible changes
Cost Efficiency:
- Cost per GB processed: < $0.10
- Cloud cost trend: Stable or decreasing
- Resource utilization: > 70%
Resources
- Frameworks Guide: references/frameworks.md
- Code Templates: references/templates.md
- Tool Documentation: references/tools.md
- Python Scripts:
scripts/ directory
Version: 2.0.0
Last Updated: December 16, 2025
Documentation Structure: Progressive disclosure with comprehensive references
Streaming Enhancement: Task #8 - Real-time streaming capabilities added
1---2name: senior-data-engineer3description: World-class data engineering skill for building scalable data pipelines, ETL/ELT systems, real-time streaming, and data infrastructure. Expertise in Python, SQL, Spark, Airflow, dbt, Kafka, Flink, Kinesis, and modern data stack. Includes data modeling, pipeline orchestration, data quality, streaming quality monitoring, and DataOps. Use when designing data architectures, building batch or streaming data pipelines, optimizing data workflows, or implementing data governance.4license: MIT5---67# Senior Data Engineer89## Core Capabilities1011- **Batch Pipeline Orchestration** - Design and implement production-ready ETL/ELT pipelines with Airflow, intelligent dependency resolution, retry logic, and comprehensive monitoring12- **Real-Time Streaming** - Build event-driven streaming pipelines with Kafka, Flink, Kinesis, and Spark Streaming with exactly-once semantics and sub-second latency13- **Data Quality Management** - Comprehensive batch and streaming data quality validation covering completeness, accuracy, consistency, timeliness, and validity14- **Streaming Quality Monitoring** - Track consumer lag, data freshness, schema drift, throughput, and dead letter queue rates for streaming pipelines15- **Performance Optimization** - Analyze and optimize pipeline performance with query optimization, Spark tuning, and cost analysis recommendations161718## Key Workflows1920### Workflow 1: Build ETL Pipeline2122**Time:** 2-4 hours2324**Steps:**251. Design pipeline architecture using Lambda, Kappa, or Medallion pattern262. Configure YAML pipeline definition with sources, transformations, targets273. Generate Airflow DAG with `pipeline_orchestrator.py`284. Define data quality validation rules295. Deploy and configure monitoring/alerting3031**Expected Output:** Production-ready ETL pipeline with 99%+ success rate, automated quality checks, and comprehensive monitoring3233### Workflow 2: Build Real-Time Streaming Pipeline3435**Time:** 3-5 days3637**Steps:**381. Select streaming architecture (Kappa vs Lambda) based on requirements392. Configure streaming pipeline YAML (sources, processing, sinks, quality)403. Generate Kafka configurations with `kafka_config_generator.py`414. Generate Flink/Spark job scaffolding with `stream_processor.py`425. Deploy and monitor with `streaming_quality_validator.py`4344**Expected Output:** Streaming pipeline processing 10K+ events/sec with P99 latency < 1s, exactly-once delivery, and real-time quality monitoring454647World-class data engineering for production-grade data systems, scalable pipelines, and enterprise data platforms.4849## Overview5051This skill provides comprehensive expertise in data engineering fundamentals through advanced production patterns. From designing medallion architectures to implementing real-time streaming pipelines, it covers the full spectrum of modern data engineering including ETL/ELT design, data quality frameworks, pipeline orchestration, and DataOps practices.5253**What This Skill Provides:**54- Production-ready pipeline templates (Airflow, Spark, dbt)55- Comprehensive data quality validation framework56- Performance optimization and cost analysis tools57- Data architecture patterns (Lambda, Kappa, Medallion)58- Complete DataOps CI/CD workflows5960**Best For:**61- Building scalable data pipelines for enterprise systems62- Implementing data quality and governance frameworks63- Optimizing ETL performance and cloud costs64- Designing modern data architectures (lake, warehouse, lakehouse)65- Production ML/AI data infrastructure6667## Quick Start6869### Pipeline Orchestration7071```bash72# Generate Airflow DAG from configuration73python scripts/pipeline_orchestrator.py --config pipeline_config.yaml --output dags/7475# Validate pipeline configuration76python scripts/pipeline_orchestrator.py --config pipeline_config.yaml --validate7778# Use incremental load template79python scripts/pipeline_orchestrator.py --template incremental --output dags/80```8182### Data Quality Validation8384```bash85# Validate CSV file with quality checks86python scripts/data_quality_validator.py --input data/sales.csv --output report.html8788# Validate database table with custom rules89python scripts/data_quality_validator.py \90 --connection postgresql://user:pass@host/db \91 --table sales_transactions \92 --rules rules/sales_validation.yaml \93 --threshold 0.9594```9596### Performance Optimization9798```bash99# Analyze pipeline performance and get recommendations100python scripts/etl_performance_optimizer.py \101 --airflow-db postgresql://host/airflow \102 --dag-id sales_etl_pipeline \103 --days 30 \104 --optimize105106# Analyze Spark job performance107python scripts/etl_performance_optimizer.py \108 --spark-history-server http://spark-history:18080 \109 --app-id app-20250115-001110```111112### Real-Time Streaming113114```bash115# Validate streaming pipeline configuration116python scripts/stream_processor.py --config streaming_config.yaml --validate117118# Generate Kafka topic and client configurations119python scripts/kafka_config_generator.py \120 --topic user-events \121 --partitions 12 \122 --replication 3 \123 --output kafka/topics/124125# Generate exactly-once producer configuration126python scripts/kafka_config_generator.py \127 --producer \128 --profile exactly-once \129 --output kafka/producer.properties130131# Generate Flink job scaffolding132python scripts/stream_processor.py \133 --config streaming_config.yaml \134 --mode flink \135 --generate \136 --output flink-jobs/137138# Monitor streaming quality139python scripts/streaming_quality_validator.py \140 --lag --consumer-group events-processor --threshold 10000 \141 --freshness --topic processed-events --max-latency-ms 5000 \142 --output streaming-health-report.html143```144145## Core Workflows146147### 1. Building Production Data Pipelines148149**Steps:**1501. **Design Architecture:** Choose pattern (Lambda, Kappa, Medallion) based on requirements1512. **Configure Pipeline:** Create YAML configuration with sources, transformations, targets1523. **Generate DAG:** `python scripts/pipeline_orchestrator.py --config config.yaml`1534. **Add Quality Checks:** Define validation rules for data quality1545. **Deploy & Monitor:** Deploy to Airflow, configure alerts, track metrics155156**Pipeline Patterns:** See [frameworks.md](references/frameworks.md) for Lambda Architecture, Kappa Architecture, Medallion Architecture (Bronze/Silver/Gold), and Microservices Data patterns.157158**Templates:** See [templates.md](references/templates.md) for complete Airflow DAG templates, Spark job templates, dbt models, and Docker configurations.159160### 2. Data Quality Management161162**Steps:**1631. **Define Rules:** Create validation rules covering completeness, accuracy, consistency1642. **Run Validation:** `python scripts/data_quality_validator.py --rules rules.yaml`1653. **Review Results:** Analyze quality scores and failed checks1664. **Integrate CI/CD:** Add validation to pipeline deployment process1675. **Monitor Trends:** Track quality scores over time168169**Quality Framework:** See [frameworks.md](references/frameworks.md) for complete Data Quality Framework covering all dimensions (completeness, accuracy, consistency, timeliness, validity).170171**Validation Templates:** See [templates.md](references/templates.md) for validation configuration examples and Python API usage.172173### 3. Data Modeling & Transformation174175**Steps:**1761. **Choose Modeling Approach:** Dimensional (Kimball), Data Vault 2.0, or One Big Table1772. **Design Schema:** Define fact tables, dimensions, and relationships1783. **Implement with dbt:** Create staging, intermediate, and mart models1794. **Handle SCD:** Implement slowly changing dimension logic (Type 1/2/3)1805. **Test & Deploy:** Run dbt tests, generate documentation, deploy181182**Modeling Patterns:** See [frameworks.md](references/frameworks.md) for Dimensional Modeling (Kimball), Data Vault 2.0, One Big Table (OBT), and SCD implementations.183184**dbt Templates:** See [templates.md](references/templates.md) for complete dbt model templates including staging, intermediate, fact tables, and SCD Type 2 logic.185186### 4. Performance Optimization187188**Steps:**1891. **Profile Pipeline:** Run performance analyzer on recent pipeline executions1902. **Identify Bottlenecks:** Review execution time breakdown and slow tasks1913. **Apply Optimizations:** Implement recommendations (partitioning, indexing, batching)1924. **Tune Spark Jobs:** Optimize memory, parallelism, and shuffle settings1935. **Measure Impact:** Compare before/after metrics, track cost savings194195**Optimization Strategies:** See [frameworks.md](references/frameworks.md) for performance best practices including partitioning strategies, query optimization, and Spark tuning.196197**Analysis Tools:** See [tools.md](references/tools.md) for complete documentation on etl_performance_optimizer.py with query analysis and Spark tuning.198199### 5. Building Real-Time Streaming Pipelines200201**Steps:**2021. **Architecture Selection:** Choose Kappa (streaming-only) or Lambda (batch + streaming) architecture2032. **Configure Pipeline:** Create YAML config with sources, processing engine, sinks, quality thresholds2043. **Generate Kafka Configs:** `python scripts/kafka_config_generator.py --topic events --partitions 12`2054. **Generate Job Scaffolding:** `python scripts/stream_processor.py --mode flink --generate`2065. **Deploy Infrastructure:** Use Docker Compose for local dev, Kubernetes for production2076. **Monitor Quality:** `python scripts/streaming_quality_validator.py --lag --freshness --throughput`208209**Streaming Patterns:** See [frameworks.md](references/frameworks.md) for stateful processing, stream joins, windowing, exactly-once semantics, and CDC patterns.210211**Templates:** See [templates.md](references/templates.md) for Flink DataStream jobs, Kafka Streams applications, PyFlink templates, and Docker Compose configurations.212213## Python Tools214215### pipeline_orchestrator.py216217Automated Airflow DAG generation with intelligent dependency resolution and monitoring.218219**Key Features:**220- Generate production-ready DAGs from YAML configuration221- Automatic task dependency resolution222- Built-in retry logic and error handling223- Multi-source support (PostgreSQL, S3, BigQuery, Snowflake)224- Integrated quality checks and alerting225226**Usage:**227```bash228# Basic DAG generation229python scripts/pipeline_orchestrator.py --config pipeline_config.yaml --output dags/230231# With validation232python scripts/pipeline_orchestrator.py --config config.yaml --validate233234# From template235python scripts/pipeline_orchestrator.py --template incremental --output dags/236```237238**Complete Documentation:** See [tools.md](references/tools.md) for full configuration options, templates, and integration examples.239240### data_quality_validator.py241242Comprehensive data quality validation framework with automated checks and reporting.243244**Capabilities:**245- Multi-dimensional validation (completeness, accuracy, consistency, timeliness, validity)246- Great Expectations integration247- Custom business rule validation248- HTML/PDF report generation249- Anomaly detection250- Historical trend tracking251252**Usage:**253```bash254# Validate with custom rules255python scripts/data_quality_validator.py \256 --input data/sales.csv \257 --rules rules/sales_validation.yaml \258 --output report.html259260# Database table validation261python scripts/data_quality_validator.py \262 --connection postgresql://host/db \263 --table sales_transactions \264 --threshold 0.95265```266267**Complete Documentation:** See [tools.md](references/tools.md) for rule configuration, API usage, and integration patterns.268269### etl_performance_optimizer.py270271Pipeline performance analysis with actionable optimization recommendations.272273**Capabilities:**274- Airflow DAG execution profiling275- Bottleneck detection and analysis276- SQL query optimization suggestions277- Spark job tuning recommendations278- Cost analysis and optimization279- Historical performance trending280281**Usage:**282```bash283# Analyze Airflow DAG284python scripts/etl_performance_optimizer.py \285 --airflow-db postgresql://host/airflow \286 --dag-id sales_etl_pipeline \287 --days 30 \288 --optimize289290# Spark job analysis291python scripts/etl_performance_optimizer.py \292 --spark-history-server http://spark-history:18080 \293 --app-id app-20250115-001294```295296**Complete Documentation:** See [tools.md](references/tools.md) for profiling options, optimization strategies, and cost analysis.297298### stream_processor.py299300Streaming pipeline configuration generator and validator for Kafka, Flink, and Kinesis.301302**Capabilities:**303- Multi-platform support (Kafka, Flink, Kinesis, Spark Streaming)304- Configuration validation with best practice checks305- Flink/Spark job scaffolding generation306- Kafka topic configuration generation307- Docker Compose for local streaming stacks308- Exactly-once semantics configuration309310**Usage:**311```bash312# Validate configuration313python scripts/stream_processor.py --config streaming_config.yaml --validate314315# Generate Kafka configurations316python scripts/stream_processor.py --config streaming_config.yaml --mode kafka --generate317318# Generate Flink job scaffolding319python scripts/stream_processor.py --config streaming_config.yaml --mode flink --generate --output flink-jobs/320321# Generate Docker Compose for local development322python scripts/stream_processor.py --config streaming_config.yaml --mode docker --generate323```324325**Complete Documentation:** See [tools.md](references/tools.md) for configuration format, validation checks, and generated outputs.326327### streaming_quality_validator.py328329Real-time streaming data quality monitoring with comprehensive health scoring.330331**Capabilities:**332- Consumer lag monitoring with thresholds333- Data freshness validation (P50/P95/P99 latency)334- Schema drift detection335- Throughput analysis (events/sec, bytes/sec)336- Dead letter queue rate monitoring337- Overall quality scoring with recommendations338- Prometheus metrics export339340**Usage:**341```bash342# Monitor consumer lag343python scripts/streaming_quality_validator.py \344 --lag --consumer-group events-processor --threshold 10000345346# Monitor data freshness347python scripts/streaming_quality_validator.py \348 --freshness --topic processed-events --max-latency-ms 5000349350# Full quality validation351python scripts/streaming_quality_validator.py \352 --lag --freshness --throughput --dlq \353 --output streaming-health-report.html354```355356**Complete Documentation:** See [tools.md](references/tools.md) for all monitoring dimensions and integration patterns.357358### kafka_config_generator.py359360Production-grade Kafka configuration generator with performance and security profiles.361362**Capabilities:**363- Topic configuration (partitions, replication, retention, compaction)364- Producer profiles (high-throughput, exactly-once, low-latency, ordered)365- Consumer profiles (exactly-once, high-throughput, batch)366- Kafka Streams configuration with state store tuning367- Security configuration (SASL-PLAIN, SASL-SCRAM, mTLS)368- Kafka Connect source/sink configurations369- Multiple output formats (properties, YAML, JSON)370371**Usage:**372```bash373# Generate topic configuration374python scripts/kafka_config_generator.py \375 --topic user-events --partitions 12 --replication 3 --retention-hours 168376377# Generate exactly-once producer378python scripts/kafka_config_generator.py \379 --producer --profile exactly-once --transactional-id producer-001380381# Generate Kafka Streams config382python scripts/kafka_config_generator.py \383 --streams --application-id events-processor --exactly-once384```385386**Complete Documentation:** See [tools.md](references/tools.md) for all profiles, security options, and Connect configurations.387388## Reference Documentation389390### Frameworks ([frameworks.md](references/frameworks.md))391392Comprehensive data engineering frameworks and patterns:393- **Architecture Patterns:** Lambda, Kappa, Medallion, Microservices data architecture394- **Data Modeling:** Dimensional (Kimball), Data Vault 2.0, One Big Table395- **ETL/ELT Patterns:** Full load, incremental load, CDC, SCD, idempotent pipelines396- **Data Quality:** Complete framework covering all quality dimensions397- **DataOps:** CI/CD for data pipelines, testing strategies, monitoring398- **Orchestration:** Airflow DAG patterns, backfill strategies399- **Real-Time Streaming:** Stateful processing, stream joins, windowing strategies, exactly-once semantics, event time processing, watermarks, backpressure, Apache Flink patterns, AWS Kinesis patterns, CDC for streaming400- **Governance:** Data catalog, lineage tracking, access control401402### Templates ([templates.md](references/templates.md))403404Production-ready code templates and examples:405- **Airflow DAGs:** Complete ETL DAG, incremental load, dynamic task generation406- **Spark Jobs:** Batch processing, streaming, optimized configurations407- **dbt Models:** Staging, intermediate, fact tables, dimensions with SCD Type 2408- **SQL Patterns:** Incremental merge (upsert), deduplication, date spine, window functions409- **Python Pipelines:** Data quality validation class, retry decorators, error handling410- **Real-Time Streaming:** Apache Flink DataStream jobs (Java), Kafka Streams applications, PyFlink jobs, AWS Kinesis consumers, Docker Compose for streaming stack411- **Kafka Configs:** Producer/consumer properties templates, topic configurations, security configurations412- **Docker:** Dockerfiles for data pipelines, Docker Compose for local development including streaming stack (Kafka, Flink, Schema Registry)413- **Configuration:** dbt project config, Spark configuration, Airflow variables, streaming pipeline YAML414- **Testing:** pytest fixtures, integration tests, data quality tests415416### Tools ([tools.md](references/tools.md))417418Python automation tool documentation:419- **pipeline_orchestrator.py:** Complete usage guide, configuration format, DAG templates420- **data_quality_validator.py:** Validation rules, dimension checks, Great Expectations integration421- **etl_performance_optimizer.py:** Performance analysis, query optimization, Spark tuning422- **stream_processor.py:** Streaming pipeline configuration, validation, job scaffolding generation423- **streaming_quality_validator.py:** Consumer lag, data freshness, schema drift, throughput monitoring424- **kafka_config_generator.py:** Topic, producer, consumer, Kafka Streams, and Connect configurations425- **Integration Patterns:** Airflow, dbt, CI/CD, monitoring systems, Prometheus426- **Best Practices:** Configuration management, error handling, performance, monitoring, streaming quality427428## Tech Stack429430**Core Technologies:**431- **Languages:** Python 3.8+, SQL, Scala (Spark), Java (Flink)432- **Orchestration:** Apache Airflow, Prefect, Dagster433- **Batch Processing:** Apache Spark, dbt, Pandas434- **Stream Processing:** Apache Kafka, Apache Flink, Kafka Streams, Spark Structured Streaming, AWS Kinesis435- **Storage:** PostgreSQL, BigQuery, Snowflake, Redshift, S3, GCS436- **Schema Management:** Confluent Schema Registry, AWS Glue Schema Registry437- **Containerization:** Docker, Kubernetes438- **Monitoring:** Datadog, Prometheus, Grafana, Kafka UI439440**Data Platforms:**441- **Cloud Data Warehouses:** Snowflake, BigQuery, Redshift442- **Data Lakes:** Delta Lake, Apache Iceberg, Apache Hudi443- **Streaming Platforms:** Apache Kafka, AWS Kinesis, Google Pub/Sub, Azure Event Hubs444- **Stream Processing Engines:** Apache Flink, Kafka Streams, Spark Structured Streaming445- **Workflow:** Airflow, Prefect, Dagster446447## Integration Points448449This skill integrates with:450- **Orchestration:** Airflow, Prefect, Dagster for workflow management451- **Transformation:** dbt for SQL transformations and testing452- **Quality:** Great Expectations for data validation453- **Monitoring:** Datadog, Prometheus for pipeline monitoring454- **BI Tools:** Looker, Tableau, Power BI for analytics455- **ML Platforms:** MLflow, Kubeflow for ML pipeline integration456- **Version Control:** Git for pipeline code and configuration457458See [tools.md](references/tools.md) for detailed integration patterns and examples.459460## Best Practices461462**Pipeline Design:**4631. Idempotent operations for safe reruns4642. Incremental processing where possible4653. Clear data lineage and documentation4664. Comprehensive error handling4675. Automated recovery mechanisms468469**Data Quality:**4701. Define quality rules early4712. Validate at every pipeline stage4723. Automate quality monitoring4734. Track quality trends over time4745. Block bad data from downstream475476**Performance:**4771. Partition large tables by date/region4782. Use columnar formats (Parquet, ORC)4793. Leverage predicate pushdown4804. Optimize for your query patterns4815. Monitor and tune regularly482483**Operations:**4841. Version control everything4852. Automate testing and deployment4863. Implement comprehensive monitoring4874. Document runbooks for incidents4885. Regular performance reviews489490## Performance Targets491492**Batch Pipeline Execution:**493- P50 latency: < 5 minutes (hourly pipelines)494- P95 latency: < 15 minutes495- Success rate: > 99%496- Data freshness: < 1 hour behind source497498**Streaming Pipeline Execution:**499- Throughput: 10K+ events/second sustained500- End-to-end latency: P99 < 1 second501- Consumer lag: < 10K records behind502- Exactly-once delivery: Zero duplicates or losses503504**Data Quality (Batch):**505- Quality score: > 95%506- Completeness: > 99%507- Timeliness: < 2 hours data lag508- Zero critical failures509510**Streaming Quality:**511- Data freshness: P95 < 5 minutes from event generation512- Late data rate: < 5% outside watermark window513- Dead letter queue rate: < 1%514- Schema compatibility: 100% backward/forward compatible changes515516**Cost Efficiency:**517- Cost per GB processed: < $0.10518- Cloud cost trend: Stable or decreasing519- Resource utilization: > 70%520521## Resources522523- **Frameworks Guide:** [references/frameworks.md](references/frameworks.md)524- **Code Templates:** [references/templates.md](references/templates.md)525- **Tool Documentation:** [references/tools.md](references/tools.md)526- **Python Scripts:** `scripts/` directory527528---529530**Version:** 2.0.0531**Last Updated:** December 16, 2025532**Documentation Structure:** Progressive disclosure with comprehensive references533**Streaming Enhancement:** Task #8 - Real-time streaming capabilities added