Наблюдаемость: трассировка, корреляция, метрики и алерты

Наблюдаемость: трассировка, корреляция, метрики и алерты

В синхронной системе отладка проста: запрос приходит, проходит через middleware, handler, service, repository - и ответ уходит. Один HTTP-запрос, одна горутина, один стектрейс. В event-driven системе одно действие пользователя может породить цепочку из десятка событий, обработанных разными подписчиками, возможно в разных сервисах, с задержкой в миллисекунды или минуты. Без специальных инструментов наблюдаемости вы будете искать причину бага, перебирая логи вручную.

Три столпа наблюдаемости (Three Pillars of Observability) - логи, метрики, трейсы - особенно критичны для event-driven архитектуры.

Correlation ID: связываем цепочку событий

Correlation ID - уникальный идентификатор, который генерируется при входящем запросе и передаётся через все события и подписчики. Он позволяет собрать все логи одного пользовательского действия в один trace.

Correlation ID путешествует через HTTP → Use-Case → Event → Subscriber; один grep собирает всю цепочку

Middleware для генерации correlation_id

// CorrelationMiddleware добавляет correlation_id в контекст запроса.
// Если клиент передал X-Correlation-ID - используем его,
// иначе генерируем новый.
func CorrelationMiddleware(next http.Handler) http.Handler {
    return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
        correlationID := r.Header.Get("X-Correlation-ID")
        if correlationID == "" {
            correlationID = ulid.Make().String()
        }

        // Сохраняем в контекст
        ctx := context.WithValue(r.Context(), correlationKey, correlationID)

        // Отдаём в заголовке ответа (удобно для отладки)
        w.Header().Set("X-Correlation-ID", correlationID)

        next.ServeHTTP(w, r.WithContext(ctx))
    })
}

// CorrelationFromCtx извлекает correlation_id из контекста.
func CorrelationFromCtx(ctx context.Context) string {
    if v, ok := ctx.Value(correlationKey).(string); ok {
        return v
    }
    return ""
}
<?php
// src/Infrastructure/Http/EventListener/CorrelationIdListener.php
declare(strict_types=1);

namespace App\Infrastructure\Http\EventListener;

use Symfony\Component\EventDispatcher\Attribute\AsEventListener;
use Symfony\Component\HttpFoundation\RequestStack;
use Symfony\Component\HttpKernel\Event\RequestEvent;
use Symfony\Component\HttpKernel\Event\ResponseEvent;
use Symfony\Component\HttpKernel\KernelEvents;
use Symfony\Component\Uid\Ulid;

// Слушаем kernel.request и kernel.response: на входе кладём correlation_id
// в атрибуты запроса (если клиент не передал - генерируем), на выходе
// добавляем тот же id в заголовок ответа.
#[AsEventListener(event: KernelEvents::REQUEST, method: 'onRequest', priority: 256)]
#[AsEventListener(event: KernelEvents::RESPONSE, method: 'onResponse')]
final class CorrelationIdListener
{
    public const string ATTRIBUTE = 'correlation_id';
    public const string HEADER = 'X-Correlation-ID';

    public function onRequest(RequestEvent $event): void
    {
        $request = $event->getRequest();
        $correlationId = $request->headers->get(self::HEADER) ?? (string) new Ulid();
        $request->attributes->set(self::ATTRIBUTE, $correlationId);
    }

    public function onResponse(ResponseEvent $event): void
    {
        $correlationId = $event->getRequest()->attributes->get(self::ATTRIBUTE);
        if (is_string($correlationId)) {
            $event->getResponse()->headers->set(self::HEADER, $correlationId);
        }
    }
}

// src/Infrastructure/Logging/CorrelationIdProcessor.php
// Monolog processor: добавляет correlation_id в каждую запись лога.
final class CorrelationIdProcessor
{
    public function __construct(
        private readonly RequestStack $requestStack,
    ) {}

    public function __invoke(array $record): array
    {
        $request = $this->requestStack->getCurrentRequest();
        if ($request !== null) {
            $record['extra']['correlation_id'] = $request->attributes->get(CorrelationIdListener::ATTRIBUTE);
        }

        return $record;
    }
}

Передача correlation_id в событие

