← Load balancing: traffic and resilience
11 / 12 · 60 MIN

Affinity, membership and state continuity

Experiment with consistent hashing under membership changes and distinguish routing stability, state presence and functional recovery.

Run and bound the experiment

Save the code below as run.py, point DR_NGINX_BIN to a local NGINX executable and run python3 run.py --output evidence.json. Recorded execution uses NGINX 1.30.5 and Python 3.13.1, one worker and three HTTP backends on IPv4 loopback. The script creates temporary ports, checks each configuration and closes the resources it created. It uses no external services and performs no financial operations. X-Lab-Key is a client-controlled teaching header. It is not a proposal for production authentication or session management.

Compare the same population

The script sends 96 keys to an A/B group and repeats the exact population. With membership and configuration stable, each key returns to its previous target. It records both complete maps for comparison instead of storing only a green indicator. Eight additional requests with the same key reach the same backend. These results show affinity in the tested configuration. They identify no users and establish no uniform cost distribution: one key may produce many requests while another generates only an occasional request.

Observe C joining

The next configuration adds C, and the script sends HUP to the master it created. It begins new sampling only after observing the expected generation header and confirming the previous worker exited. Some keys move to C while the others retain their previous targets in this experiment. The exact count depends on the keys and addresses used in the ring, including temporary ports. Repetition therefore validates relationships between maps without demanding an identical percentage. Sizing a service also requires measuring request frequency, duration and cost per key.

Separate target and data

The experiment restores A/B, selects a key mapped to A and writes a synthetic value into that backend’s local dictionary. It then removes A from selection. B responds to the same request with HTTP 200, but the value is absent; the original copy remains on A. This sequence is useful for an APS incident where a portal appears available while a user loses access to a draft. Before declaring deletion, locate the state and its valid version. The transport response alone does not establish journey continuity.

Interpret the teaching copy

One fixture line explicitly copies the value from A to B, and the next read finds it. The load balancer did not perform that transfer. The example changes one variable at a time: first the target, then the presence of the value. There is no persistence, concurrent replication, session expiration or conflict resolution. For a real solution, the team must define state authority, how a new target accesses it and what happens when that dependency fails. Restoring A can also reintroduce a stale copy.

Bring evidence into the change

For a fictional banking infrastructure project, add an acceptance journey that creates state, changes target and confirms the expected outcome. Record configuration, observed generations, members and approved correlation identifiers without copying session tokens into tickets. The project manager confirms owners and criteria with APS and development; RUN must be able to repeat the diagnosis. The sentence “Routing checks passed; state continuity still needs representative validation” communicates the present limit in English. These examples do not describe internal BNP Paribas procedures.

#!/usr/bin/env python3
"""Original bounded loopback fixture. No external services or financial operations."""
import argparse
from collections import Counter
from concurrent.futures import ThreadPoolExecutor
import hashlib
from http.client import HTTPConnection
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
import json
import os
from pathlib import Path
import platform
import re
import signal
import socket
import subprocess
import tempfile
import threading
import time


def digest(path):
 return hashlib.sha256(Path(path).read_bytes).hexdigest


