Why This Matters
A data lake on AWS is a centralized repository that stores structured, semi-structured, and unstructured data at any scale. Built primarily on Amazon S3, it provides a foundation for advanced analytics, machine learning, and real-time processing. Data lakes enable organizations to break down data silos and derive insights from diverse data sources. Understanding data lake patterns is critical for data engineers because the Medallion architecture, proper partitioning, and governance controls directly impact data quality, query performance, and cost efficiency. Mastering these patterns enables you to build scalable, secure data platforms that serve analysts, data scientists, and business applications.
Data Lake Architecture
Data Lake vs Data Warehouse
| Feature | Data Lake | Data Warehouse |
|---|---|---|
| Data Type | Raw, unprocessed | Cleaned, transformed |
| Schema | Schema-on-read | Schema-on-write |
| Storage Cost | Low (S3) | High (compute-optimized) |
| Query Performance | Variable | Optimized for BI |
| Users | Data Engineers, Scientists | Business Analysts |
| Format | Any format | Structured only |
Real-World Project Structure
aws-data-lake/
âââ infrastructure/
â âââ terraform/
â â âââ s3.tf # Buckets, lifecycle, versioning
â â âââ lake_formation.tf # Governance, permissions
â â âââ glue.tf # Crawlers, catalogs
â â âââ iam.tf # Roles and policies
â âââ cloudformation/
â âââ data-lake-stack.yaml # Full stack definition
âââ pipelines/
â âââ ingestion/
â â âââ bronze_ingest.py # Raw data ingestion
â â âââ streaming_ingest.py # Kinesis to Bronze
â âââ transformation/
â â âââ bronze_to_silver.py # Dedup, validate, cleanse
â â âââ silver_to_gold.py # Aggregate, model
â âââ orchestration/
â âââ dag_pipeline.py # Step Functions / Airflow
âââ governance/
â âââ lake_formation/
â â âââ permissions.json # Table/column permissions
â â âââ row_filters.json # Row-level security
â âââ classification/
â â âââ pii_detection.py # Auto-classify PII
â âââ audit/
â âââ access_logging.py # CloudTrail integration
âââ analytics/
â âââ athena/
â â âââ views/
â â â âââ daily_metrics.sql # Standard views
â â âââ queries/
â â âââ ad_hoc.sql # Ad-hoc analysis
â âââ redshift/
â â âââ spectrum_queries.sql # External table queries
â âââ quicksight/
â âââ dashboard_config.json # BI configurations
âââ tests/
â âââ data_quality/
â â âââ bronze_checks.sql # Raw data validation
â â âââ silver_checks.sql # Cleansed data validation
â â âââ gold_checks.sql # Aggregate validation
â âââ integration/
â âââ pipeline_test.py # End-to-end tests
âââ monitoring/
âââ cloudwatch/
â âââ alarms.json # Data freshness, quality
â âââ dashboards.json # Operational metrics
âââ alerts/
âââ notification_config.json # SNS alert routing
Bronze Layer Implementation
"""
AWS Glue ETL Job - Bronze Layer Ingestion
Ingests raw data from multiple sources into the Bronze layer.
"""
import sys
import logging
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from pyspark.sql.functions import current_timestamp, input_file_name, lit, col
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
args = getResolvedOptions(sys.argv, ['JOB_NAME', 'SOURCE_PATH', 'DATABASE_NAME'])
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
try:
logger.info(f"Starting Bronze ingestion from {args['SOURCE_PATH']}")
raw_df = spark.read.json(args['SOURCE_PATH'])
record_count = raw_df.count()
logger.info(f"Read {record_count} records from source")
if record_count == 0:
logger.warning("No records found. Skipping ingestion.")
sys.exit(0)
bronze_df = raw_df \
.withColumn("ingestion_timestamp", current_timestamp()) \
.withColumn("source_file", input_file_name()) \
.withColumn("batch_id", lit(f"batch_{current_timestamp().strftime('%Y%m%d_%H%M%S')}"))
bronze_df.write.mode("append") \
.partitionBy("event_date") \
.parquet("s3://data-lake/bronze/events/")
logger.info(f"Successfully ingested {record_count} records to Bronze layer")
except Exception as e:
logger.error(f"Bronze ingestion failed: {str(e)}")
raise
job.commit()
Silver Layer Processing
"""
Bronze to Silver transformation with deduplication and quality checks.
"""
from pyspark.sql.functions import col, row_number, when, count, sum as spark_sum
from pyspark.sql.window import Window
def process_silver(spark, bronze_path, silver_path):
"""Transform Bronze data to Silver with quality validation."""
try:
bronze_df = spark.read.parquet(bronze_path)
logger.info(f"Read {bronze_df.count()} records from Bronze")
# Deduplicate by primary key
window_spec = Window.partitionBy("event_id").orderBy(col("ingestion_timestamp").desc())
deduped_df = bronze_df \
.withColumn("row_num", row_number().over(window_spec)) \
.filter(col("row_num") == 1) \
.drop("row_num")
# Type casting and validation
validated_df = deduped_df \
.withColumn("event_id", col("event_id").cast("string")) \
.withColumn("user_id", col("user_id").cast("long")) \
.withColumn("amount", col("amount").cast("decimal(10,2)"))
# Data quality filters
valid_df = validated_df.filter(
col("event_id").isNotNull() &
col("user_id").isNotNull() &
(col("amount") >= 0) &
(col("amount") <= 1000000)
)
invalid_count = deduped_df.count() - valid_df.count()
if invalid_count > 0:
logger.warning(f"Filtered {invalid_count} invalid records")
valid_df.write.mode("overwrite") \
.partitionBy("event_date", "event_type") \
.parquet(silver_path)
logger.info(f"Successfully wrote {valid_df.count()} records to Silver layer")
except Exception as e:
logger.error(f"Silver processing failed: {str(e)}")
raise
Gold Layer Aggregation
"""
Silver to Gold transformation with business aggregations.
"""
from pyspark.sql.functions import sum as spark_sum, count, avg, date_trunc
def process_gold(spark, silver_path, gold_path):
"""Aggregate Silver data to Gold layer for BI consumption."""
try:
silver_df = spark.read.parquet(silver_path)
logger.info(f"Read {silver_df.count()} records from Silver")
# Daily revenue aggregation
daily_revenue = silver_df.groupBy("event_date", "product_category") \
.agg(
spark_sum("amount").alias("total_revenue"),
count("event_id").alias("transaction_count"),
avg("amount").alias("avg_order_value")
)
daily_revenue.write.mode("overwrite") \
.partitionBy("event_date") \
.parquet(f"{gold_path}/daily_revenue/")
# Dimension tables
product_dim = silver_df.select(
"product_id", "product_name", "product_category"
).distinct()
product_dim.write.mode("overwrite") \
.parquet(f"{gold_path}/dimensions/product_dim/")
logger.info("Successfully wrote Gold layer aggregations")
except Exception as e:
logger.error(f"Gold processing failed: {str(e)}")
raise
S3 Data Lake Configuration
#!/bin/bash
# Configure S3 data lake bucket with best practices
BUCKET_NAME="company-data-lake-prod"
REGION="us-east-1"
# Create bucket
aws s3api create-bucket \
--bucket "$BUCKET_NAME" \
--region "$REGION" \
--create-bucket-configuration LocationConstraint="$REGION"
# Enable versioning
aws s3api put-bucket-versioning \
--bucket "$BUCKET_NAME" \
--versioning-configuration Status=Enabled
# Enable default encryption
aws s3api put-bucket-encryption \
--bucket "$BUCKET_NAME" \
--server-side-encryption-configuration '{
"Rules": [
{
"ApplyServerSideEncryptionByDefault": {
"SSEAlgorithm": "aws:kms",
"KMSMasterKeyID": "arn:aws:kms:us-east-1:ACCOUNT:key/KEY_ID"
}
}
]
}'
# Lifecycle policies for cost optimization
aws s3api put-bucket-lifecycle-configuration \
--bucket "$BUCKET_NAME" \
--lifecycle-configuration '{
"Rules": [
{
"ID": "TieredStorage",
"Status": "Enabled",
"Transitions": [
{"Days": 90, "StorageClass": "STANDARD_IA"},
{"Days": 180, "StorageClass": "GLACIER"},
{"Days": 365, "StorageClass": "DEEP_ARCHIVE"}
]
}
]
}'
# Block public access
aws s3api put-public-access-block \
--bucket "$BUCKET_NAME" \
--public-access-block-configuration \
BlockPublicAcls=true,IgnorePublicAcls=true,BlockPublicPolicy=true,RestrictPublicBuckets=true
Mathematical Formulas
Storage Cost Optimization
Partition Pruning Benefit
Data Lake ROI
Performance Considerations
| Optimization | Impact | Implementation |
|---|---|---|
| Partitioning | 10-100x faster queries | Partition by date, region, business unit |
| Columnar Formats | 2-10x faster reads | Parquet/ORC with compression |
| File Size | Avoid small files problem | Target 256MB-1GB per file |
| Compaction | Reduce metadata overhead | Regular compaction jobs |
| Z-Ordering | Multi-dimensional queries | Order by frequently filtered columns |
| Bucketing | Faster joins | Hash high-cardinality columns |
| Predicate Pushdown | Less data scanned | Filter at source |
Security Considerations
| Layer | Controls | Implementation |
|---|---|---|
| Network | VPC endpoints, private subnets | S3 Gateway endpoint |
| Identity | IAM roles, Lake Formation | Table/column-level permissions |
| Data | SSE-KMS encryption | Customer-managed keys |
| Governance | Lake Formation tags | Tag-based access control |
| Audit | CloudTrail, access logging | Log all data access |
| PII | Macie detection, masking | Auto-classify and protect |
| Compliance | Config Rules | Continuous compliance checks |
Interview Questions & Answers
Q1: What is the Medallion Architecture and why use it?
Answer:
The Medallion Architecture is a multi-hop data processing pattern that organizes data into three layers: Bronze (raw ingestion), Silver (cleansed and conformed), and Gold (business-level aggregates).
Key benefits:
- Clear separation of concerns at each processing stage
- Data quality improves progressively through layers
- Each layer can be independently processed and tested
- Enables different consumption patterns (ad-hoc, reporting, ML)
- Provides lineage tracking and data governance
Q2: How do you optimize query performance in a data lake?
Answer:
- Partitioning: Partition by frequently filtered columns (date, region) to enable partition pruning
- File Format: Use columnar formats like Parquet or ORC for efficient compression and predicate pushdown
- File Size: Maintain optimal file sizes (256MB-1GB) to avoid small file problems
- Compaction: Regularly compact small files using AWS Glue or Spark
- Bucketing: Use bucketing for joins on high-cardinality columns
- Z-Ordering: Apply Z-ordering on multiple columns for multi-dimensional queries
Q3: Explain data lake governance with AWS Lake Formation.
Answer:
AWS Lake Formation provides centralized governance for data lakes:
- Permission Management: Granular permissions at database, table, and column levels
- Tag-based Access Control: Apply tags for simplified permission management
- Row-level Security: Filter rows based on user attributes
- Column-level Security: Hide sensitive columns from unauthorized users
- Audit Logging: Track all data access through CloudTrail integration
- Data Sharing: Secure cross-account data sharing without data movement
Q4: How do you handle schema evolution in a data lake?
Answer:
Schema evolution strategies include:
- Partition Evolution: Add new partitions with evolved schema
- Column Addition: Use Parquet's schema evolution to add nullable columns
- Schema Registry: Use Glue Schema Registry for schema versioning and compatibility checks
- Backward Compatibility: Maintain backward compatibility through schema validation
- Schema-on-Read: Leverage schema-on-read for flexible interpretation
Q5: What are the key S3 data lake best practices?
Answer:
- Bucket Structure: Organize by zone (bronze/silver/gold) and use prefixes for partitioning
- Lifecycle Policies: Implement tiered storage (Standard -> IA -> Glacier -> Deep Archive)
- Encryption: Enable default encryption with SSE-KMS
- Versioning: Enable versioning for data recovery and audit requirements
- Access Logging: Enable server access logging for compliance
- Replication: Configure cross-region replication for disaster recovery
- Object Lock: Use S3 Object Lock for WORM compliance requirements
Q6: How do you implement real-time analytics on a data lake?
Answer:
- Ingestion: Kinesis Data Streams or MSK for real-time data capture
- Processing: Kinesis Data Analytics or Flink for stream processing
- Storage: Write to S3 Bronze layer in near real-time
- Query: Use Athena for ad-hoc queries or Redshift for complex analytics
- Visualization: QuickSight SPICE for real-time dashboards
Q7: How do you monitor and troubleshoot data lake performance?
Answer:
- CloudWatch Metrics: Track S3 request rates, latency, and error rates
- Glue Job Monitoring: Monitor ETL job metrics, duration, and failures
- Query Performance: Use Athena query metrics and CloudWatch dashboards
- Data Quality: Implement data quality checks with Glue DataBrew
- Cost Monitoring: Use Cost Explorer and Budgets for storage cost tracking
- Alerting: Set up CloudWatch Alarms for critical metrics
Q8: Explain data lake security best practices.
Answer:
- Network Security: VPC endpoints for S3 access, private subnets for compute
- IAM Policies: Least privilege access with resource-based policies
- Encryption: At-rest (KMS) and in-transit (TLS 1.2+) encryption
- Lake Formation: Column and row-level security for fine-grained access
- Macie: Automated PII detection and classification
- Security Hub: Centralized security findings and compliance
- Audit: CloudTrail for API activity logging
Common Pitfalls
| Pitfall | Impact | Solution |
|---|---|---|
| Small files problem | Slow queries, high metadata overhead | Compact to 256MB-1GB files |
| No partitioning | Full table scans, high costs | Partition by date/region |
| Missing governance | Data exposure, compliance violations | Implement Lake Formation |
| No lifecycle policies | Unnecessary storage costs | Tiered storage classes |
| Ignoring data quality | Bad data in, bad insights out | Quality checks at each layer |
| Over-permissive IAM | Security vulnerabilities | Least privilege, resource constraints |
| No versioning | Data loss, no audit trail | Enable S3 versioning |
| Ignoring compaction | Degraded query performance | Regular compaction jobs |