// PublishEvent включает correlation_id в payload события.
func (p *EventPublisher) PublishEvent(ctx context.Context, eventType string, data any) error {
    envelope := EventEnvelope{
        EventID:       ulid.Make().String(),
        Type:          eventType,
        CorrelationID: CorrelationFromCtx(ctx),
        OccurredAt:    time.Now().UTC(),
        Data:          data,
    }

    payload, err := json.Marshal(envelope)
    if err != nil {
        return fmt.Errorf("marshal event: %w", err)
    }

    return p.broker.Publish(ctx, eventType, payload)
}

// EventEnvelope - конверт события с метаданными.
type EventEnvelope struct {
    EventID       string    `json:"event_id"`
    Type          string    `json:"type"`
    CorrelationID string    `json:"correlation_id"`
    OccurredAt    time.Time `json:"occurred_at"`
    Data          any       `json:"data"`
}
<?php
// src/Infrastructure/Messaging/EventEnvelope.php
declare(strict_types=1);

namespace App\Infrastructure\Messaging;

// EventEnvelope - конверт события с метаданными.
final readonly class EventEnvelope
{
    public function __construct(
        public string $eventId,
        public string $type,
        public string $correlationId,
        public \DateTimeImmutable $occurredAt,
        public array $data,
    ) {}
}

// src/Infrastructure/Messaging/EventPublisher.php
namespace App\Infrastructure\Messaging;

use App\Infrastructure\Http\EventListener\CorrelationIdListener;
use Symfony\Component\HttpFoundation\RequestStack;
use Symfony\Component\Uid\Ulid;

final class EventPublisher
{
    public function __construct(
        private readonly BrokerInterface $broker,
        private readonly RequestStack $requestStack,
    ) {}

    // publishEvent включает correlation_id в payload события.
    public function publishEvent(string $eventType, array $data): void
    {
        $envelope = new EventEnvelope(
            eventId: (string) new Ulid(),
            type: $eventType,
            correlationId: $this->correlationId(),
            occurredAt: new \DateTimeImmutable(),
            data: $data,
        );

        $payload = json_encode($envelope, JSON_THROW_ON_ERROR);
        $this->broker->publish($eventType, $payload);
    }

    private function correlationId(): string
    {
        $request = $this->requestStack->getCurrentRequest();

        return $request?->attributes->get(CorrelationIdListener::ATTRIBUTE) ?? '';
    }
}

В Symfony Messenger вместо собственного EventEnvelope обычно используется Symfony\Component\Messenger\Envelope + stamp (например, CorrelationIdStamp), а middleware на bus добавляет stamp в каждое исходящее сообщение и копирует его на дочерние сообщения внутри handler-а.

Подписчик, получив событие, извлекает correlation_id и кладёт его в свой контекст. Все логи этого подписчика будут содержать тот же ID. В Kibana/Grafana Loki достаточно отфильтровать по correlation_id, чтобы увидеть всю цепочку.

Structured Logging с slog

Каждый подписчик логирует обработку событий структурированно. Не printf/echo, а структурный logger с атрибутами:

// HandleEvent обрабатывает событие с полным структурированным логированием.
func (h *AchievementHandler) HandleEvent(ctx context.Context, envelope EventEnvelope) error {
    logger := slog.With(
        slog.String("correlation_id", envelope.CorrelationID),
        slog.String("event_id", envelope.EventID),
        slog.String("event_type", envelope.Type),
        slog.String("handler", "achievement_checker"),
    )

    start := time.Now()
    logger.InfoContext(ctx, "event processing started")

    err := h.process(ctx, envelope)

    duration := time.Since(start)
    if err != nil {
        logger.ErrorContext(ctx, "event processing failed",
            slog.String("err", err.Error()),
            slog.Duration("duration", duration),
        )
        return err
    }

    logger.InfoContext(ctx, "event processing completed",
        slog.Duration("duration", duration),
    )
    return nil
}
<?php
// src/Application/EventHandler/AchievementHandler.php
declare(strict_types=1);

namespace App\Application\EventHandler;

use App\Infrastructure\Messaging\EventEnvelope;
use Psr\Log\LoggerInterface;

// HandleEvent обрабатывает событие с полным структурированным логированием.
// В Symfony Monolog настраивается с JsonFormatter - каждая запись уходит
// как одна JSON-строка с атрибутами в context (для ELK/Loki).
final class AchievementHandler
{
    public function __construct(
        private readonly LoggerInterface $logger,
        private readonly AchievementProcessor $processor,
    ) {}

