Skip to main content
senior-data-engineer 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.
Jump to install Skills Marketplace Discover and explore AI skills built by the community.
Install with Codex or Claude Copy this prompt, paste it into Codex, Claude, or another assistant, and let it review the skill page and install it for you.
Copy promptShow prompt details A direct command skips the review prompt. Inspect the source before running it.
npx skills add https://github.com/benchflow-ai/skillsbench --skill senior-data-engineerThe command stays on one line. Scroll horizontally to inspect it before copying.
Prefer a local copy? Download the files currently available to SkillsMP.
Download Zip Downloading... More from this repository SkillsBench contribution workflow. Use when: (1) Creating benchmark tasks, (2) Understanding repo structure, (3) Preparing PRs for task submission.
SkillsBench task authoring — walk a contributor from idea to submission-ready task following CONTRIBUTING.md and the task-implementation rubric. Use when the user wants to create a new SkillsBench task, scaffold a task from an existing workflow (notebook, Excel workbook, document, dataset), convert a prompt or a benchmark item into a SkillsBench task, write skills for a task, or prepare a SkillsBench PR. Pairs with `task-review` (run that as a self-check before submitting).
SkillsBench task PR review — classifies the task track (standard / research / multimodal), runs static policy checks against the track-specific rubric, benchmarks the task across oracle plus Claude and Codex (with and without skills), audits trajectories for cheating and skill invocation, and produces a `pr-N-task-timestamp-run.txt` review report alongside a `prN.zip` bundle of trajectories. Use when reviewing a SkillsBench task PR (by number, branch, or local task path), when the user asks to review a task, run benchmarks on a PR, audit a submission, classify a task as research or multimodal track, or prepare a comment to post on a SkillsBench PR.
Related occupations SOC
Based on SOC occupation classification
name senior-data-engineer title Senior Data Engineer Skill Package description 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. domain engineering subdomain data-engineering difficulty advanced time-saved TODO: Quantify time savings frequency TODO: Estimate usage frequency use-cases ["Designing data pipelines for ETL/ELT processes","Building data warehouses and data lakes","Implementing data quality and governance frameworks","Creating analytics dashboards and reporting","Building real-time streaming pipelines with Kafka and Flink","Implementing exactly-once streaming semantics","Monitoring streaming quality (consumer lag, data freshness, schema drift)"] related-agents [] related-skills [] related-commands [] orchestrated-by [] dependencies {"scripts":[],"references":[],"assets":[]} compatibility Python 3.8+; platforms: macos, linux, windows tech-stack ["Python","SQL","Apache Spark","Airflow","dbt","Apache Kafka","Apache Flink","AWS Kinesis","Spark Structured Streaming","Kafka Streams","PostgreSQL","BigQuery","Snowflake","Docker","Schema Registry"] examples [{"title":"Example Usage","input":"TODO: Add example input for senior-data-engineer","output":"TODO: Add expected output"}] stats {"downloads":0,"stars":0,"rating":0,"reviews":0} version v2.0.0 author Claude Skills Team contributors [] created "2025-10-20T00:00:00.000Z" updated "2025-12-16T00:00:00.000Z" license MIT tags ["architecture","data","design","engineer","engineering","senior","streaming","kafka","flink","real-time"] featured false verified true
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
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
python scripts/pipeline_orchestrator.py --config pipeline_config.yaml --output dags/
python scripts/pipeline_orchestrator.py --config pipeline_config.yaml --validate
python scripts/pipeline_orchestrator.py --template incremental --output dags/
Data Quality Validation
python scripts/data_quality_validator.py --input data/sales.csv --output report.html
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
python scripts/etl_performance_optimizer.py \
--airflow-db postgresql://host/airflow \
--dag-id sales_etl_pipeline \
--days 30 \
--optimize
python scripts/etl_performance_optimizer.py \
--spark-history-server http://spark-history:18080 \
--app-id app-20250115-001
Real-Time Streaming
python scripts/stream_processor.py --config streaming_config.yaml --validate
python scripts/kafka_config_generator.py \
--topic user-events \
--partitions 12 \
--replication 3 \
--output kafka/topics/
python scripts/kafka_config_generator.py \
--producer \
--profile exactly-once \
--output kafka/producer.properties
python scripts/stream_processor.py \
--config streaming_config.yaml \
--mode flink \
--generate \
--output flink-jobs/
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
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
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
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
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
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.
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
python scripts/pipeline_orchestrator.py --config pipeline_config.yaml --output dags/
python scripts/pipeline_orchestrator.py --config config.yaml --validate
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.
Multi-dimensional validation (completeness, accuracy, consistency, timeliness, validity)
Great Expectations integration
Custom business rule validation
HTML/PDF report generation
Anomaly detection
Historical trend tracking
python scripts/data_quality_validator.py \
--input data/sales.csv \
--rules rules/sales_validation.yaml \
--output report.html
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.
Airflow DAG execution profiling
Bottleneck detection and analysis
SQL query optimization suggestions
Spark job tuning recommendations
Cost analysis and optimization
Historical performance trending
python scripts/etl_performance_optimizer.py \
--airflow-db postgresql://host/airflow \
--dag-id sales_etl_pipeline \
--days 30 \
--optimize
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.
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
python scripts/stream_processor.py --config streaming_config.yaml --validate
python scripts/stream_processor.py --config streaming_config.yaml --mode kafka --generate
python scripts/stream_processor.py --config streaming_config.yaml --mode flink --generate --output flink-jobs/
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.
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
python scripts/streaming_quality_validator.py \
--lag --consumer-group events-processor --threshold 10000
python scripts/streaming_quality_validator.py \
--freshness --topic processed-events --max-latency-ms 5000
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.
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)
python scripts/kafka_config_generator.py \
--topic user-events --partitions 12 --replication 3 --retention-hours 168
python scripts/kafka_config_generator.py \
--producer --profile exactly-once --transactional-id producer-001
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 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
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
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
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
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
Idempotent operations for safe reruns
Incremental processing where possible
Clear data lineage and documentation
Comprehensive error handling
Automated recovery mechanisms
Define quality rules early
Validate at every pipeline stage
Automate quality monitoring
Track quality trends over time
Block bad data from downstream
Partition large tables by date/region
Use columnar formats (Parquet, ORC)
Leverage predicate pushdown
Optimize for your query patterns
Monitor and tune regularly
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
Quality score: > 95%
Completeness: > 99%
Timeliness: < 2 hours data lag
Zero critical failures
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 per GB processed: < $0.10
Cloud cost trend: Stable or decreasing
Resource utilization: > 70%
Resources
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