def run(binary):
 version = subprocess.run([binary, '-V'], capture_output=True, text=True, check=True).stderr
 started, release, completed = threading.Event, threading.Event, threading.Event
 events, stores = [], {name: {} for name in 'ABC'}
 lock = threading.Lock
 servers, threads, checks, generations = [], [], {}, []
 master = None
 evidence = dict(scriptSha256=digest(__file__), binarySha256=digest(binary), nginxVersion=version,
 pythonVersion=platform.python_version, platform=platform.platform,
 scope='One NGINX worker, three synthetic HTTP loopback backends, 96 selected keys and explicit in-memory fixture state; not a throughput benchmark, durable session store, identity system, production drain or business cancellation test.')
 tmp = tempfile.TemporaryDirectory(prefix='dr-lb-affinity-')
 root = Path(tmp.name)
 pool = ThreadPoolExecutor(max_workers=1)

 def handler(name):
 class Handler(BaseHTTPRequestHandler):
 protocol_version = 'HTTP/1.1'

 def log_message(self, *args):
 pass

 def do_POST(self):
 self.do_GET

 def do_GET(self):
 key = self.headers.get('X-Lab-Key', '')
 with lock:
 events.append(dict(backend=name, path=self.path, key=key, method=self.command))
 if self.path == '/hold' and name == 'A':
 started.set
 if not release.wait(12):
 return
 # Deliberately independent of socket liveness. No production side effects.
 completed.set
 if self.path == '/session':
 if self.command == 'POST':
 value = self.rfile.read(int(self.headers.get('Content-Length', '0'))).decode
 with lock:
 stores[name][key] = value
 with lock:
 value = stores[name].get(key)
 else:
 value = None
 data = json.dumps(dict(backend=name, value=value)).encode
 try:
 self.send_response(200)
 self.send_header('Content-Length', str(len(data)))
 self.send_header('Connection', 'close')
 self.end_headers
 self.wfile.write(data)
 except (BrokenPipeError, ConnectionResetError):
 pass
 self.close_connection = True
 return Handler

 try:
 for name in 'ABC':
 server = ThreadingHTTPServer(('127.0.0.1', 0), handler(name))
 server.daemon_threads = False
 thread = threading.Thread(target=server.serve_forever, kwargs={'poll_interval':.05})
 thread.start
 servers.append(server)
 threads.append(thread)
 ports = dict(zip('ABC', [s.server_port for s in servers]))
 with socket.socket as reserve:
 reserve.bind(('127.0.0.1', 0))
 port = reserve.getsockname[1]
 conf = root / 'nginx.conf'

 def configure(members, hold, generation):
 lines = '\n'.join(f' server 127.0.0.1:{ports[n]};' for n in members)
 config = f'''daemon off;
worker_processes 1;
worker_shutdown_timeout 1s;
pid {root}/master.pid;
error_log {root}/error.log notice;
events {{ worker_connections 256; }}
http {{
 log_format lab '$request_uri|$status|$upstream_status|$upstream_addr'
 access_log {root}/access.log lab;
 upstream affinity {{ hash $http_x_lab_key consistent; {lines} keepalive 0; }}
 upstream held {{ server 127.0.0.1:{ports[hold]}; keepalive 0; }}
 server {{
 listen 127.0.0.1:{port};
 add_header X-Lab-Generation {generation} always;
 proxy_http_version 1.1;
 proxy_set_header Connection close;
 proxy_connect_timeout 2s;
 proxy_read_timeout 10s;
 proxy_next_upstream off;
 location / {{ proxy_pass http://affinity; }}
 location = /hold {{ proxy_pass http://held; }}
 }}
}}
'''
 conf.write_text(config)
 result = subprocess.run([binary, '-p', str(root) + '/', '-c', str(conf), '-t'], capture_output=True, text=True)
 assert result.returncode == 0, result.stderr
 generations.append(dict(generation=generation, members=list(members), hold=hold, configuration=config))

 def request(path='/', key='probe', method='GET', body=None):
 conn = HTTPConnection('127.0.0.1', port, timeout=6)
 try:
 conn.request(method, path, body=body, headers={'Connection': 'close', 'X-Lab-Key': key})
 res = conn.getresponse
 data = res.read
 return dict(status=res.status, generation=res.getheader('X-Lab-Generation'), **json.loads(data))
 except Exception as exc:
 return dict(error=type(exc).__name__)
 finally:
 conn.close

 def wait_generation(generation):
 deadline = time.monotonic + 5
 while time.monotonic < deadline:
 r = request('/ready')
 if r.get('generation') == generation:
 return r
 time.sleep(.02)
 raise AssertionError('new configuration generation not observed')

 def reload(members, hold, generation, settle=True):
 prior = re.findall(r'start worker process (\d+)', (root / 'error.log').read_text)[-1]
 configure(members, hold, generation)
 master.send_signal(signal.SIGHUP)
 result = wait_generation(generation)
 # Observing a new worker does not prove the previous one has stopped accepting.
 # Settle membership experiments; preserve overlap deliberately for the deadline test.
 if settle:
 deadline = time.monotonic + 5
 marker = f'worker process {prior} exited with code 0'
 while marker not in (root / 'error.log').read_text:
 assert time.monotonic < deadline, 'previous worker did not exit'
 time.sleep(.02)
 return result

 configure('AB', 'A', 'ab')
 master = subprocess.Popen([binary, '-p', str(root) + '/', '-c', str(conf)], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
 wait_generation('ab')
 keys = [f'key-{i:03}' for i in range(96)]
 before = {k: request(key=k)['backend'] for k in keys}
 repeat = {k: request(key=k)['backend'] for k in keys}
 checks['stableMembership'] = dict(keys=96, identicalMappings=before == repeat)
 checks['sampleCoverage'] = dict(backends=sorted(set(before.values)))
 # Same routing key does not represent separate authenticated identities.
 shared = [request(key=keys[0])['backend'] for _ in range(8)]
 checks['sharedKey'] = dict(requests=8, distinctBackends=len(set(shared)), authenticationImplemented=False)
 reload('ABC', 'A', 'abc')
 after = {k: request(key=k)['backend'] for k in keys}
 moved = [k for k in keys if before[k]!= after[k]]
 checks['membershipRemap'] = dict(someMoved=0 < len(moved) < 96, movedOnlyToAddedBackend=all(after[k] == 'C' for k in moved))
 checks['unmovedOwners'] = dict(retainedOldOwner=all(before[k] == after[k] for k in keys if after[k]!= 'C'), allThreeObserved=set(after.values) == set('ABC'))
 reload('AB', 'A', 'session-ab')
 key = next(k for k in keys if before[k] == 'A')
 written = request('/session', key, 'POST', 'synthetic-session-v1')
 read = request('/session', key)
 checks['localState'] = dict(writeBackend=written['backend'], readBackend=read['backend'], valuePresent=read['value'] == 'synthetic-session-v1')
 reload('B', 'A', 'session-b')
 missing = request('/session', key)
 checks['remappedState'] = dict(backend=missing['backend'], valueMissing=missing['value'] is None, originalStillInA=stores['A'][key] == 'synthetic-session-v1')
 # Explicit test-fixture copy, not a replication capability of NGINX.
 with lock:
 stores['B'][key] = stores['A'][key]
 recovered = request('/session', key)
 checks['explicitFixtureCopy'] = dict(backend=recovered['backend'], valuePresent=recovered['value'] == 'synthetic-session-v1', nginxReplication=False)
 future = pool.submit(request, '/hold', 'held-work')
 assert started.wait(3), 'backend A did not start held request'
 reload('B', 'B', 'deadline-b', settle=False)
 new = request('/hold', 'new-work')
 checks['newTraffic'] = dict(backend=new['backend'], status=new['status'], oldClientPending=not future.done)
 old = future.result(timeout=5)
 checks['shutdownDeadline'] = dict(clientFailed='error' in old, clientError=old.get('error'), backendReleased=release.is_set)
 checks['backendPending'] = dict(started=started.is_set, completedBeforeRelease=completed.is_set)
 release.set
 checks['completionAfterDisconnect'] = dict(completedAfterRelease=completed.wait(3), financialSideEffects=False)
 assert checks['stableMembership']['identicalMappings']
 assert checks['sampleCoverage']['backends'] == ['A', 'B']
 assert checks['sharedKey']['distinctBackends'] == 1
 assert all(checks['membershipRemap'].values) and all(checks['unmovedOwners'].values)
 assert checks['localState'] == dict(writeBackend='A', readBackend='A', valuePresent=True)
 assert checks['remappedState'] == dict(backend='B', valueMissing=True, originalStillInA=True)
 assert checks['explicitFixtureCopy'] == dict(backend='B', valuePresent=True, nginxReplication=False)
 assert checks['newTraffic'] == dict(backend='B', status=200, oldClientPending=True)
 assert checks['shutdownDeadline']['clientFailed']
 assert checks['shutdownDeadline']['backendReleased'] is False
 assert checks['backendPending'] == dict(started=True, completedBeforeRelease=False)
 assert checks['completionAfterDisconnect']['completedAfterRelease']
 evidence.update(checks=checks, passed=len(checks), failed=0, mappingBefore=before, mappingAfter=after,
 remappedCount=len(moved), mappingCountsBefore=dict(Counter(before.values)),
 mappingCountsAfter=dict(Counter(after.values)), generations=generations,
 backendEvents=events, clientAfterDeadline=old)
 finally:
 release.set
 pool.shutdown(wait=True)
 if master is not None and master.poll is None:
 master.send_signal(signal.SIGQUIT)
 try:
 master.wait(timeout=5)
 except subprocess.TimeoutExpired:
 master.terminate
 master.wait(timeout=3)
 for server in servers:
 server.shutdown
 server.server_close
 for thread in threads:
 thread.join(timeout=3)
 evidence['childExited'] = master is not None and master.poll is not None
 evidence['backendThreadsStopped'] = all(not thread.is_alive for thread in threads)
 evidence['errorLog'] = (root / 'error.log').read_text if (root / 'error.log').exists else ''
 evidence['accessLog'] = (root / 'access.log').read_text.splitlines if (root / 'access.log').exists else []
 tmp.cleanup
 evidence['temporaryDirectoryRemoved'] = not root.exists
 return evidence


if __name__ == '__main__':
 parser = argparse.ArgumentParser
 parser.add_argument('--output', required=True)
 args = parser.parse_args
 data = run(os.environ['DR_NGINX_BIN'])
 Path(args.output).write_text(json.dumps(data, indent=2) + '\n')
 print(json.dumps(dict(passed=data['passed'], failed=data['failed'], remappedCount=data['remappedCount'], checks=data['checks'])))
IN PRACTICE

A fictional portal adds C: some sessions change target and a 200 hides missing drafts. The team compares maps and reads state before deciding recovery.

Common pitfalls

Confusing hashing with authentication, key balance with load balance or a teaching copy with durable replication; declaring recovery solely because a target returned.

Related topics: Algorithms, affinity, and state · Retries, limits, and client origin · Draining, releases, and operations

Take this idea with you

Affinity selects a target. Continuity requires that target to obtain valid state and complete the expected journey, including after membership changes.

Create account

Reference: NGINX upstream hashing and membership · BigSavant load balancing 2026-09; selected NGINX, HAProxy 3.2, Kubernetes and AWS ALB behavior