    public function handleEvent(EventEnvelope $envelope): void
    {
        $ctx = [
            'correlation_id' => $envelope->correlationId,
            'event_id' => $envelope->eventId,
            'event_type' => $envelope->type,
            'handler' => 'achievement_checker',
        ];

        $start = microtime(true);
        $this->logger->info('event processing started', $ctx);

        try {
            $this->processor->process($envelope);
        } catch (\Throwable $e) {
            $this->logger->error('event processing failed', $ctx + [
                'err' => $e->getMessage(),
                'duration_ms' => (int) round((microtime(true) - $start) * 1000),
            ]);
            throw $e;
        }

        $this->logger->info('event processing completed', $ctx + [
            'duration_ms' => (int) round((microtime(true) - $start) * 1000),
        ]);
    }
}

В PHP нет прямого аналога slog - структурированный JSON-лог собирается через Monolog c JsonFormatter (или LogstashFormatter) + processor-ы (PsrLogMessageProcessor, кастомный CorrelationIdProcessor). Атрибуты передаются через массив context, а Logger::withName() - ближайший аналог slog.With() для долгоживущего набора полей.

Результат в JSON (для ELK/Loki):

{
  "time": "2026-05-03T14:22:01Z",
  "level": "INFO",
  "msg": "event processing completed",
  "correlation_id": "01JXYZ...",
  "event_id": "01JABC...",
  "event_type": "lesson.completed",
  "handler": "achievement_checker",
  "duration": "2.341ms"
}

Теперь по correlation_id в Kibana вы найдёте: HTTP-запрос -> use-case -> outbox insert -> poller publish -> subscriber receive -> subscriber process. Вся цепочка.

Метрики с Prometheus

Числовые метрики показывают здоровье системы в реальном времени. Для event-driven системы критичны:

import "github.com/prometheus/client_golang/prometheus"

var (
    // Счётчик опубликованных событий
    eventsPublished = prometheus.NewCounterVec(
        prometheus.CounterOpts{
            Name: "events_published_total",
            Help: "Total number of events published to broker",
        },
        []string{"event_type"},
    )

    // Счётчик обработанных событий
    eventsProcessed = prometheus.NewCounterVec(
        prometheus.CounterOpts{
            Name: "events_processed_total",
            Help: "Total number of events processed by subscribers",
        },
        []string{"event_type", "handler", "status"}, // status: success, error
    )

    // Гистограмма времени обработки
    processingDuration = prometheus.NewHistogramVec(
        prometheus.HistogramOpts{
            Name:    "event_processing_duration_seconds",
            Help:    "Time spent processing an event",
            Buckets: []float64{0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1, 5},
        },
        []string{"event_type", "handler"},
    )

    // Gauge: количество событий в outbox, ожидающих отправки
    outboxPending = prometheus.NewGauge(
        prometheus.GaugeOpts{
            Name: "outbox_pending_events",
            Help: "Number of events in outbox waiting to be published",
        },
    )
)

func init() {
    prometheus.MustRegister(eventsPublished, eventsProcessed, processingDuration, outboxPending)
}
<?php
// src/Infrastructure/Metrics/EventMetrics.php
declare(strict_types=1);

namespace App\Infrastructure\Metrics;

use Prometheus\CollectorRegistry;
use Prometheus\Counter;
use Prometheus\Gauge;
use Prometheus\Histogram;

// Регистрация метрик через promphp/prometheus_client_php
// (CollectorRegistry, Counter, Gauge, Histogram).
final class EventMetrics
{
    private readonly Counter $eventsPublished;
    private readonly Counter $eventsProcessed;
    private readonly Histogram $processingDuration;
    private readonly Gauge $outboxPending;

    public function __construct(CollectorRegistry $registry)
    {
        // Счётчик опубликованных событий
        $this->eventsPublished = $registry->getOrRegisterCounter(
            'app',
            'events_published_total',
            'Total number of events published to broker',
            ['event_type'],
        );

        // Счётчик обработанных событий (status: success, error)
        $this->eventsProcessed = $registry->getOrRegisterCounter(
            'app',
            'events_processed_total',
            'Total number of events processed by subscribers',
            ['event_type', 'handler', 'status'],
        );

        // Гистограмма времени обработки (секунды)
        $this->processingDuration = $registry->getOrRegisterHistogram(
            'app',
            'event_processing_duration_seconds',
            'Time spent processing an event',
            ['event_type', 'handler'],
            [0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1, 5],
        );

        // Gauge: количество событий в outbox, ожидающих отправки
        $this->outboxPending = $registry->getOrRegisterGauge(
            'app',
            'outbox_pending_events',
            'Number of events in outbox waiting to be published',
        );
    }

