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

Source facts

Repository
benchflow-ai/skillsbench
Last source activity
June 5, 2026 at 23:46
Detected SKILL.md language
English
Stars
1,816
Forks
368

Install options

The review-first prompt is selected by default. You can switch to a direct command or download a local copy.

Review the source files

Read SKILL.md and any companion files shown by SkillsMP before deciding whether to install.

File Explorer
14 files

Showing SKILL.md

SKILL.md
Source instructions · Read-only preview
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:** 1. Design pipeline architecture using Lambda, Kappa, or Medallion pattern 2. Configure YAML pipeline definition with sources, transformations, targets 3. Generate Airflow DAG with `pipeline_orchestrator.py` 4. Define data quality validation rules 5. 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:** 1. Select streaming architecture (Kappa vs Lambda) based on requirements 2. Configure streaming pipeline YAML (sources, processing, sinks, quality) 3. Generate Kafka configurations with `kafka_config_generator.py` 4. Generate Flink/Spark job scaffolding with `stream_processor.py` 5. 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 ```bash # 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 ```bash # 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 ```bash # 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 ```bash # 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:** 1. **Design Architecture:** Choose pattern (Lambda, Kappa, Medallion) based on requirements 2. **Configure Pipeline:** Create YAML configuration with sources, transformations, targets 3. **Generate DAG:** `python scripts/pipeline_orchestrator.py --config config.yaml` 4. **Add Quality Checks:** Define validation rules for data quality 5. **Deploy & Monitor:** Deploy to Airflow, configure alerts, track metrics **Pipeline Patterns:** See [frameworks.md](references/frameworks.md) for Lambda Architecture, Kappa Architecture, Medallion Architecture (Bronze/Silver/Gold), and Microservices Data patterns. **Templates:** See [templates.md](references/templates.md) for complete Airflow DAG templates, Spark job templates, dbt models, and Docker configurations. ### 2. Data Quality Management **Steps:** 1. **Define Rules:** Create validation rules covering completeness, accuracy, consistency 2. **Run Validation:** `python scripts/data_quality_validator.py --rules rules.yaml` 3. **Review Results:** Analyze quality scores and failed checks 4. **Integrate CI/CD:** Add validation to pipeline deployment process 5. **Monitor Trends:** Track quality scores over time **Quality Framework:** See [frameworks.md](references/frameworks.md) for complete Data Quality Framework covering all dimensions (completeness, accuracy, consistency, timeliness, validity). **Validation Templates:** See [templates.md](references/templates.md) for validation configuration examples and Python API usage. ### 3. Data Modeling & Transformation **Steps:** 1. **Choose Modeling Approach:** Dimensional (Kimball), Data Vault 2.0, or One Big Table 2. **Design Schema:** Define fact tables, dimensions, and relationships 3. **Implement with dbt:** Create staging, intermediate, and mart models 4. **Handle SCD:** Implement slowly changing dimension logic (Type 1/2/3) 5. **Test & Deploy:** Run dbt tests, generate documentation, deploy **Modeling Patterns:** See [frameworks.md](references/frameworks.md) for Dimensional Modeling (Kimball), Data Vault 2.0, One Big Table (OBT), and SCD implementations. **dbt Templates:** See [templates.md](references/templates.md) for complete dbt model templates including staging, intermediate, fact tables, and SCD Type 2 logic. ### 4. Performance Optimization **Steps:** 1. **Profile Pipeline:** Run performance analyzer on recent pipeline executions 2. **Identify Bottlenecks:** Review execution time breakdown and slow tasks 3. **Apply Optimizations:** Implement recommendations (partitioning, indexing, batching) 4. **Tune Spark Jobs:** Optimize memory, parallelism, and shuffle settings 5. **Measure Impact:** Compare before/after metrics, track cost savings **Optimization Strategies:** See [frameworks.md](references/frameworks.md) for performance best practices including partitioning strategies, query optimization, and Spark tuning. **Analysis Tools:** See [tools.md](references/tools.md) for complete documentation on etl_performance_optimizer.py with query analysis and Spark tuning. ### 5. Building Real-Time Streaming Pipelines **Steps:** 1. **Architecture Selection:** Choose Kappa (streaming-only) or Lambda (batch + streaming) architecture 2. **Configure Pipeline:** Create YAML config with sources, processing engine, sinks, quality thresholds 3. **Generate Kafka Configs:** `python scripts/kafka_config_generator.py --topic events --partitions 12` 4. **Generate Job Scaffolding:** `python scripts/stream_processor.py --mode flink --generate` 5. **Deploy Infrastructure:** Use Docker Compose for local dev, Kubernetes for production 6. **Monitor Quality:** `python scripts/streaming_quality_validator.py --lag --freshness --throughput` **Streaming Patterns:** See [frameworks.md](references/frameworks.md) for stateful processing, stream joins, windowing, exactly-once semantics, and CDC patterns. **Templates:** See [templates.md](references/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:** ```bash # 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](references/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:** ```bash # 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](references/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:** ```bash # 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](references/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:** ```bash # 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](references/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:** ```bash # 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](references/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:** ```bash # 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](references/tools.md) for all profiles, security options, and Connect configurations. ## Reference Documentation ### Frameworks ([frameworks.md](references/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
View on GitHub
This SKILL.md is very large, so SkillsMP previews the first section here. View on GitHub