Консистентность в распределенных системах
В распределенной системе данные существуют на нескольких узлах. Консистентность определяет, когда изменение на одном узле становится видимым на других.
CAP-теорема
В распределенной системе при сетевом разделении (Partition) можно гарантировать только одно из двух: Consistency или Availability.
| Выбор | Поведение при partition | Примеры |
|---|---|---|
| CP (Consistency + Partition) | Отказ от обслуживания | PostgreSQL, MongoDB (default) |
| AP (Availability + Partition) | Обслуживание с возможно устаревшими данными | Cassandra, DynamoDB |
Важно: CAP -- это не выбор "одного из трех". Это выбор между C и A при наличии P (сетевых проблем). В нормальном режиме можно иметь и C, и A.
Strong Consistency (Строгая консистентность)
Каждое чтение возвращает результат последней записи. Все узлы видят одинаковые данные в один момент времени.
<?php
declare(strict_types=1);
/**
* Strong consistency: all reads see the latest write
*
* Implementation: synchronous replication
* Write is not considered successful until ALL replicas confirm
*
* Use when: financial transactions, inventory, bookings
*/
final class StrongConsistencyRepository
{
public function __construct(
private readonly \PDO $primary,
/** @var array<\PDO> */
private readonly array $replicas,
) {}
/**
* Synchronous write: wait for all replicas to confirm.
* PostgreSQL supports this with synchronous_commit and
* synchronous_standby_names configuration.
*
* Application-level equivalent:
*/
public function writeWithConfirmation(string $sql, array $params): void
{
// Write to primary
$this->primary->prepare($sql)->execute($params);
// Verify data reached replicas (synchronous replication)
// In PostgreSQL, this is handled at DB config level:
// synchronous_commit = on
// synchronous_standby_names = 'replica1,replica2'
// Application-level verification (if needed):
$maxWait = 5000; // 5 seconds
$waited = 0;
foreach ($this->replicas as $replica) {
while ($waited < $maxWait) {
$stmt = $replica->prepare(
'SELECT pg_last_wal_replay_lsn() >= pg_current_wal_lsn() AS synced'
);
$stmt->execute();
if ($stmt->fetchColumn()) {
break;
}
usleep(100_000); // 100ms
$waited += 100;
}
}
}
}
Eventual Consistency (Конечная консистентность)
Если новые записи прекратятся, со временем все узлы придут к одинаковому состоянию. Но в промежутке чтения могут возвращать устаревшие данные.
<?php
declare(strict_types=1);
/**
* Eventual consistency: reads may return stale data temporarily
*
* Common in: social feeds, analytics, notifications, search indexes
*
* The key question: what's the acceptable staleness window?
* - Social feed: seconds to minutes (OK)
* - Search results: minutes to hours (OK)
* - Account balance: NOT acceptable (use strong consistency)
*/
// Read-your-writes consistency: compromise between strong and eventual
final class ReadYourWritesRepository
{
public function __construct(
private readonly \PDO $primary,
private readonly \PDO $replica,
private readonly \Redis $cache,
) {}
public function write(string $userId, string $sql, array $params): void
{
$this->primary->prepare($sql)->execute($params);
// Record that this user just wrote something
$this->cache->setex("recent_write:$userId", 10, '1');
}
/**
* Read with "read-your-writes" guarantee:
* If this user recently wrote, read from primary.
* Otherwise, read from replica (faster, closer).
*/
public function read(string $userId, string $sql, array $params): array
{
$recentWrite = $this->cache->get("recent_write:$userId");
$db = $recentWrite !== false
? $this->primary // Recent write: use primary for consistency
: $this->replica; // No recent write: use replica for performance
$stmt = $db->prepare($sql);
$stmt->execute($params);
return $stmt->fetchAll(\PDO::FETCH_ASSOC);
}
}
Идемпотентность
Идемпотентная операция дает одинаковый результат при повторном выполнении. Это критично для надежных распределенных систем, где сетевые ошибки приводят к повторным запросам.
Idempotency Keys
<?php
declare(strict_types=1);
/**
* Idempotency key pattern
*
* Problem: Network failure after server processed request but before
* client received response. Client retries. Without idempotency,
* the operation executes twice (double charge, double order).
*
* Solution: Client sends a unique key with each request.
* Server stores the result and returns it on retry.
*/
final class IdempotentPaymentService
{
public function __construct(
private readonly \PDO $db,
private readonly PaymentGateway $gateway,
) {}
/**
* Process payment idempotently.
* Same idempotency key always returns the same result.
*/
public function processPayment(PaymentRequest $request): PaymentResult
{
// Step 1: Check if we already processed this request
$existing = $this->findByIdempotencyKey($request->idempotencyKey);
if ($existing !== null) {
// Already processed: return stored result
return $existing;
}
// Step 2: Lock the idempotency key to prevent concurrent processing
$this->db->beginTransaction();
try {
// Insert idempotency record with status "processing"
// UNIQUE constraint on key prevents duplicates
$this->db->prepare(
"INSERT INTO idempotency_keys (key, status, created_at)
VALUES (:key, 'processing', NOW())"
)->execute(['key' => $request->idempotencyKey]);
// Step 3: Execute the actual payment
$gatewayResult = $this->gateway->charge(
$request->amount,
$request->currency,
$request->paymentMethod,
);
$result = new PaymentResult(
paymentId: $gatewayResult->id,
status: $gatewayResult->status,
amount: $request->amount,
currency: $request->currency,
);
// Step 4: Store the result for future idempotent lookups
$this->db->prepare(
"UPDATE idempotency_keys
SET status = 'completed',
response = :response,
completed_at = NOW()
WHERE key = :key"
)->execute([
'key' => $request->idempotencyKey,
'response' => json_encode($result),
]);
$this->db->commit();
return $result;
} catch (\Throwable $e) {
$this->db->rollBack();
// Mark as failed for retry
$this->db->prepare(
"UPDATE idempotency_keys
SET status = 'failed', error = :error
WHERE key = :key"
)->execute([
'key' => $request->idempotencyKey,
'error' => $e->getMessage(),
]);
throw $e;
}
}
private function findByIdempotencyKey(string $key): ?PaymentResult
{
$stmt = $this->db->prepare(
"SELECT status, response FROM idempotency_keys WHERE key = :key"
);
$stmt->execute(['key' => $key]);
$row = $stmt->fetch(\PDO::FETCH_ASSOC);
if ($row === false) {
return null;
}
if ($row['status'] === 'processing') {
// Another request is processing: wait or reject
throw new ConflictException('Request is being processed');
}
if ($row['status'] === 'completed' && $row['response'] !== null) {
return PaymentResult::fromJson($row['response']);
}
return null; // Failed: allow retry
}
}
Optimistic Locking
Вместо блокировки ресурса проверяем при записи, что данные не изменились с момента чтения.
<?php
declare(strict_types=1);
/**
* Optimistic Locking: check version on write, no lock on read
*
* Best for: low contention, read-heavy workloads
* How: version column or timestamp, check on UPDATE
*/
final class OptimisticLockingRepository
{
public function __construct(
private readonly \PDO $db,
) {}
public function findById(string $id): ?array
{
$stmt = $this->db->prepare(
'SELECT id, name, price, version FROM products WHERE id = :id'
);
$stmt->execute(['id' => $id]);
return $stmt->fetch(\PDO::FETCH_ASSOC) ?: null;
}
/**
* Update with version check.
* If version doesn't match, someone else modified the data.
*
* @throws OptimisticLockException
*/
public function update(string $id, array $data, int $expectedVersion): void
{
$stmt = $this->db->prepare(
'UPDATE products
SET name = :name, price = :price, version = version + 1, updated_at = NOW()
WHERE id = :id AND version = :expected_version'
);
$stmt->execute([
'id' => $id,
'name' => $data['name'],
'price' => $data['price'],
'expected_version' => $expectedVersion,
]);
if ($stmt->rowCount() === 0) {
throw new OptimisticLockException(
"Product $id was modified by another process. " .
"Expected version $expectedVersion."
);
}
}
/**
* Update with retry on conflict.
* Reads fresh data and reapplies the change.
*/
public function updateWithRetry(string $id, callable $modifier, int $maxRetries = 3): void
{
for ($attempt = 1; $attempt <= $maxRetries; $attempt++) {
$current = $this->findById($id);
if ($current === null) {
throw new \RuntimeException("Product $id not found");
}
$updated = $modifier($current);
try {
$this->update($id, $updated, $current['version']);
return; // Success
} catch (OptimisticLockException $e) {
if ($attempt === $maxRetries) {
throw $e; // Exhausted retries
}
// Retry with fresh data
usleep(100_000 * $attempt); // Backoff: 100ms, 200ms, 300ms
}
}
}
}
Pessimistic Locking
Блокируем ресурс перед модификацией. Другие процессы ждут освобождения.
<?php
declare(strict_types=1);
/**
* Pessimistic Locking: lock resource before reading/modifying
*
* Best for: high contention, short operations
* How: SELECT ... FOR UPDATE, advisory locks, distributed locks
*/
final class PessimisticLockingRepository
{
public function __construct(
private readonly \PDO $db,
) {}
/**
* Decrement inventory with pessimistic lock.
* Prevents overselling under concurrent access.
*/
public function decrementInventory(string $productId, int $quantity): bool
{
$this->db->beginTransaction();
try {
// Lock the row: other transactions wait here
$stmt = $this->db->prepare(
'SELECT quantity FROM inventory
WHERE product_id = :productId
FOR UPDATE'
);
$stmt->execute(['productId' => $productId]);
$row = $stmt->fetch(\PDO::FETCH_ASSOC);
if ($row === false || (int) $row['quantity'] < $quantity) {
$this->db->rollBack();
return false; // Insufficient inventory
}
// Safely decrement (we hold the lock)
$this->db->prepare(
'UPDATE inventory SET quantity = quantity - :qty, updated_at = NOW()
WHERE product_id = :productId'
)->execute([
'productId' => $productId,
'qty' => $quantity,
]);
$this->db->commit();
return true;
} catch (\Throwable $e) {
$this->db->rollBack();
throw $e;
}
}
}
/**
* Distributed lock with Redis (for cross-service locking)
*/
final class DistributedLock
{
public function __construct(
private readonly \Redis $redis,
private readonly int $defaultTtlMs = 10000,
) {}
/**
* Acquire a lock. Returns lock token or null if failed.
*/
public function acquire(string $resource, int $ttlMs = 0): ?string
{
$ttlMs = $ttlMs ?: $this->defaultTtlMs;
$token = bin2hex(random_bytes(16));
$acquired = $this->redis->set(
"lock:$resource",
$token,
['NX' => true, 'PX' => $ttlMs],
);
return $acquired ? $token : null;
}
/**
* Release a lock. Only releases if we own it (token matches).
*/
public function release(string $resource, string $token): bool
{
// Atomic check-and-delete via Lua script
$script = <<<'LUA'
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
end
return 0
LUA;
return (bool) $this->redis->eval($script, ["lock:$resource", $token], 1);
}
/**
* Execute a callback while holding a lock.
*
* @template T
* @param callable(): T $callback
* @return T
*/
public function withLock(string $resource, callable $callback): mixed
{
$token = $this->acquire($resource);
if ($token === null) {
throw new LockAcquisitionException("Could not acquire lock for: $resource");
}
try {
return $callback();
} finally {
$this->release($resource, $token);
}
}
}
Сравнение подходов
| Паттерн | Конкуренция | Производительность | Сложность | Когда использовать |
|---|---|---|---|---|
| Optimistic Lock | Низкая-средняя | Высокая (no lock wait) | Низкая | Профили, товары |
| Pessimistic Lock | Высокая | Средняя (lock wait) | Низкая | Inventory, booking |
| Distributed Lock | Кросс-сервис | Низкая (network) | Высокая | Распределенные системы |
| Idempotency Key | Любая | Средняя | Средняя | Платежи, заказы |
Выводы
Консистентность и идемпотентность -- это не бинарный выбор, а спектр. Для каждой операции определите требуемый уровень гарантий и используйте соответствующий паттерн. Начинайте с простейшего (optimistic locking) и усложняйте по мере необходимости.