    public function publishedInc(string $eventType): void
    {
        $this->eventsPublished->inc([$eventType]);
    }

    public function processedInc(string $eventType, string $handler, string $status): void
    {
        $this->eventsProcessed->inc([$eventType, $handler, $status]);
    }

    public function observeDuration(string $eventType, string $handler, float $seconds): void
    {
        $this->processingDuration->observe($seconds, [$eventType, $handler]);
    }

    public function setOutboxPending(int $count): void
    {
        $this->outboxPending->set($count);
    }
}

Использование в обработчике:

// InstrumentedHandler оборачивает обработчик и записывает метрики.
func InstrumentedHandler(handler, eventType string, fn func(ctx context.Context) error) func(ctx context.Context) error {
    return func(ctx context.Context) error {
        start := time.Now()

        err := fn(ctx)

        duration := time.Since(start).Seconds()
        status := "success"
        if err != nil {
            status = "error"
        }

        eventsProcessed.WithLabelValues(eventType, handler, status).Inc()
        processingDuration.WithLabelValues(eventType, handler).Observe(duration)

        return err
    }
}
<?php
// src/Infrastructure/Messaging/InstrumentedHandler.php
declare(strict_types=1);

namespace App\Infrastructure\Messaging;

use App\Infrastructure\Metrics\EventMetrics;

// InstrumentedHandler оборачивает обработчик и записывает метрики.
final class InstrumentedHandler
{
    public function __construct(
        private readonly EventMetrics $metrics,
        private readonly string $handler,
        private readonly string $eventType,
        /** @var callable(EventEnvelope): void */
        private $inner,
    ) {}

    public function __invoke(EventEnvelope $envelope): void
    {
        $start = microtime(true);
        $status = 'success';

        try {
            ($this->inner)($envelope);
        } catch (\Throwable $e) {
            $status = 'error';
            throw $e;
        } finally {
            $duration = microtime(true) - $start;
            $this->metrics->processedInc($this->eventType, $this->handler, $status);
            $this->metrics->observeDuration($this->eventType, $this->handler, $duration);
        }
    }
}

Distributed Tracing: OpenTelemetry

Метрики показывают агрегированную картину. Трейсы показывают конкретный путь конкретного события. OpenTelemetry позволяет создать span на каждом этапе: publish -> broker -> consume -> process.

import "go.opentelemetry.io/otel"

var tracer = otel.Tracer("event-processing")

// PublishWithTrace создаёт span при публикации события.
func (p *Publisher) PublishWithTrace(ctx context.Context, eventType string, payload []byte) error {
    ctx, span := tracer.Start(ctx, "event.publish",
        trace.WithAttributes(
            attribute.String("event.type", eventType),
        ),
    )
    defer span.End()

    // Inject trace context в headers сообщения
    // (для Kafka - в message headers, для NATS - в message metadata)
    err := p.broker.Publish(ctx, eventType, payload)
    if err != nil {
        span.SetStatus(codes.Error, err.Error())
    }
    return err
}

// ConsumeWithTrace создаёт span при обработке события.
func (s *Subscriber) ConsumeWithTrace(ctx context.Context, msg Message) error {
    // Extract trace context из headers сообщения
    ctx, span := tracer.Start(ctx, "event.process",
        trace.WithAttributes(
            attribute.String("event.type", msg.Type),
            attribute.String("event.id", msg.EventID),
        ),
    )
    defer span.End()

    err := s.handler(ctx, msg)
    if err != nil {
        span.SetStatus(codes.Error, err.Error())
    }
    return err
}
<?php
// src/Infrastructure/Tracing/TracedPublisher.php
declare(strict_types=1);

namespace App\Infrastructure\Tracing;

use App\Infrastructure\Messaging\BrokerInterface;
use OpenTelemetry\API\Globals;
use OpenTelemetry\API\Trace\SpanKind;
use OpenTelemetry\API\Trace\StatusCode;
use OpenTelemetry\Context\Context;

// TracedPublisher создаёт span при публикации события.
// Trace context инжектится в headers сообщения (TextMapPropagator)
// и извлекается в consumer-е - так получается единый trace publisher -> consumer.
final class TracedPublisher
{
    public function __construct(
        private readonly BrokerInterface $broker,
    ) {}

