| name | Bulkhead Pattern |
| description | Resilience pattern isolating system components to prevent cascade failures and enable partial degradation |
| category | software-development |
Bulkhead Pattern
What I do
I provide a pattern that isolates system components to prevent failures in one part from affecting the entire system. Named after ship bulkheads that contain flooding, the bulkhead pattern allocates separate resources (threads, connections, capacity) to different components. If one component fails or becomes slow, it doesn't exhaust resources for other components, enabling graceful degradation and improved system resilience.
When to use me
Use bulkheads when you have multiple services or components that share common resources, when you need to prioritize certain operations over others, or when preventing cascade failures is critical. They're essential in high-traffic systems, microservices architectures, and applications with varying priority operations. Bulkheads are valuable when different operations have different resource requirements or failure tolerance.
Core Concepts
- Isolation Levels: Complete, pooled, or hybrid isolation
- Thread Pool: Separate pools for different operations
- Connection Pooling: Dedicated database connections per service
- Rate Limiting: Controlling request rates per component
- Priority Queuing: Processing high-priority first
- Resource Partitioning: Dividing resources between components
- Semaphore: Controlling concurrent access
- Failure Containment: Containing failures within partitions
- Graceful Degradation: Partial functionality during failures
- Resource Exhaustion: Preventing total system failure
Code Examples
Thread Pool Bulkhead
import queue
import threading
import time
from abc import ABC, abstractmethod
from dataclasses import dataclass
from typing import Callable, Generic, TypeVar
from uuid import uuid4
from enum import Enum
T = TypeVar("T")
class TaskPriority(Enum):
HIGH = 0
MEDIUM = 1
LOW = 2
@dataclass
class Task(Generic[T]):
id: str
func: Callable[..., T]
args: tuple
kwargs: dict
priority: TaskPriority
created_at: float = time.time()
class PriorityThreadPoolExecutor:
def __init__(
self,
max_workers: int,
max_high_priority: int = 10,
max_medium_priority: int = 20,
max_low_priority: int = 50
):
self._max_workers = max_workers
self._queues: dict[TaskPriority, queue.PriorityQueue] = {
TaskPriority.HIGH: queue.PriorityQueue(maxsize=max_high_priority),
TaskPriority.MEDIUM: queue.PriorityQueue(maxsize=max_medium_priority),
TaskPriority.LOW: queue.PriorityQueue(maxsize=max_low_priority)
}
._workers: [threading.Thread] = []
._running =
._metrics = {
: ,
: ,
:
}
() -> :
task_id = (uuid4())
task = Task(
=task_id,
func=func,
args=args,
kwargs=kwargs,
priority=priority
)
q = ._queues[priority]
q.full():
._metrics[] +=
QueueFullError()
q.put((priority.value, task))
._metrics[] +=
task_id
() -> :
_ (._max_workers):
worker = threading.Thread(target=._worker_loop)
worker.daemon =
worker.start()
._workers.append(worker)
() -> :
._running:
priority TaskPriority:
q = ._queues[priority]
:
_, task = q.get(timeout=)
result = task.func(*task.args, **task.k kwargs)
._metrics[] +=
queue.Empty:
Exception e:
()
() -> :
._running =
q ._queues.values():
q.join()
() -> :
._metrics.copy()
():
() -> :
time.sleep()
{: order_id, : }
():
executor = PriorityThreadPoolExecutor(
max_workers=,
max_high_priority=,
max_medium_priority=,
max_low_priority=
)
executor.start()
:
high_id = executor.submit(
process_order,
,
priority=TaskPriority.HIGH
)
low_id = executor.submit(
process_order,
,
priority=TaskPriority.LOW
)
()
:
time.sleep()
executor.shutdown()
()
Connection Pool Bulkhead
import threading
import queue
from abc import ABC, abstractmethod
from dataclasses import dataclass
from typing import Protocol, Generic, TypeVar
from contextlib import contextmanager
from uuid import uuid4
T = TypeVar("T", bound='PooledConnection')
@dataclass
class ConnectionConfig:
max_connections: int
min_idle: int = 0
connection_timeout: float = 30.0
idle_timeout: float = 600.0
class PooledConnection(Protocol):
@property
def id(self) -> str:
pass
def is_healthy(self) -> bool:
pass
def execute(self, query: str) -> dict:
pass
def close(self) -> None:
pass
class ([T]):
():
._factory = factory
._config = config
._pool: queue.Queue[T] = queue.Queue()
._in_use: [, T] = {}
._lock = threading.Lock()
._metrics = {
: ,
: ,
:
}
():
conn = ._get_connection()
._in_use[conn.] = conn
._metrics[] +=
:
conn
:
._return_connection(conn)
._in_use[conn.]
._metrics[] -=
() -> T:
:
conn = ._pool.get_nowait()
conn.is_healthy():
conn
conn.close()
queue.Empty:
._lock:
._metrics[] < ._config.max_connections:
conn = ._factory()
._metrics[] +=
conn
conn = ._pool.get(timeout=._config.connection_timeout)
conn
() -> :
conn.is_healthy():
:
._pool.put_nowait(conn)
queue.Full:
conn.close()
:
conn.close()
() -> :
_ ((count, ._config.max_connections)):
conn = ._factory()
._pool.put(conn)
._metrics[] +=
() -> :
._pool.empty():
:
conn = ._pool.get_nowait()
conn.close()
queue.Empty:
conn ._in_use.values():
conn.close()
() -> :
{
**._metrics,
: ._pool.qsize(),
: (._in_use)
}
:
():
. = (uuid4())
.db_name = db_name
._healthy =
() -> :
._healthy
() -> :
{: query, : }
() -> :
._healthy =
pool = ConnectionPool(
factory=: DatabaseConnection(),
config=ConnectionConfig(max_connections=, min_idle=)
)
pool.prewarm()
pool.acquire() conn:
result = conn.execute()
()
()
Semaphore-Based Bulkhead
import asyncio
import threading
from typing import Callable
from dataclasses import dataclass
from datetime import datetime
@dataclass
class BulkheadConfig:
max_concurrent_calls: int
max_waiting: int = 0
timeout_seconds: float = 30.0
class ThreadBulkhead:
def __init__(self, config: BulkheadConfig):
self._semaphore = threading.BoundedSemaphore(config.max_concurrent_calls)
self._waiting = 0
self._lock = threading.Lock()
self._metrics = {
"total_calls": 0,
"rejected_calls": 0,
"active_calls": 0
}
def execute(self, func: Callable, *args, **kwargs):
with self._lock:
self._metrics["total_calls"]
acquired = self._semaphore.acquire(timeout=30)
if not acquired:
with self._lock:
._metrics[] +=
BulkheadRejectedError()
:
._lock:
._metrics[] +=
func(*args, **kwargs)
:
._lock:
._metrics[] -=
._semaphore.release()
() -> :
._metrics.copy()
:
():
._semaphore = asyncio.Semaphore(config.max_concurrent_calls)
._waiting =
._metrics = {
: ,
: ,
:
}
():
._metrics[] +=
:
asyncio.timeout():
._semaphore.acquire()
._metrics[] +=
:
func(*args, **kwargs)
:
._metrics[] -=
._semaphore.release()
asyncio.TimeoutError:
._metrics[] +=
BulkheadRejectedError()
():
() -> :
time
time.sleep()
{: service_name, : }
bulkhead = ThreadBulkhead(BulkheadConfig(max_concurrent_calls=))
results = []
i ():
:
result = bulkhead.execute(external_api_call, )
results.append(result)
BulkheadRejectedError:
results.append()
()
Service Partitioning
from abc import ABC, abstractmethod
from dataclasses import dataclass
from typing import Protocol
@dataclass
class ServicePartition:
name: str
max_concurrent: int
max_queue: int
priority: int
class ServiceRouter(ABC):
@abstractmethod
def get_partition(self, request: 'Request') -> 'ServicePartition':
pass
class TieredServiceBulkhead:
def __init__(self, partitions: list[ServicePartition]):
self._partitions = {
p.name: {
"config": p,
"semaphore": asyncio.Semaphore(p.max_concurrent) if hasattr(asyncio, 'Semaphore') else None,
"queue": queue.Queue(p.max_queue),
"metrics": {
"calls": 0,
"rejected": 0,
"executed": 0
}
}
for p in partitions
}
() -> :
._partitions.get(partition_name, {})
() -> :
partition = router.get_partition(request)
partition_info = ._partitions.get(partition.name)
partition_info:
PartitionResult(status=)
:
partition_info[].put_nowait(request)
PartitionResult(status=)
queue.Full:
partition_info[][] +=
PartitionResult(status=, reason=)
() -> :
{
name: info[]
name, info ._partitions.items()
}
:
():
. =
.priority = priority
.data = data
:
():
.status = status
.reason = reason
partitions = [
ServicePartition(, max_concurrent=, max_queue=, priority=),
ServicePartition(, max_concurrent=, max_queue=, priority=),
ServicePartition(, max_concurrent=, max_queue=, priority=)
]
bulkhead = TieredServiceBulkhead(partitions)
Adaptive Bulkhead
import asyncio
import time
from dataclasses import dataclass
from typing import Protocol
@dataclass
class AdaptiveConfig:
initial_capacity: int
min_capacity: int
max_capacity: int
scale_up_threshold: float
scale_down_threshold: float
scale_interval_seconds: float
class MetricsCollector(Protocol):
def get_latency_percentile(self, percentile: float) -> float:
pass
def get_error_rate(self) -> float:
pass
class AdaptiveBulkhead:
def __init__(self, config: AdaptiveConfig, metrics: MetricsCollector):
self._config = config
self._metrics = metrics
self._current_capacity = config.initial_capacity
self._semaphore = asyncio.Semaphore(self._current_capacity)
self._scale_lock = asyncio.Lock()
@property
def capacity(self) -> int:
return ._current_capacity
():
._scale_lock:
._maybe_scale()
._semaphore.acquire()
:
func(*args, **kwargs)
:
._semaphore.release()
() -> :
latency = ._metrics.get_latency_percentile()
error_rate = ._metrics.get_error_rate()
error_rate > ._config.scale_up_threshold latency > :
._current_capacity < ._config.max_capacity:
._current_capacity +=
new_semaphore = asyncio.Semaphore(._current_capacity)
_ (._current_capacity - ):
new_semaphore.acquire()
._semaphore = new_semaphore
()
latency < error_rate < :
._current_capacity > ._config.min_capacity:
._current_capacity -=
._semaphore = asyncio.Semaphore(._current_capacity)
()
() -> :
:
asyncio.sleep(._config.scale_interval_seconds)
._scale_lock:
._maybe_scale()
Circuit Breaker Integration
from enum import Enum
from dataclasses import dataclass
class HealthStatus(Enum):
HEALTHY = "healthy"
DEGRADED = "degraded"
CRITICAL = "critical"
@dataclass
class HealthCheckResult:
status: HealthStatus
latency_ms: float
error_rate: float
class IsolatedComponent:
def __init__(
self,
name: str,
max_concurrent: int,
circuit_breaker: CircuitBreaker
):
self.name = name
self._semaphore = asyncio.Semaphore(max_concurrent)
self._circuit_breaker = circuit_breaker
self._health_status = HealthStatus.HEALTHY
self._metrics = {"calls": 0, "failures": 0, "latencies": []}
async def execute(self, func: Callable, *args, **kwargs):
self._metrics["calls"] += 1
if self._circuit_breaker.state == CircuitState.OPEN:
raise ComponentIsolatedError(self.name)
start = time.time()
:
._semaphore:
result = func(*args, **kwargs)
result
Exception e:
._metrics[] +=
:
latency = (time.time() - start) *
._metrics[].append(latency)
() -> HealthCheckResult:
._metrics[]:
HealthCheckResult(HealthStatus.HEALTHY, , )
avg_latency = (._metrics[]) / (._metrics[])
error_rate = ._metrics[] / ._metrics[]
error_rate > avg_latency > :
status = HealthStatus.CRITICAL
error_rate > avg_latency > :
status = HealthStatus.DEGRADED
:
status = HealthStatus.HEALTHY
HealthCheckResult(status, avg_latency, error_rate)
:
():
._components: [, IsolatedComponent] = {}
() -> :
._components[component.name] = component
() -> IsolatedComponent:
._components[name]
() -> [, HealthCheckResult]:
{
name: component.get_health()
name, component ._components.items()
}
() -> :
name ._components:
._components[name]._health_status = HealthStatus.CRITICAL
():
():
.component_name = component_name
().__init__()
Best Practices
- Size Pools Appropriately: Don't over-allocate resources
- Monitor Actively: Track utilization and saturation
- Prioritize Work: Use priorities for critical operations
- Combine with Circuit Breakers: Layer multiple resilience patterns
- Graceful Backpressure: Reject requests when at capacity
- Separate Failure Domains: Isolate different services/components
- Testing: Test behavior under overload
- Alerting: Notify when pools are saturated
- Graceful Degradation: Reduce functionality when constrained
- Fair Distribution: Prevent single operation from hogging resources
- Dynamic Adjustment: Adapt to changing load patterns
- Documentation: Document partition strategies and limits