Comprehensive Python data engineering patterns for AWS Data Lake, including PySpark, Pandas, Apache Airflow, AWS Glue, ETL pipelines, data quality, schema management, performance optimization,...
Production-ready Python data engineering patterns for building scalable, reliable data pipelines on AWS Data Lake infrastructure. Covers PySpark, Pandas, Airflow, Glue, FastAPI, and modern data engineering best practices.
Auto-activates when working with:
API/Service Layer (FastAPI) ā Orchestration (Airflow) ā Processing (PySpark/Pandas) ā Storage (S3)
from fastapi import FastAPI, HTTPException, Depends
from pydantic import BaseModel, Field, validator
from typing import List, Optional
import boto3
from datetime import datetime
app = FastAPI(title="Data Lake API")
class DataRequest(BaseModel):
dataset_name: str = Field(..., regex="^[a-z0-9_]+$")
partition_date: datetime
limit: Optional[int] = Field(1000, gt=0, le=10000)
@validator('partition_date')
def validate_date(cls, v):
if v > datetime.now():
raise ValueError('partition_date cannot be in the future')
return v
@app.get("/datasets/{dataset_name}/data")
async def get_dataset_data(dataset_name: str, request: DataRequest = Depends()):
"""Query data from S3 Data Lake with Athena."""
try:
athena = boto3.client('athena')
query = f"""
SELECT * FROM {dataset_name}
WHERE partition_date = '{request.partition_date.strftime('%Y-%m-%d')}'
LIMIT {request.limit}
"""
response = athena.start_query_execution(
QueryString=query,
QueryExecutionContext={'Database': 'data_lake'},
ResultConfiguration={'OutputLocation': 's3://query-results/'}
)
return {"query_execution_id": response['QueryExecutionId']}
except Exception as e:
logging.error(f"Query failed: {e}")
raise HTTPException(status_code=500, detail=str(e))
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when, lit, current_timestamp
from pyspark.sql.types import StructType, StructField, StringType, IntegerType
import logging
class DataLakeETL:
def __init__(self, app_name: str):
self.spark = (SparkSession.builder
.appName(app_name)
.config("spark.sql.adaptive.enabled", "true")
.config("spark.sql.adaptive.coalescePartitions.enabled", "true")
.getOrCreate())
self.logger = logging.getLogger(__name__)
def read_bronze(self, path: str, schema: StructType):
"""Read from raw/bronze zone with schema enforcement."""
return (self.spark.read
.schema(schema)
.option("mode", "PERMISSIVE")
.option("columnNameOfCorruptRecord", "_corrupt_record")
.parquet(path))
def transform_to_silver(self, df):
"""Clean and validate data for silver zone."""
return (df
.filter(col("_corrupt_record").isNull()) # Remove corrupt records
.dropDuplicates(["id"]) # Deduplicate
.withColumn("processed_at", current_timestamp())
.withColumn("is_valid",
when(col("value").isNotNull() & (col("value") > 0), True)
.otherwise(False))
.filter(col("is_valid"))) # Only valid records
def write_silver(self, df, path: str, partition_cols: List[str]):
"""Write to processed/silver zone with partitioning."""
(df.write
.mode("overwrite")
.partitionBy(*partition_cols)
.parquet(path))
self.logger.info(f"Written {df.count()} records to {path}")
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.amazon.aws.operators.glue import GlueJobOperator
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
from datetime import datetime, timedelta
default_args = {
'owner': 'data-engineering',
'depends_on_past': False,
'start_date': datetime(2025, 1, 1),
'email_on_failure': True,
'email_on_retry': False,
'retries': 3,
'retry_delay': timedelta(minutes=5),
'retry_exponential_backoff': True,
}
with DAG(
'data_lake_bronze_to_silver',
default_args=default_args,
description='Process raw data to silver zone',
schedule_interval='0 2 * * *', # Daily at 2 AM
catchup=False,
tags=['data-lake', 'etl', 'silver'],
) as dag:
# Wait for source data
wait_for_data = S3KeySensor(
task_id='wait_for_bronze_data',
bucket_name='data-lake-bronze',
bucket_key='events/{{ ds }}/*.parquet',
timeout=3600,
poke_interval=300,
)
# Run Glue job
process_data = GlueJobOperator(
task_id='bronze_to_silver',
job_name='bronze-to-silver-etl',
script_args={
'--input_path': 's3://data-lake-bronze/events/{{ ds }}/',
'--output_path': 's3://data-lake-silver/events/',
'--partition_date': '{{ ds }}',
},
)
# Data quality checks
def validate_silver_data(**context):
"""Validate silver zone data quality."""
import boto3
athena = boto3.client('athena')
# Row count check
query = f"""
SELECT COUNT(*) as row_count
FROM silver.events
WHERE partition_date = '{context['ds']}'
"""
# Execute and validate (simplified)
# Add full implementation with quality checks
quality_check = PythonOperator(
task_id='data_quality_check',
python_callable=validate_silver_data,
)
wait_for_data >> process_data >> quality_check
services/python/
āāā data-ingestion/
ā āāā app/
ā ā āāā api/ # FastAPI routes
ā ā āāā services/ # Business logic
ā ā āāā models/ # Pydantic models
ā āāā tests/
ā āāā Dockerfile
ā āāā requirements.txt
āāā etl-pipeline/
ā āāā jobs/ # PySpark jobs
ā āāā glue/ # Glue-specific code
ā āāā tests/
ā āāā requirements.txt
āāā airflow/
āāā dags/ # Airflow DAGs
āāā plugins/ # Custom operators
āāā requirements.txt
For detailed information on specific topics:
def process_partition(partition_date: str, force_reprocess: bool = False):
"""Process data idempotently."""
output_path = f"s3://silver/{partition_date}/"
# Check if already processed
if not force_reprocess and path_exists(output_path):
logger.info(f"Partition {partition_date} already processed, skipping")
return
# Process data
df = read_bronze(partition_date)
df_clean = transform_to_silver(df)
# Atomic write (write to temp, then move)
temp_path = f"{output_path}_temp/"
df_clean.write.parquet(temp_path)
# Move to final location (atomic on S3)
move(temp_path, output_path)
def process_record_safe(record):
"""Process with error handling and DLQ."""
try:
# Validation
validated = validate_schema(record)
# Processing
result = transform(validated)
return ("success", result)
except ValidationError as e:
logger.warning(f"Validation failed: {e}")
return ("dlq", {"record": record, "error": str(e), "type": "validation"})
except Exception as e:
logger.error(f"Processing failed: {e}")
return ("dlq", {"record": record, "error": str(e), "type": "processing"})
# Process batch
results = df.rdd.map(process_record_safe).collect()
success_records = [r for status, r in results if status == "success"]
dlq_records = [r for status, r in results if status == "dlq"]
# Write DLQ records for investigation
if dlq_records:
write_to_dlq(dlq_records)
import logging
from pythonjsonlogger import jsonlogger
# Structured logging
logger = logging.getLogger()
logHandler = logging.StreamHandler()
formatter = jsonlogger.JsonFormatter(
'%(timestamp)s %(level)s %(name)s %(message)s %(correlation_id)s'
)
logHandler.setFormatter(formatter)
logger.addHandler(logHandler)
# Usage with context
def process_with_context(data_id: str):
with log_context(correlation_id=data_id):
logger.info("Processing started", extra={"record_count": len(data)})
# Process data
logger.info("Processing completed")
Status: Production-Ready Last Updated: 2025-11-04 Follows: Anthropic 500-line rule, progressive disclosure pattern