Что такое идемпотентность
Идемпотентность -- свойство операции, при котором многократное выполнение дает тот же результат, что и однократное. В математике: f(f(x)) = f(x). В распределенных системах это критически важная гарантия.
Зачем нужна идемпотентность
В реальных системах сбои неизбежны:
Клиент ──► Сервер: POST /api/payments
Сервер обработал, платёж списан
Клиент ◄── Сервер: 200 OK
⚡ Сетевой разрыв! Ответ потерян
Клиент ──► Сервер: POST /api/payments (retry)
❌ Без идемпотентности: двойное списание!
✅ С идемпотентностью: вернул сохранённый ответ
Типичные причины дублирования запросов:
| Причина | Описание | Частота |
|---|---|---|
| Network timeout | Клиент не получил ответ | Высокая |
| Client retry | Браузер/мобильное приложение повторяет | Высокая |
| Load balancer retry | Nginx/HAProxy ретраит при 502/503 | Средняя |
| Queue redelivery | RabbitMQ/SQS повторно доставляет | Высокая |
| User double-click | Двойной клик на кнопке | Очень высокая |
HTTP-методы и идемпотентность
RFC 7231 определяет идемпотентность HTTP-методов:
| Метод | Идемпотентный | Безопасный | Пояснение |
|---|---|---|---|
| GET | Да | Да | Чтение, никогда не меняет данные |
| HEAD | Да | Да | Как GET, но без тела ответа |
| OPTIONS | Да | Да | Метаинформация |
| PUT | Да | Нет | Полная замена ресурса |
| DELETE | Да | Нет | Удаление (повторное -- 404, но состояние то же) |
| POST | Нет | Нет | Создание ресурса, побочные эффекты |
| PATCH | Нет | Нет | Частичное обновление (зависит от реализации) |
PUT идемпотентен по определению: PUT /users/123 {name: "Alice"} -- сколько бы раз ни вызвали, результат один. POST -- нет: каждый POST /orders создает новый заказ.
Почему PATCH не идемпотентен
PATCH /account/balance {"operation": "add", "amount": 100}
# Первый вызов: баланс 0 → 100
# Второй вызов: баланс 100 → 200 ← результат отличается!
# Но если PATCH устанавливает значение напрямую:
PATCH /users/123 {"name": "Alice"}
# Первый вызов: name = Alice
# Второй вызов: name = Alice ← идемпотентно!
PATCH может быть идемпотентным, если описывает конечное состояние, а не дельту.
Idempotency Key Pattern
Основная идея: клиент генерирует уникальный ключ для каждой логической операции и передает его в заголовке. Сервер сохраняет результат и при повторном запросе возвращает сохраненный ответ.
Как это работает
Idempotency Key Flow
Клиент Сервер
│ │
│ POST /api/payments │
│ Idempotency-Key: idk_abc123 │
│ ───────────────────────────────────────►│
│ │ 1. Проверить ключ в БД
│ │ 2. Ключ не найден → обработать
│ │ 3. Сохранить: key → response
│ HTTP 200 {payment_id: "pay_xyz"} │
│ ◄───────────────────────────────────────│
│ │
│ ⚡ Timeout! Retry... │
│ │
│ POST /api/payments │
│ Idempotency-Key: idk_abc123 │
│ ───────────────────────────────────────►│
│ │ 1. Проверить ключ в БД
│ │ 2. Ключ НАЙДЕН → вернуть
│ HTTP 200 {payment_id: "pay_xyz"} │ сохранённый ответ
│ ◄───────────────────────────────────────│
Стандарт и реализации в индустрии
| Компания | Заголовок | Формат ключа | TTL |
|---|---|---|---|
| Stripe | Idempotency-Key |
UUID v4 | 24 часа |
| PayPal | PayPal-Request-Id |
UUID v4 | Зависит от API |
| Amazon Pay | x-amz-pay-idempotency-key |
UUID v4 | 24 часа |
| IETF Draft | Idempotency-Key |
Произвольная строка | Рекомендован |
PostgreSQL schema для deduplication
-- Deduplication table for idempotency keys
CREATE TABLE idempotency_keys (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
key VARCHAR(255) NOT NULL,
user_id UUID NOT NULL,
-- Request fingerprint
method VARCHAR(10) NOT NULL,
path VARCHAR(512) NOT NULL,
request_hash VARCHAR(64), -- SHA-256 of request body
-- Stored response
status_code SMALLINT NOT NULL,
response_headers JSONB DEFAULT '{}',
response_body JSONB,
-- Lifecycle
locked_at TIMESTAMPTZ, -- Processing lock
completed_at TIMESTAMPTZ, -- When processing finished
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
expires_at TIMESTAMPTZ NOT NULL DEFAULT NOW() + INTERVAL '24 hours',
-- Unique per user to prevent cross-user collision
CONSTRAINT uq_idempotency_key_user UNIQUE (key, user_id)
);
-- Index for cleanup job
CREATE INDEX idx_idempotency_keys_expires ON idempotency_keys (expires_at)
WHERE completed_at IS NOT NULL;
-- Index for lock detection (find stale locks)
CREATE INDEX idx_idempotency_keys_locked ON idempotency_keys (locked_at)
WHERE locked_at IS NOT NULL AND completed_at IS NULL;
-- TTL cleanup: run periodically via cron or pg_cron
-- DELETE FROM idempotency_keys WHERE expires_at < NOW();
PHP: Symfony Event Subscriber для Idempotency-Key
<?php
declare(strict_types=1);
namespace App\EventSubscriber;
use Symfony\Component\EventDispatcher\EventSubscriberInterface;
use Symfony\Component\HttpFoundation\JsonResponse;
use Symfony\Component\HttpFoundation\Request;
use Symfony\Component\HttpFoundation\Response;
use Symfony\Component\HttpKernel\Event\RequestEvent;
use Symfony\Component\HttpKernel\Event\ResponseEvent;
use Symfony\Component\HttpKernel\KernelEvents;
/**
* Handles Idempotency-Key header for POST/PATCH requests.
*
* Flow:
* 1. onKernelRequest: check if key exists in DB, return cached response
* 2. Controller processes request normally
* 3. onKernelResponse: store response in DB for future deduplication
*/
final class IdempotencySubscriber implements EventSubscriberInterface
{
private const HEADER = 'Idempotency-Key';
private const MAX_KEY_LENGTH = 255;
private const LOCK_TIMEOUT_SECONDS = 30;
public function __construct(
private readonly \PDO $db,
private readonly \Psr\Log\LoggerInterface $logger,
) {}
public static function getSubscribedEvents(): array
{
return [
KernelEvents::REQUEST => ['onKernelRequest', 100],
KernelEvents::RESPONSE => ['onKernelResponse', -100],
];
}
public function onKernelRequest(RequestEvent $event): void
{
$request = $event->getRequest();
if (!$event->isMainRequest()) {
return;
}
// Only apply to non-idempotent methods
if (!in_array($request->getMethod(), ['POST', 'PATCH'], true)) {
return;
}
$idempotencyKey = $request->headers->get(self::HEADER);
if ($idempotencyKey === null) {
return; // No key provided: process normally
}
// Validate key format
if (strlen($idempotencyKey) > self::MAX_KEY_LENGTH) {
$event->setResponse(new JsonResponse(
['error' => 'Idempotency key too long'],
Response::HTTP_BAD_REQUEST,
));
return;
}
$userId = $this->resolveUserId($request);
$existing = $this->findExistingKey($idempotencyKey, $userId);
if ($existing === null) {
// New key: lock it for processing
$this->lockKey($idempotencyKey, $userId, $request);
$request->attributes->set('_idempotency_key', $idempotencyKey);
return;
}
// Key exists but still processing (concurrent request)
if ($existing['completed_at'] === null) {
if ($this->isLockStale($existing['locked_at'])) {
// Stale lock: allow retry by re-locking
$this->relockKey($idempotencyKey, $userId);
$request->attributes->set('_idempotency_key', $idempotencyKey);
return;
}
$event->setResponse(new JsonResponse(
['error' => 'Request is already being processed'],
Response::HTTP_CONFLICT,
));
return;
}
// Verify request body matches (prevent key reuse for different payload)
$currentHash = $this->hashRequestBody($request);
if ($existing['request_hash'] !== null && $existing['request_hash'] !== $currentHash) {
$event->setResponse(new JsonResponse(
['error' => 'Idempotency key reused with different request body'],
Response::HTTP_UNPROCESSABLE_ENTITY,
));
return;
}
// Return cached response
$this->logger->info('Returning cached idempotent response', [
'key' => $idempotencyKey,
]);
$event->setResponse(new JsonResponse(
json_decode($existing['response_body'], true),
(int) $existing['status_code'],
json_decode($existing['response_headers'], true),
));
}
public function onKernelResponse(ResponseEvent $event): void
{
$request = $event->getRequest();
if (!$event->isMainRequest()) {
return;
}
$idempotencyKey = $request->attributes->get('_idempotency_key');
if ($idempotencyKey === null) {
return;
}
$response = $event->getResponse();
$userId = $this->resolveUserId($request);
// Store the response for future deduplication
$this->completeKey(
$idempotencyKey,
$userId,
$response->getStatusCode(),
$this->extractHeaders($response),
$response->getContent(),
$this->hashRequestBody($request),
);
}
private function findExistingKey(string $key, string $userId): ?array
{
$stmt = $this->db->prepare(
'SELECT status_code, response_headers, response_body,
request_hash, locked_at, completed_at
FROM idempotency_keys
WHERE key = :key AND user_id = :userId AND expires_at > NOW()'
);
$stmt->execute(['key' => $key, 'userId' => $userId]);
return $stmt->fetch(\PDO::FETCH_ASSOC) ?: null;
}
private function lockKey(string $key, string $userId, Request $request): void
{
$this->db->prepare(
'INSERT INTO idempotency_keys (key, user_id, method, path, locked_at)
VALUES (:key, :userId, :method, :path, NOW())
ON CONFLICT (key, user_id) DO NOTHING'
)->execute([
'key' => $key,
'userId' => $userId,
'method' => $request->getMethod(),
'path' => $request->getPathInfo(),
]);
}
private function completeKey(
string $key,
string $userId,
int $statusCode,
string $headers,
string $body,
string $requestHash,
): void {
$this->db->prepare(
'UPDATE idempotency_keys
SET status_code = :statusCode,
response_headers = :headers,
response_body = :body,
request_hash = :requestHash,
completed_at = NOW()
WHERE key = :key AND user_id = :userId'
)->execute([
'statusCode' => $statusCode,
'headers' => $headers,
'body' => $body,
'requestHash' => $requestHash,
'key' => $key,
'userId' => $userId,
]);
}
private function relockKey(string $key, string $userId): void
{
$this->db->prepare(
'UPDATE idempotency_keys
SET locked_at = NOW(), completed_at = NULL
WHERE key = :key AND user_id = :userId'
)->execute(['key' => $key, 'userId' => $userId]);
}
private function isLockStale(?string $lockedAt): bool
{
if ($lockedAt === null) {
return true;
}
$locked = new \DateTimeImmutable($lockedAt);
$threshold = new \DateTimeImmutable("-" . self::LOCK_TIMEOUT_SECONDS . " seconds");
return $locked < $threshold;
}
private function resolveUserId(Request $request): string
{
// Extract user ID from JWT token or session
return $request->attributes->get('_user_id', 'anonymous');
}
private function hashRequestBody(Request $request): string
{
return hash('sha256', $request->getContent());
}
private function extractHeaders(Response $response): string
{
$headers = [];
foreach (['Content-Type', 'X-Request-Id'] as $name) {
$value = $response->headers->get($name);
if ($value !== null) {
$headers[$name] = $value;
}
}
return json_encode($headers);
}
}
Go: HTTP Middleware для Idempotency-Key
package middleware
import (
"context"
"crypto/sha256"
"database/sql"
"encoding/hex"
"encoding/json"
"io"
"log/slog"
"net/http"
"strings"
"time"
)
// IdempotencyRecord stores the cached response for a given key.
type IdempotencyRecord struct {
StatusCode int `json:"status_code"`
ResponseHeaders map[string]string `json:"response_headers"`
ResponseBody json.RawMessage `json:"response_body"`
RequestHash string `json:"request_hash"`
CompletedAt *time.Time `json:"completed_at"`
LockedAt *time.Time `json:"locked_at"`
}
// IdempotencyStore abstracts storage for idempotency keys.
type IdempotencyStore interface {
Find(ctx context.Context, key, userID string) (*IdempotencyRecord, error)
Lock(ctx context.Context, key, userID, method, path string) error
Complete(ctx context.Context, key, userID string, rec IdempotencyRecord) error
}
const (
idempotencyHeader = "Idempotency-Key"
maxKeyLength = 255
lockTimeoutSeconds = 30
)
// Idempotency returns middleware that handles Idempotency-Key header.
func Idempotency(store IdempotencyStore, logger *slog.Logger) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
// Only apply to non-idempotent methods
if r.Method != http.MethodPost && r.Method != http.MethodPatch {
next.ServeHTTP(w, r)
return
}
key := r.Header.Get(idempotencyHeader)
if key == "" {
next.ServeHTTP(w, r)
return
}
if len(key) > maxKeyLength {
writeJSON(w, http.StatusBadRequest, map[string]string{
"error": "Idempotency key too long",
})
return
}
userID := resolveUserID(r)
// Read body for hashing (buffer it for re-read)
body, err := io.ReadAll(r.Body)
if err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{
"error": "Failed to read request body",
})
return
}
r.Body = io.NopCloser(strings.NewReader(string(body)))
bodyHash := hashBody(body)
existing, err := store.Find(r.Context(), key, userID)
if err != nil {
logger.Error("idempotency store lookup failed", "error", err)
next.ServeHTTP(w, r)
return
}
if existing != nil {
// Still processing by another request
if existing.CompletedAt == nil {
if !isLockStale(existing.LockedAt) {
writeJSON(w, http.StatusConflict, map[string]string{
"error": "Request is already being processed",
})
return
}
} else {
// Verify request body matches
if existing.RequestHash != "" && existing.RequestHash != bodyHash {
writeJSON(w, http.StatusUnprocessableEntity, map[string]string{
"error": "Idempotency key reused with different request body",
})
return
}
// Return cached response
logger.Info("returning cached idempotent response", "key", key)
for k, v := range existing.ResponseHeaders {
w.Header().Set(k, v)
}
w.WriteHeader(existing.StatusCode)
w.Write(existing.ResponseBody) //nolint:errcheck
return
}
}
// Lock the key for processing
if err := store.Lock(r.Context(), key, userID, r.Method, r.URL.Path); err != nil {
logger.Error("failed to lock idempotency key", "error", err)
next.ServeHTTP(w, r)
return
}
// Wrap ResponseWriter to capture the response
recorder := &responseRecorder{
ResponseWriter: w,
statusCode: http.StatusOK,
}
next.ServeHTTP(recorder, r)
// Store the response
rec := IdempotencyRecord{
StatusCode: recorder.statusCode,
ResponseHeaders: extractHeaders(recorder.Header()),
ResponseBody: recorder.body,
RequestHash: bodyHash,
}
if err := store.Complete(r.Context(), key, userID, rec); err != nil {
logger.Error("failed to store idempotent response", "error", err)
}
})
}
}
// responseRecorder captures the response for storage.
type responseRecorder struct {
http.ResponseWriter
statusCode int
body []byte
}
func (r *responseRecorder) WriteHeader(code int) {
r.statusCode = code
r.ResponseWriter.WriteHeader(code)
}
func (r *responseRecorder) Write(b []byte) (int, error) {
r.body = append(r.body, b...)
return r.ResponseWriter.Write(b)
}
func hashBody(body []byte) string {
h := sha256.Sum256(body)
return hex.EncodeToString(h[:])
}
func isLockStale(lockedAt *time.Time) bool {
if lockedAt == nil {
return true
}
return time.Since(*lockedAt) > lockTimeoutSeconds*time.Second
}
func resolveUserID(r *http.Request) string {
if uid := r.Context().Value("user_id"); uid != nil {
return uid.(string)
}
return "anonymous"
}
func extractHeaders(h http.Header) map[string]string {
result := make(map[string]string)
for _, name := range []string{"Content-Type", "X-Request-Id"} {
if v := h.Get(name); v != "" {
result[name] = v
}
}
return result
}
func writeJSON(w http.ResponseWriter, status int, data any) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
json.NewEncoder(w).Encode(data) //nolint:errcheck
}
Подключение Go middleware
package main
import (
"net/http"
"github.com/go-chi/chi/v5"
)
func main() {
r := chi.NewRouter()
store := NewPostgresIdempotencyStore(db)
logger := slog.Default()
// Apply idempotency middleware to API routes
r.Route("/api/v1", func(r chi.Router) {
r.Use(middleware.Idempotency(store, logger))
r.Post("/payments", handleCreatePayment)
r.Post("/orders", handleCreateOrder)
})
http.ListenAndServe(":8080", r)
}
Redis-based Deduplication (SETNX + TTL)
Для высоконагруженных систем, где обращение к PostgreSQL при каждом запросе слишком дорого, используют Redis.
<?php
declare(strict_types=1);
/**
* Redis-based idempotency store using SETNX + TTL.
*
* Faster than PostgreSQL but less durable:
* - Redis may lose data on restart (unless AOF is configured)
* - Suitable for non-critical deduplication
* - Use PostgreSQL for financial transactions
*/
final class RedisIdempotencyStore
{
private const PREFIX = 'idem:';
private const LOCK_SUFFIX = ':lock';
private const DEFAULT_TTL = 86400; // 24 hours
public function __construct(
private readonly \Redis $redis,
private readonly int $ttl = self::DEFAULT_TTL,
private readonly int $lockTtl = 30,
) {}
/**
* Try to acquire processing lock for the key.
* Returns true if this is a new request.
*/
public function tryLock(string $key, string $userId): bool
{
$lockKey = self::PREFIX . $userId . ':' . $key . self::LOCK_SUFFIX;
// SETNX: only sets if key doesn't exist (atomic)
$acquired = $this->redis->set($lockKey, time(), [
'NX' => true, // Only set if Not eXists
'EX' => $this->lockTtl, // Auto-expire lock
]);
return (bool) $acquired;
}
/**
* Store the response for future deduplication.
*/
public function store(string $key, string $userId, array $response): void
{
$dataKey = self::PREFIX . $userId . ':' . $key;
$lockKey = $dataKey . self::LOCK_SUFFIX;
// Atomic: store response and remove lock in pipeline
$this->redis->multi();
$this->redis->setex($dataKey, $this->ttl, json_encode($response));
$this->redis->del($lockKey);
$this->redis->exec();
}
/**
* Get cached response if exists.
*/
public function find(string $key, string $userId): ?array
{
$dataKey = self::PREFIX . $userId . ':' . $key;
$data = $this->redis->get($dataKey);
if ($data === false) {
return null;
}
return json_decode($data, true);
}
/**
* Check if key is currently being processed.
*/
public function isLocked(string $key, string $userId): bool
{
$lockKey = self::PREFIX . $userId . ':' . $key . self::LOCK_SUFFIX;
return (bool) $this->redis->exists($lockKey);
}
}
Сравнение: PostgreSQL vs Redis для дедупликации
| Критерий | PostgreSQL | Redis |
|---|---|---|
| Надежность | Высокая (ACID) | Средняя (может потерять при рестарте) |
| Скорость | ~2-5ms | ~0.1-0.5ms |
| Подходит для | Платежи, финансы | Некритичные операции |
| Масштабирование | Вертикальное | Горизонтальное (Cluster) |
| TTL cleanup | Нужен cron | Встроенный (EXPIRE) |
| Concurrent locks | FOR UPDATE / Advisory | SETNX (атомарный) |
Go: Redis Idempotency Store
package idempotency
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
// RedisStore implements IdempotencyStore using Redis SETNX.
type RedisStore struct {
client *redis.Client
ttl time.Duration
lockTTL time.Duration
}
func NewRedisStore(client *redis.Client) *RedisStore {
return &RedisStore{
client: client,
ttl: 24 * time.Hour,
lockTTL: 30 * time.Second,
}
}
func (s *RedisStore) Find(ctx context.Context, key, userID string) (*IdempotencyRecord, error) {
dataKey := s.dataKey(key, userID)
data, err := s.client.Get(ctx, dataKey).Bytes()
if err == redis.Nil {
// Check if locked (being processed)
lockKey := s.lockKey(key, userID)
exists, _ := s.client.Exists(ctx, lockKey).Result()
if exists > 0 {
now := time.Now()
return &IdempotencyRecord{LockedAt: &now}, nil
}
return nil, nil
}
if err != nil {
return nil, fmt.Errorf("redis get: %w", err)
}
var rec IdempotencyRecord
if err := json.Unmarshal(data, &rec); err != nil {
return nil, fmt.Errorf("unmarshal record: %w", err)
}
now := time.Now()
rec.CompletedAt = &now
return &rec, nil
}
func (s *RedisStore) Lock(ctx context.Context, key, userID, method, path string) error {
lockKey := s.lockKey(key, userID)
ok, err := s.client.SetNX(ctx, lockKey, time.Now().Unix(), s.lockTTL).Result()
if err != nil {
return fmt.Errorf("redis setnx: %w", err)
}
if !ok {
return fmt.Errorf("key already locked: %s", key)
}
return nil
}
func (s *RedisStore) Complete(ctx context.Context, key, userID string, rec IdempotencyRecord) error {
dataKey := s.dataKey(key, userID)
lockKey := s.lockKey(key, userID)
data, err := json.Marshal(rec)
if err != nil {
return fmt.Errorf("marshal record: %w", err)
}
// Pipeline: store data + remove lock atomically
pipe := s.client.Pipeline()
pipe.Set(ctx, dataKey, data, s.ttl)
pipe.Del(ctx, lockKey)
_, err = pipe.Exec(ctx)
return err
}
func (s *RedisStore) dataKey(key, userID string) string {
return fmt.Sprintf("idem:%s:%s", userID, key)
}
func (s *RedisStore) lockKey(key, userID string) string {
return fmt.Sprintf("idem:%s:%s:lock", userID, key)
}
Retry Safety: клиентская сторона
Идемпотентность на сервере бесполезна без правильных ретраев на клиенте.
Exponential Backoff с Retry-After
package httpclient
import (
"context"
"fmt"
"math"
"math/rand/v2"
"net/http"
"strconv"
"time"
)
// RetryConfig configures retry behavior.
type RetryConfig struct {
MaxRetries int
BaseDelay time.Duration
MaxDelay time.Duration
RetryableCodes map[int]bool
}
func DefaultRetryConfig() RetryConfig {
return RetryConfig{
MaxRetries: 3,
BaseDelay: 200 * time.Millisecond,
MaxDelay: 30 * time.Second,
RetryableCodes: map[int]bool{
http.StatusRequestTimeout: true, // 408
http.StatusTooManyRequests: true, // 429
http.StatusInternalServerError: true, // 500
http.StatusBadGateway: true, // 502
http.StatusServiceUnavailable: true, // 503
http.StatusGatewayTimeout: true, // 504
},
}
}
// DoWithRetry executes an HTTP request with retries and exponential backoff.
// Respects Retry-After header from the server.
func DoWithRetry(
ctx context.Context,
client *http.Client,
req *http.Request,
cfg RetryConfig,
) (*http.Response, error) {
var lastErr error
for attempt := 0; attempt <= cfg.MaxRetries; attempt++ {
if attempt > 0 {
delay := calculateBackoff(attempt-1, cfg.BaseDelay, cfg.MaxDelay)
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(delay):
}
}
resp, err := client.Do(req)
if err != nil {
lastErr = err
continue // Network error: retry
}
// Success or non-retryable status
if !cfg.RetryableCodes[resp.StatusCode] {
return resp, nil
}
// Check Retry-After header
if retryAfter := resp.Header.Get("Retry-After"); retryAfter != "" {
if seconds, err := strconv.Atoi(retryAfter); err == nil {
delay := time.Duration(seconds) * time.Second
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(delay):
}
}
}
lastErr = fmt.Errorf("HTTP %d from %s", resp.StatusCode, req.URL)
resp.Body.Close()
}
return nil, fmt.Errorf("max retries exceeded: %w", lastErr)
}
// calculateBackoff returns delay with exponential backoff + full jitter.
func calculateBackoff(attempt int, base, max time.Duration) time.Duration {
exp := math.Pow(2, float64(attempt))
delay := time.Duration(float64(base) * exp)
if delay > max {
delay = max
}
// Full jitter: random between 0 and calculated delay
jitter := time.Duration(rand.Int64N(int64(delay)))
return jitter
}
Идемпотентность в очередях сообщений
Очереди гарантируют at-least-once delivery -- сообщение может прийти дважды. Consumer обязан быть идемпотентным.
RabbitMQ: Message Deduplication
<?php
declare(strict_types=1);
/**
* Idempotent message consumer for RabbitMQ.
*
* RabbitMQ guarantees at-least-once delivery:
* - Consumer crashes after processing but before ack → redelivery
* - Network partition → redelivery
* - Consumer must handle duplicates safely
*/
final class IdempotentConsumer
{
public function __construct(
private readonly \PDO $db,
private readonly \Psr\Log\LoggerInterface $logger,
) {}
/**
* Process a message idempotently using message_id as dedup key.
*
* @param callable(array): void $handler
*/
public function handle(string $messageId, array $payload, callable $handler): void
{
// Check if already processed
if ($this->isProcessed($messageId)) {
$this->logger->info('Skipping duplicate message', [
'message_id' => $messageId,
]);
return; // Already processed: skip
}
// Process within transaction
$this->db->beginTransaction();
try {
// Mark as processing (INSERT with unique constraint)
$this->markAsProcessing($messageId);
// Execute business logic
$handler($payload);
// Mark as completed
$this->markAsCompleted($messageId);
$this->db->commit();
} catch (\Throwable $e) {
$this->db->rollBack();
// Remove the processing mark so message can be retried
$this->removeProcessingMark($messageId);
throw $e;
}
}
private function isProcessed(string $messageId): bool
{
$stmt = $this->db->prepare(
'SELECT 1 FROM processed_messages
WHERE message_id = :id AND status = :status'
);
$stmt->execute(['id' => $messageId, 'status' => 'completed']);
return $stmt->fetchColumn() !== false;
}
private function markAsProcessing(string $messageId): void
{
$this->db->prepare(
'INSERT INTO processed_messages (message_id, status, created_at)
VALUES (:id, :status, NOW())
ON CONFLICT (message_id) DO NOTHING'
)->execute(['id' => $messageId, 'status' => 'processing']);
}
private function markAsCompleted(string $messageId): void
{
$this->db->prepare(
'UPDATE processed_messages
SET status = :status, completed_at = NOW()
WHERE message_id = :id'
)->execute(['id' => $messageId, 'status' => 'completed']);
}
private function removeProcessingMark(string $messageId): void
{
$this->db->prepare(
'DELETE FROM processed_messages
WHERE message_id = :id AND status = :status'
)->execute(['id' => $messageId, 'status' => 'processing']);
}
}
// Schema for processed messages:
// CREATE TABLE processed_messages (
// message_id VARCHAR(255) PRIMARY KEY,
// status VARCHAR(20) NOT NULL DEFAULT 'processing',
// created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
// completed_at TIMESTAMPTZ
// );
SQS: встроенная дедупликация
Amazon SQS FIFO очереди имеют встроенную дедупликацию с 5-минутным окном:
Producer → SQS FIFO (MessageDeduplicationId: "order-123-payment")
↓ Окно 5 минут
Повторное сообщение с тем же DeduplicationId → автоматически отклонено
Идемпотентность в Saga Pattern
В Saga каждый шаг должен быть идемпотентным, потому что compensating transaction может выполняться повторно.
<?php
declare(strict_types=1);
/**
* Saga step with idempotent execution and compensation.
*
* Each step tracks its execution state.
* Re-executing a completed step returns cached result.
* Re-compensating a compensated step is a no-op.
*/
final class IdempotentSagaStep
{
public function __construct(
private readonly string $sagaId,
private readonly string $stepName,
private readonly \PDO $db,
) {}
/**
* Execute step idempotently.
* Returns cached result if step was already completed.
*/
public function execute(callable $action): mixed
{
$state = $this->getState();
// Already completed: return cached result
if ($state === 'completed') {
return $this->getCachedResult();
}
// Already compensated: cannot re-execute
if ($state === 'compensated') {
throw new \LogicException("Cannot execute compensated step: {$this->stepName}");
}
$this->db->beginTransaction();
try {
$this->setState('executing');
$result = $action();
$this->setState('completed', $result);
$this->db->commit();
return $result;
} catch (\Throwable $e) {
$this->db->rollBack();
$this->setState('failed');
throw $e;
}
}
/**
* Compensate step idempotently.
* No-op if already compensated.
*/
public function compensate(callable $compensation): void
{
$state = $this->getState();
// Not executed or already compensated: nothing to do
if ($state === null || $state === 'compensated') {
return;
}
$this->db->beginTransaction();
try {
$compensation();
$this->setState('compensated');
$this->db->commit();
} catch (\Throwable $e) {
$this->db->rollBack();
throw $e;
}
}
private function getState(): ?string
{
$stmt = $this->db->prepare(
'SELECT state FROM saga_steps
WHERE saga_id = :sagaId AND step_name = :stepName'
);
$stmt->execute([
'sagaId' => $this->sagaId,
'stepName' => $this->stepName,
]);
return $stmt->fetchColumn() ?: null;
}
private function setState(string $state, mixed $result = null): void
{
$this->db->prepare(
'INSERT INTO saga_steps (saga_id, step_name, state, result, updated_at)
VALUES (:sagaId, :stepName, :state, :result, NOW())
ON CONFLICT (saga_id, step_name)
DO UPDATE SET state = :state2, result = :result2, updated_at = NOW()'
)->execute([
'sagaId' => $this->sagaId,
'stepName' => $this->stepName,
'state' => $state,
'result' => $result !== null ? json_encode($result) : null,
'state2' => $state,
'result2' => $result !== null ? json_encode($result) : null,
]);
}
private function getCachedResult(): mixed
{
$stmt = $this->db->prepare(
'SELECT result FROM saga_steps
WHERE saga_id = :sagaId AND step_name = :stepName'
);
$stmt->execute([
'sagaId' => $this->sagaId,
'stepName' => $this->stepName,
]);
$result = $stmt->fetchColumn();
return $result !== false ? json_decode($result, true) : null;
}
}
Anti-patterns: когда идемпотентность НЕ нужна
Не все операции требуют идемпотентности. Добавление сложности без необходимости -- это тоже ошибка.
| Сценарий | Нужна идемпотентность | Почему |
|---|---|---|
| Платежи, переводы | Да | Дублирование = финансовые потери |
| Создание заказов | Да | Двойной заказ = проблемы |
| Отправка email | Да | Двойное письмо = спам |
| Чтение данных (GET) | Нет | Уже идемпотентно по определению |
| Логирование/аналитика | Нет | Дублирование безвредно (обычно) |
| Обновление счетчика лайков | Зависит | Повторный лайк уже ограничен бизнес-логикой |
| Генерация отчетов | Нет | Повторная генерация = тот же результат |
| DELETE с фиксированным ID | Нет | HTTP DELETE уже идемпотентен |
Ошибки при реализации
Anti-pattern 1: Idempotency key без TTL
─────────────────────────────────────────
❌ Ключи накапливаются бесконечно → БД растёт
✅ Всегда устанавливайте TTL (обычно 24 часа)
Anti-pattern 2: Один ключ на все запросы
─────────────────────────────────────────
❌ Клиент переиспользует ключ для разных операций
✅ Генерировать уникальный ключ для каждой операции
Anti-pattern 3: Idempotency без проверки тела запроса
─────────────────────────────────────────
❌ Тот же ключ, другое тело → неправильный ответ
✅ Хеш тела запроса как дополнительная проверка
Anti-pattern 4: Блокировка без timeout
─────────────────────────────────────────
❌ Процесс упал → ключ залочен навсегда
✅ TTL на lock + проверка stale locks
Выводы
Идемпотентность -- фундаментальное свойство надежных распределенных систем. Для POST/PATCH запросов используйте Idempotency-Key паттерн с deduplication table. Для очередей -- отслеживайте message_id обработанных сообщений. Для Saga -- каждый шаг должен быть безопасен при повторном выполнении. Начинайте с PostgreSQL для критичных данных и переходите на Redis для высоконагруженных сценариев.