    public function publishWithTrace(string $eventType, string $payload, array $headers = []): void
    {
        $tracer = Globals::tracerProvider()->getTracer('event-processing');
        $span = $tracer->spanBuilder('event.publish')
 ->setSpanKind(SpanKind::KIND_PRODUCER)
 ->setAttribute('event.type', $eventType)
 ->startSpan();

        $scope = $span->activate();
        try {
            // Inject trace context в headers сообщения
            Globals::propagator()->inject($headers);
            $this->broker->publish($eventType, $payload, $headers);
        } catch (\Throwable $e) {
            $span->setStatus(StatusCode::STATUS_ERROR, $e->getMessage());
            throw $e;
        } finally {
            $scope->detach();
            $span->end();
        }
    }
}

// src/Infrastructure/Tracing/TracedSubscriber.php
namespace App\Infrastructure\Tracing;

// TracedSubscriber создаёт span при обработке события.
final class TracedSubscriber
{
    /**
     * @param  callable(string): void  $handler
     */
    public function __construct(
        private $handler,
    ) {}

    public function consumeWithTrace(string $eventType, string $eventId, array $headers, string $body): void
    {
        // Extract trace context из headers сообщения
        $parent = Globals::propagator()->extract($headers);

        $tracer = Globals::tracerProvider()->getTracer('event-processing');
        $span = $tracer->spanBuilder('event.process')
 ->setParent($parent)
 ->setSpanKind(SpanKind::KIND_CONSUMER)
 ->setAttribute('event.type', $eventType)
 ->setAttribute('event.id', $eventId)
 ->startSpan();

        $scope = $span->activate();
        try {
            ($this->handler)($body);
        } catch (\Throwable $e) {
            $span->setStatus(StatusCode::STATUS_ERROR, $e->getMessage());
            throw $e;
        } finally {
            $scope->detach();
            $span->end();
        }
    }
}

В Symfony Messenger пакет open-telemetry/sdk + middleware на bus делают то же автоматически: на отправке кладут traceparent в stamp, на consume извлекают и стартуют дочерний span. Ручной inject/extract обычно нужен только при работе с «голым» Kafka/AMQP-клиентом.

В Jaeger/Tempo вы увидите единый trace: HTTP request -> publish (3ms) -> consume (1ms) -> process (45ms) -> achievement granted. С таймингами каждого этапа.

На что настраивать алерты

Метрики бесполезны без алертов. Ключевые сигналы для event-driven системы:

Consumer lag - количество необработанных сообщений в очереди. Если lag растёт, consumers не справляются. Алерт: lag > 1000 в течение 5 минут.

Error rate - процент ошибок при обработке. Алерт: error rate > 5% за последние 10 минут.

Processing latency p99 - 99-й перцентиль времени обработки. Алерт: p99 > 5 секунд.

Outbox pending count - количество неотправленных событий в outbox. Если растёт - poller не работает или брокер недоступен. Алерт: pending > 100 в течение 2 минут.

Dead Letter Queue size - события, которые не удалось обработать после всех retry. Алерт: любое событие в DLQ требует внимания.

Dead Letter Queue: последний рубеж

Если событие не удалось обработать после N попыток, оно попадает в Dead Letter Queue (DLQ). DLQ - это не мусорка, а карантин. События в DLQ требуют ручного анализа: баг в коде, невалидные данные, недоступный внешний сервис. На уровне RabbitMQ DLQ настраивается аргументом x-dead-letter-exchange у основной очереди (см. урок CQRS).

// ProcessWithRetry обрабатывает событие с retry и отправкой в DLQ.
func (s *Subscriber) ProcessWithRetry(ctx context.Context, msg Message, maxRetries int) error {
    for attempt := 1; attempt <= maxRetries; attempt++ {
        err := s.handler(ctx, msg)
        if err == nil {
            return nil
        }

        slog.WarnContext(ctx, "event processing retry",
            slog.String("event_id", msg.EventID),
            slog.Int("attempt", attempt),
            slog.Int("max_retries", maxRetries),
            slog.String("err", err.Error()),
        )

        // Exponential backoff
        time.Sleep(time.Duration(attempt*attempt) * 100 * time.Millisecond)
    }

    // Все попытки исчерпаны - отправляем в DLQ
    slog.ErrorContext(ctx, "event moved to DLQ",
        slog.String("event_id", msg.EventID),
        slog.String("event_type", msg.Type),
    )
    return s.dlq.Send(ctx, msg)
}
<?php
// src/Infrastructure/Messaging/ProcessWithRetry.php
declare(strict_types=1);

