Наблюдаемость: трассировка, корреляция, метрики и алерты
Наблюдаемость: трассировка, корреляция, метрики и алерты
В синхронной системе отладка проста: запрос приходит, проходит через middleware, handler, service, repository - и ответ уходит. Один HTTP-запрос, одна горутина, один стектрейс. В event-driven системе одно действие пользователя может породить цепочку из десятка событий, обработанных разными подписчиками, возможно в разных сервисах, с задержкой в миллисекунды или минуты. Без специальных инструментов наблюдаемости вы будете искать причину бага, перебирая логи вручную.
Три столпа наблюдаемости (Three Pillars of Observability) - логи, метрики, трейсы - особенно критичны для event-driven архитектуры.
Correlation ID: связываем цепочку событий
Correlation ID - уникальный идентификатор, который генерируется при входящем запросе и передаётся через все события и подписчики. Он позволяет собрать все логи одного пользовательского действия в один trace.
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 - консольные команды для разбора инцидентов.
Типичные ошибки
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 алерта для своей системы: на что срабатывают и какой порог