- Project Architecture
- Code Structure
- Extending the Project
- Adding New Analyses
- Custom Visualizations
- Performance Optimization
- Testing
- Deployment
┌─────────────────────────────────────────┐
│ Presentation Layer │
│ (Streamlit Web UI - app.py) │
└─────────────────────────────────────────┘
↓
┌─────────────────────────────────────────┐
│ Business Logic Layer │
│ (Analysis Functions - utils.py) │
└─────────────────────────────────────────┘
↓
┌─────────────────────────────────────────┐
│ Data Layer │
│ (PySpark Processing - traffic_...) │
└─────────────────────────────────────────┘
↓
┌─────────────────────────────────────────┐
│ Data Source │
│ (CSV - traffic_accidents_50000.csv) │
└─────────────────────────────────────────┘
Responsibilities:
- UI rendering and layout
- User interaction handling
- Data visualization display
- Session state management
Key Functions:
main() # Main application entry point
get_spark_session() # Get cached Spark session
load_and_cache_data() # Load and cache dataStructure:
- Page configuration
- Session state initialization
- Navigation routing
- Page implementations (6 pages)
Responsibilities:
- PySpark operations
- Data loading and cleaning
- Analysis execution
- Data conversion
Key Functions:
initialize_spark() # Create Spark session
load_data() # Load CSV
clean_data() # Clean dataset
analyze_data() # Run analysesResponsibilities:
- Centralized configuration
- Environment settings
- Performance tuning
Key Classes:
ApplicationConfig # App-level settings
SparkConfig # Spark optimizations
DataConfig # Data processing settings
VisualizationConfig # Plot configurationsResponsibilities:
- Batch processing
- Report generation
- Graph saving
Workflow:
- Initialize Spark
- Load data
- Clean data
- Analyze data
- Visualize results
- Generate reports
Step 1: Add analysis function in utils.py
def analyze_new_metric(df):
"""
Analyze new metric from dataset
Args:
df (DataFrame): Cleaned Spark DataFrame
Returns:
DataFrame: Analysis results
"""
results = df.groupBy('new_column') \
.agg(count('*').alias('count')) \
.orderBy(col('count').desc())
return resultsStep 2: Add to analysis pipeline in analyze_data()
def analyze_data(df):
analysis_results = {}
# ... existing analyses ...
# Add new analysis
analysis_results['new_metric'] = analyze_new_metric(df)
return analysis_resultsStep 3: Display in app.py
# In Analysis tab
new_data = analysis_results['new_metric'].toPandas()
st.dataframe(new_data, use_container_width=True)Step 1: Create visualization function
def visualize_custom_chart(df, analysis_results):
"""
Create custom visualization
"""
data = analysis_results['new_metric'].toPandas()
fig, ax = plt.subplots(figsize=Config.FIG_SIZE_LARGE)
# Your plotting code here
ax.plot(data['x'], data['y'], marker='o', linewidth=2)
ax.set_xlabel('X Label', fontsize=12, fontweight='bold')
ax.set_ylabel('Y Label', fontsize=12, fontweight='bold')
ax.set_title('Custom Chart Title', fontsize=14, fontweight='bold')
return figStep 2: Add to Streamlit app
elif viz_type == "Custom Chart":
st.subheader("📊 Custom Chart Title")
fig = visualize_custom_chart(df_clean, analysis_results)
st.pyplot(fig)
# Print insights
st.success("✓ Custom insight here")# In config.py - SparkConfig
DRIVER_MEMORY = "8g" # Increase for large operations
SQL_SHUFFLE_PARTITIONS = 400 # Increase for more parallelism# Cache frequently accessed DataFrames
df.cache()
df.count() # Force evaluation
# Reuse multiple times
df.groupBy(...).agg(...)
df.filter(...).count()
df.select(...).show()# Good - lazy evaluation
df_filtered = df.filter(col('city') == 'BANGALORE')
df_grouped = df_filtered.groupBy('weather').count()
# Bad - triggers evaluation prematurely
df_filtered.show() # Avoid here unless necessary# Repartition for large joins
df1 = df1.repartition(200, 'city')
df2 = df2.repartition(200, 'city')
result = df1.join(df2, 'city')# Filter early to reduce data size
df = df.filter(col('severity') != 'Unknown') # Do this first
df = df.filter(col('latitude') != 0) # Then other filters
# Then aggregations and groupingAdd to config.py:
class MetricsConfig:
"""Custom metrics configuration"""
# Define custom metrics
METRICS = {
'accident_density': 'accidents_per_1000_population',
'severity_index': 'weighted_severity_score',
'risk_factor': 'casualties_per_accident'
}# Add to utils.py
from pyspark.ml import Pipeline
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.clustering import KMeans
def predict_accident_clusters(df):
"""Predict accident clusters using KMeans"""
assembler = VectorAssembler(
inputCols=['latitude', 'longitude', 'hour'],
outputCol='features'
)
kmeans = KMeans(k=5, seed=42)
pipeline = Pipeline(stages=[assembler, kmeans])
model = pipeline.fit(df)
predictions = model.transform(df)
return predictionsCreate tests/test_utils.py:
import unittest
from utils import clean_data, analyze_data
class TestDataCleaning(unittest.TestCase):
def setUp(self):
"""Set up test fixtures"""
self.spark = initialize_spark()
self.sample_data = load_data(self.spark, 'test_data.csv')
def test_clean_data_removes_nulls(self):
"""Test that clean_data removes null values"""
df_clean = clean_data(self.sample_data)
null_count = df_clean.filter(col('city').isNull()).count()
self.assertEqual(null_count, 0)
def test_analyze_data_returns_dict(self):
"""Test that analyze_data returns dictionary"""
df_clean = clean_data(self.sample_data)
results = analyze_data(df_clean)
self.assertIsInstance(results, dict)
self.assertIn('accidents_per_city', results)
if __name__ == '__main__':
unittest.main()python -m pytest tests/ -v# Run Streamlit app
streamlit run app.pyCreate Dockerfile:
FROM python:3.9
WORKDIR /app
COPY requirements.txt .
RUN pip install -r requirements.txt
COPY . .
EXPOSE 8501
CMD ["streamlit", "run", "app.py"]Build and run:
docker build -t traffic-analysis .
docker run -p 8501:8501 traffic-analysis- Push code to GitHub
- Visit https://share.streamlit.io
- Connect repository
- Deploy app
# In config.py
class LoggingConfig:
LOG_LEVEL = 'DEBUG' # Change from INFO# In utils.py
def analyze_data(df):
# Print schema
print("DataFrame Schema:")
df.printSchema()
# Show sample data
print("Sample Data:")
df.show(10)
# Print execution plan
print("Execution Plan:")
df.groupBy('city').agg(count('*')).explain()# Run with debug mode
streamlit run app.py --logger.level=debug# Format with black
black *.py
# Check with flake8
flake8 *.pyfrom typing import Dict, List, Tuple
def analyze_data(df: DataFrame) -> Dict[str, DataFrame]:
"""Type hints for better IDE support"""
pass# Feature branch
git checkout -b feature/new-analysis
# Make changes
git add .
git commit -m "Add new analysis for weather impact"
# Push and create PR
git push origin feature/new-analysisvenv/
__pycache__/
*.pyc
.DS_Store
output_graphs/
logs/
.streamlit/secrets.toml
def analyze_data(df: DataFrame) -> Dict[str, DataFrame]:
"""
Perform comprehensive Big Data analysis on cleaned dataset
This function executes multiple analyses:
- City-wise accident distribution
- Time-based patterns
- Weather impact assessment
- Severity analysis
Args:
df (DataFrame): Cleaned Spark DataFrame with validated data
Returns:
dict: Analysis results with keys:
- 'accidents_per_city': Top 15 cities
- 'accidents_by_hour': Hourly distribution
- ... other analyses ...
Raises:
Exception: If DataFrame is empty or invalid
Example:
>>> spark = initialize_spark()
>>> df = load_data(spark, 'data.csv')
>>> df_clean = clean_data(df)
>>> results = analyze_data(df_clean)
>>> print(results['accidents_per_city'].show())
"""
passimport time
def timed_operation(func):
"""Decorator to measure execution time"""
def wrapper(*args, **kwargs):
start = time.time()
result = func(*args, **kwargs)
elapsed = time.time() - start
print(f"{func.__name__} took {elapsed:.2f} seconds")
return result
return wrapper
@timed_operation
def analyze_data(df):
return df.groupBy('city').count()Happy developing! 🚀