namespace App\Infrastructure\Messaging;

use Psr\Log\LoggerInterface;

// ProcessWithRetry обрабатывает событие с retry и отправкой в DLQ.
final class ProcessWithRetry
{
    public function __construct(
        /** @var callable(Message): void */
        private $handler,
        private readonly DeadLetterQueue $dlq,
        private readonly LoggerInterface $logger,
        private readonly int $maxRetries = 3,
    ) {}

    public function process(Message $msg): void
    {
        for ($attempt = 1; $attempt <= $this->maxRetries; $attempt++) {
            try {
                ($this->handler)($msg);

                return;
            } catch (\Throwable $e) {
                $this->logger->warning('event processing retry', [
                    'event_id' => $msg->eventId,
                    'attempt' => $attempt,
                    'max_retries' => $this->maxRetries,
                    'err' => $e->getMessage(),
                ]);

                // Exponential backoff (квадратичный, миллисекунды)
                usleep($attempt * $attempt * 100_000);
            }
        }

        // Все попытки исчерпаны - отправляем в DLQ
        $this->logger->error('event moved to DLQ', [
            'event_id' => $msg->eventId,
            'event_type' => $msg->type,
        ]);

        $this->dlq->send($msg);
    }
}

В Symfony Messenger вся retry-логика + DLQ делается декларативно: retry_strategy.max_retries, failure_transport - сообщение после исчерпания попыток автоматически уезжает в failure-транспорт (он же DLQ), а messenger:failed:show/messenger:failed:retry - консольные команды для разбора инцидентов.

Первое, что нужно сделать при переходе на event-driven архитектуру - добавить correlation_id в каждое событие и каждый лог. Без этого вы не сможете связать HTTP-запрос пользователя с цепочкой событий, которые он вызвал. Добавьте его до того, как система выйдет в production.

Типичные ошибки

  • correlation_id теряется при ре-публикации - subscriber обрабатывает событие A, генерирует событие B, кладёт новый correlation_id. Цепочка обрывается. Пробрасывай тот же correlation_id из входящего события в исходящие - единый «trace» через всю цепочку.
  • Алёрт на «есть события в DLQ» - каждое случайное падение consumer-а будит дежурного. Алёртить не на сам факт, а на rate (≥ N в минуту) или age старейшего письма (висит > 30 мин - действительно проблема).
  • Метрика events_processed_total без label-а handler - все subscriber-ы сливаются в один счётчик, не понять, какой именно тормозит. Размечай event_type, handler, result=ok|error|retry.
  • outbox_pending алёрт по абсолютному значению - при ночном batch-импорте 10K событий ложно срабатывает. Алёртить по age старейшей не-опубликованной записи (now - min(created_at) > 5min).
  • Логирование payload целиком - событие с email/телефоном/токеном попадает в логи, потом в централизованный сборщик. GDPR-инцидент. Логируй event_type, event_id, aggregate_id, метаданные - не содержимое payload.
  • Tracing включён только для HTTP, не для event consumer-ов - в Jaeger видно, что HTTP-запрос отдал ответ за 30ms, а почему email пришёл через 10 минут - никаких следов. OTel-инструментирование нужно и в consumer-ах: tracer.Start() на каждой обработке + извлечение span context из header события.
  • DLQ как «folder для забытых писем» - события попадают, никто не смотрит. Через год - десятки тысяч писем, разбор невозможен. Дежурный обязан проверять DLQ ежедневно, корневые причины фиксить, не накапливать.
  • Метрики обновляются «когда-нибудь после обработки» - defer metrics.Inc() без учёта успех/ошибка. Отметка error всё равно ставит processed_total++. Делай отдельные processed_ok_total и processed_error_total, без агрегации.

Мини-задание

  • Добавь correlation_id в EventEnvelope и HTTP-middleware, который его генерирует
  • Настрой slog так, чтобы каждый лог содержал correlation_id, event_id и event_type
  • Зарегистрируй Prometheus-метрики: events_processed_total, event_processing_duration_seconds, outbox_pending_events
  • Реализуй ProcessWithRetry с exponential backoff и отправкой в DLQ после 3 попыток
  • Опиши 3 алерта для своей системы: на что срабатывают и какой порог

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.