| name | data-engineering-study-material |
| description | Comprehensive study guide covering data engineering concepts, tools, and best practices for learning and reference |
| triggers | ["explain data engineering concepts","show me data engineering study materials","what are data engineering best practices","help me learn data engineering","guide me through data engineering topics","data engineering interview preparation","overview of data engineering tools","data engineering learning resources"] |
Data Engineering Study Material
Skill by ara.so — Data Skills collection.
Overview
This project is a comprehensive study guide and reference repository for data engineering concepts, tools, and practices. It serves as a centralized resource for learning core data engineering principles, understanding modern data stack components, and preparing for data engineering roles.
The repository covers:
- Data engineering fundamentals and architecture patterns
- ETL/ELT pipeline design and implementation
- Data warehousing and lake architectures
- Streaming and batch processing frameworks
- Cloud data platforms (AWS, GCP, Azure)
- Data quality, governance, and observability
- Infrastructure as Code and orchestration tools
- Interview preparation and best practices
Installation
This is a study material repository, not an installable package. Clone it to access the materials:
git clone https://github.com/Ahmeduddin3403/data-engineering-study-material.git
cd data-engineering-study-material
Repository Structure
The materials are typically organized by topic area:
data-engineering-study-material/
├── fundamentals/ # Core concepts and principles
├── tools/ # Tool-specific guides
├── architectures/ # Design patterns and architectures
├── pipelines/ # ETL/ELT examples
├── cloud-platforms/ # Cloud-specific implementations
├── streaming/ # Real-time processing
├── batch-processing/ # Batch job patterns
├── data-quality/ # Testing and validation
├── orchestration/ # Workflow management
├── interview-prep/ # Interview questions and answers
└── projects/ # Hands-on project examples
Core Data Engineering Concepts
ETL Pipeline Example (Python)
import pandas as pd
from sqlalchemy import create_engine
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
class ETLPipeline:
"""Simple ETL pipeline for extracting, transforming, and loading data"""
def __init__(self, source_path, target_conn_string):
self.source_path = source_path
self.engine = create_engine(target_conn_string)
def extract(self):
"""Extract data from source"""
logger.info(f"Extracting data from {self.source_path}")
df = pd.read_csv(self.source_path)
logger.info(f"Extracted {len(df)} rows")
return df
def transform(self, df):
"""Transform data: clean, deduplicate, enrich"""
logger.info("Transforming data")
df = df.drop_duplicates()
df = df.fillna({
'numeric_column': 0,
'string_column': 'Unknown'
})
df['created_date'] = pd.to_datetime(df['timestamp']).dt.date
df = df[df['amount'] > ]
logger.info()
df
():
logger.info()
df.to_sql(table_name, .engine, if_exists=, index=)
logger.info()
():
:
df = .extract()
df_transformed = .transform(df)
.load(df_transformed, table_name)
logger.info()
Exception e:
logger.error()
__name__ == :
pipeline = ETLPipeline(
source_path=,
target_conn_string=
)
pipeline.run()
Data Quality Checks
import great_expectations as ge
def validate_data_quality(df):
"""Implement data quality checks using Great Expectations"""
ge_df = ge.from_pandas(df)
expectations = {
'id': lambda col: col.expect_column_values_to_be_unique(),
'email': lambda col: col.expect_column_values_to_match_regex(r'^[\w\.-]+@[\w\.-]+\.\w+$'),
'amount': lambda col: col.expect_column_values_to_be_between(min_value=0, max_value=1000000),
'created_at': lambda col: col.expect_column_values_to_not_be_null(),
'status': lambda col: col.expect_column_values_to_be_in_set(['active', 'inactive', 'pending'])
}
results = []
for column, expectation_func in expectations.items():
if column in ge_df.columns:
result = expectation_func(ge_df[column])
results.append(result)
all_passed = all(r.success for r in results)
return all_passed, results
Apache Airflow DAG Example
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.providers.amazon.aws.transfers.s3_to_redshift import S3ToRedshiftOperator
from datetime import datetime, timedelta
default_args = {
'owner': 'data-engineering',
'depends_on_past': False,
'start_date': datetime(2024, 1, 1),
'email_on_failure': True,
'email_on_retry': False,
'retries': 2,
'retry_delay': timedelta(minutes=5)
}
def extract_from_api(**context):
"""Extract data from external API"""
import requests
import json
response = requests.get('https://api.example.com/data')
data = response.json()
s3_path = f"s3://my-bucket/raw/{context['ds']}/data.json"
return s3_path
def transform_data(**context):
"""Transform extracted data"""
import pandas as pd
s3_path = context[].xcom_pull(task_ids=)
df = pd.read_json(s3_path)
df_transformed = df.drop_duplicates()
df_transformed[] = context[]
output_path =
df_transformed.to_parquet(output_path)
output_path
DAG(
,
default_args=default_args,
description=,
schedule_interval=,
catchup=,
tags=[, ]
) dag:
extract_task = PythonOperator(
task_id=,
python_callable=extract_from_api,
provide_context=
)
transform_task = PythonOperator(
task_id=,
python_callable=transform_data,
provide_context=
)
load_task = S3ToRedshiftOperator(
task_id=,
s3_bucket=,
s3_key=,
schema=,
table=,
copy_options=[]
)
data_quality_check = PostgresOperator(
task_id=,
postgres_conn_id=,
sql=
)
extract_task >> transform_task >> load_task >> data_quality_check
Spark Batch Processing Example
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when, sum, avg, count, to_date
from pyspark.sql.window import Window
def process_batch_data():
"""Process large-scale batch data with Apache Spark"""
spark = SparkSession.builder \
.appName("BatchDataProcessing") \
.config("spark.sql.adaptive.enabled", "true") \
.getOrCreate()
df = spark.read \
.format("parquet") \
.load("s3a://data-lake/raw/transactions/")
df_transformed = df \
.withColumn("transaction_date", to_date(col("timestamp"))) \
.filter(col("amount") > 0) \
.dropDuplicates(["transaction_id"])
daily_summary = df_transformed.groupBy("transaction_date", "category") \
.agg(
count("transaction_id").alias("transaction_count"),
sum("amount").alias("total_amount"),
avg("amount").alias("avg_amount")
)
window_spec = Window.partitionBy("transaction_date").orderBy(col("total_amount").desc())
ranked_summary = daily_summary \
.withColumn("rank", dense_rank().over(window_spec)) \
.(col() <= )
ranked_summary.write \
.() \
.mode() \
.partitionBy() \
.save()
spark.stop()
__name__ == :
process_batch_data()
Streaming Pipeline Example (Kafka + Spark Structured Streaming)
from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col, window
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType
def create_streaming_pipeline():
"""Real-time data processing with Spark Structured Streaming"""
spark = SparkSession.builder \
.appName("RealTimeStreaming") \
.getOrCreate()
schema = StructType([
StructField("event_id", StringType(), True),
StructField("user_id", StringType(), True),
StructField("event_type", StringType(), True),
StructField("amount", DoubleType(), True),
StructField("timestamp", TimestampType(), True)
])
df_stream = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "events") \
.option("startingOffsets", "latest") \
.load()
df_parsed = df_stream.select(
from_json(col("value").cast("string"), schema).alias("data")
).select("data.*")
df_windowed = df_parsed \
.withWatermark("timestamp", "10 minutes") \
.groupBy(
window(col("timestamp"), , ),
col()
) \
.agg(
count().alias(),
().alias()
)
query = df_windowed.writeStream \
.outputMode() \
.() \
.option(, ) \
.option(, ) \
.trigger(processingTime=) \
.start()
query.awaitTermination()
__name__ == :
create_streaming_pipeline()
Common Patterns
Incremental Data Loading
def incremental_load(table_name, watermark_column, last_watermark):
"""Load only new/updated records since last run"""
query = f"""
SELECT *
FROM {table_name}
WHERE {watermark_column} > '{last_watermark}'
ORDER BY {watermark_column}
"""
new_data = pd.read_sql(query, source_conn)
if not new_data.empty:
new_watermark = new_data[watermark_column].max()
new_data.to_sql('target_table', target_conn, if_exists='append', index=False)
save_watermark(table_name, new_watermark)
return len(new_data)
SCD Type 2 Implementation
def scd_type2_merge(source_df, target_table, business_key, effective_date):
"""Implement Slowly Changing Dimension Type 2"""
from datetime import datetime
current_df = pd.read_sql(f"SELECT * FROM {target_table} WHERE is_current = 1", conn)
merged = source_df.merge(
current_df,
on=business_key,
how='left',
suffixes=('_new', '_old')
)
changed = merged[
(merged.apply(lambda row: row_has_changes(row), axis=1))
]
if not changed.empty:
expire_query = f"""
UPDATE {target_table}
SET is_current = 0,
end_date = '{effective_date}'
WHERE {business_key} IN ({','.join(map(str, changed[business_key].tolist()))})
AND is_current = 1
"""
conn.execute(expire_query)
new_records = source_df[source_df[business_key].isin(changed[business_key])]
new_records['start_date'] = effective_date
new_records['end_date'] = '9999-12-31'
new_records['is_current'] = 1
new_records.to_sql(target_table, conn, if_exists='append', index=False)
Configuration
Database Connection Configuration
import os
DATABASE_CONFIG = {
'source': {
'host': os.getenv('SOURCE_DB_HOST'),
'port': os.getenv('SOURCE_DB_PORT', 5432),
'database': os.getenv('SOURCE_DB_NAME'),
'user': os.getenv('SOURCE_DB_USER'),
'password': os.getenv('SOURCE_DB_PASSWORD')
},
'warehouse': {
'host': os.getenv('WAREHOUSE_HOST'),
'port': os.getenv('WAREHOUSE_PORT', 5439),
'database': os.getenv('WAREHOUSE_DB'),
'user': os.getenv('WAREHOUSE_USER'),
'password': os.getenv('WAREHOUSE_PASSWORD')
}
}
SPARK_CONFIG = {
'spark.executor.memory': '4g',
'spark.driver.memory': '2g',
'spark.sql.adaptive.enabled': 'true',
'spark.sql.adaptive.coalescePartitions.enabled': 'true'
}
AIRFLOW_CONFIG = {
'concurrency': 16,
'max_active_runs': 3,
'dagbag_import_timeout': 30
}
Troubleshooting
Common Issues and Solutions
Issue: Out of Memory in Spark Jobs
spark = SparkSession.builder \
.config("spark.executor.memory", "8g") \
.config("spark.driver.memory", "4g") \
.config("spark.sql.shuffle.partitions", "200") \
.config("spark.default.parallelism", "200") \
.getOrCreate()
from pyspark.sql.functions import broadcast
result = large_df.join(broadcast(small_df), "key")
Issue: Slow Incremental Loads
df = spark.read.parquet("s3://data/table/") \
.where(f"partition_date >= '{start_date}'")
Issue: Data Quality Failures
def validate_and_quarantine(df, rules):
"""Separate valid and invalid records"""
valid_df = df
invalid_records = []
for rule_name, rule_func in rules.items():
mask = rule_func(valid_df)
invalid = valid_df[~mask].copy()
invalid['failed_rule'] = rule_name
invalid_records.append(invalid)
valid_df = valid_df[mask]
if invalid_records:
pd.concat(invalid_records).to_sql(
'data_quality_quarantine',
conn,
if_exists='append'
)
return valid_df
Best Practices
- Idempotency: Ensure pipelines can be re-run safely
- Monitoring: Implement comprehensive logging and alerting
- Data Quality: Validate data at every stage
- Partitioning: Use appropriate partitioning strategies for performance
- Documentation: Document data lineage and transformations
- Version Control: Track schema changes and pipeline versions
- Testing: Test pipelines with sample data before production
- Security: Use IAM roles, encryption, and secure credential management
Interview Preparation
Common data engineering interview topics covered:
- SQL optimization and query tuning
- Distributed systems concepts
- Data modeling (star schema, snowflake, Data Vault)
- ETL vs ELT trade-offs
- CAP theorem and consistency models
- Data quality frameworks
- Cloud platform services (S3, Redshift, BigQuery, Databricks)
- Orchestration tools (Airflow, Prefect, Dagster)
- Streaming architectures (Kafka, Kinesis, Pub/Sub)