Архитектура

Построение event-driven системы на Go с Redis Pub/Sub и PostgreSQL

Ruslan Ismailov Опубликовано 18 мин чтения
П

Введение в event-driven архитектуру

Event-driven архитектура (EDA) — это подход к проектированию систем, при котором компоненты общаются через события, а не через прямые вызовы. Вместо того чтобы сервис A напрямую вызывал метод сервиса B, он публикует событие «что-то произошло», а сервис B (или несколько сервисов) реагирует на него независимо.

Преимущества такого подхода очевидны для любого опытного разработчика:

  • Слабая связанность: издатель не знает о подписчиках и не зависит от их доступности.
  • Горизонтальная масштабируемость: подписчики масштабируются независимо от издателей.
  • Отказоустойчивость: временная недоступность одного компонента не блокирует остальные.
  • Аудит и воспроизводимость: сохранённые события позволяют восстановить состояние системы в любой момент.

В 2026 году event-driven архитектура на Go стала стандартом де-факто для высоконагруженных микросервисных систем. В этой статье мы построим полноценную систему с нуля, используя Redis Pub/Sub в качестве шины событий и PostgreSQL для персистентности.

Обзор инструментов: Go, Redis Pub/Sub, PostgreSQL

Каждый компонент нашего стека выполняет строго определённую роль.

Go — ядро системы

Go идеально подходит для event-driven систем благодаря нативной поддержке конкурентности через горутины и каналы, минималистичной стандартной библиотеке и низкому потреблению памяти. Компилируемый тип, строгая типизация и быстрый старт делают его отличным выбором для реализации как publisher, так и subscriber.

Redis Pub/Sub — шина событий

Redis в режиме Pub/Sub обеспечивает механизм «опубликовал и забыл» (fire-and-forget) с минимальными задержками. Он не сохраняет сообщения — если подписчик недоступен в момент публикации, сообщение теряется. Это ключевое ограничение, которое мы учтём в архитектуре.

PostgreSQL — хранилище событий

PostgreSQL выступает в роли надёжного персистентного хранилища: мы будем записывать каждое событие в таблицу audit log до его публикации в Redis. Это даёт нам возможность восстановить пропущенные события и обеспечить гарантии доставки.

Проектирование событий: структура, схема, версионирование

Событие — это неизменяемый факт, произошедший в системе. Хорошо спроектированное событие содержит достаточно данных для обработки без дополнительных запросов.

Базовая структура события в Go:

package events

import (
    "time"
    "github.com/google/uuid"
)

// EventType задаёт тип события
type EventType string

const (
    UserCreated   EventType = "user.created"
    UserUpdated   EventType = "user.updated"
    OrderPlaced   EventType = "order.placed"
    PaymentFailed EventType = "payment.failed"
)

// BaseEvent — общая обёртка для всех событий
type BaseEvent struct {
    ID          string      `json:"id"`           // уникальный идентификатор события
    Type        EventType   `json:"type"`         // тип события
    Version     int         `json:"version"`      // версия схемы
    OccurredAt  time.Time   `json:"occurred_at"` // время возникновения
    Source      string      `json:"source"`       // сервис-источник
    Payload     interface{} `json:"payload"`      // данные события
}

// NewEvent создаёт новое событие с заполненными метаданными
func NewEvent(eventType EventType, source string, version int, payload interface{}) BaseEvent {
    return BaseEvent{
        ID:         uuid.New().String(),
        Type:       eventType,
        Version:    version,
        OccurredAt: time.Now().UTC(),
        Source:     source,
        Payload:    payload,
    }
}

// UserCreatedPayload — payload для события создания пользователя
type UserCreatedPayload struct {
    UserID    string `json:"user_id"`
    Email     string `json:"email"`
    Name      string `json:"name"`
    CreatedAt string `json:"created_at"`
}

Версионирование событий критично для долгоживущих систем. Используйте поле version и обрабатывайте миграцию схем явно в обработчиках. Рекомендуется следовать принципу аддитивных изменений: добавляйте новые поля, не удаляя старые.

Реализация publisher на Go

Publisher отвечает за публикацию событий в Redis-канал. Важный паттерн — сначала сохраняем событие в PostgreSQL (outbox pattern), затем публикуем в Redis. Это обеспечивает durability.

package publisher

import (
    "context"
    "encoding/json"
    "fmt"
    "log"

    "github.com/go-redis/redis/v9"
    "github.com/jmoiron/sqlx"

    "myapp/events"
)

// Publisher публикует события в Redis и сохраняет их в PostgreSQL
type Publisher struct {
    redis *redis.Client
    db    *sqlx.DB
}

func NewPublisher(redisClient *redis.Client, db *sqlx.DB) *Publisher {
    return &Publisher{
        redis: redisClient,
        db:    db,
    }
}

// Publish сохраняет событие в БД и публикует его в Redis
func (p *Publisher) Publish(ctx context.Context, event events.BaseEvent) error {
    // 1. Сериализуем событие
    payload, err := json.Marshal(event)
    if err != nil {
        return fmt.Errorf("marshal event: %w", err)
    }

    // 2. Сохраняем в PostgreSQL (outbox / audit log)
    if err := p.saveEventToDB(ctx, event, payload); err != nil {
        return fmt.Errorf("save event to db: %w", err)
    }

    // 3. Публикуем в Redis Pub/Sub
    channel := string(event.Type)
    if err := p.redis.Publish(ctx, channel, payload).Err(); err != nil {
        // Не фатально — событие уже в БД, relay-worker переотправит
        log.Printf("WARN: failed to publish to Redis channel %s: %v", channel, err)
        return nil
    }

    log.Printf("INFO: published event %s (id=%s) to channel %s", event.Type, event.ID, channel)
    return nil
}

// saveEventToDB сохраняет событие в таблицу event_outbox
func (p *Publisher) saveEventToDB(ctx context.Context, event events.BaseEvent, payload []byte) error {
    query := `
        INSERT INTO event_outbox (
            id, event_type, version, source, occurred_at, payload, published
        ) VALUES (
            $1, $2, $3, $4, $5, $6, false
        )
    `
    _, err := p.db.ExecContext(ctx, query,
        event.ID,
        string(event.Type),
        event.Version,
        event.Source,
        event.OccurredAt,
        payload,
    )
    return err
}

Relay Worker: от outbox к Redis

Чтобы события, которые не попали в Redis из-за сбоя, всё же были доставлены, реализуем фоновый worker, который опрашивает outbox и ретранслирует неотправленные события:

package publisher

import (
    "context"
    "log"
    "time"
)

// RelayWorker переотправляет события из outbox в Redis
func (p *Publisher) RelayWorker(ctx context.Context) {
    ticker := time.NewTicker(5 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return
        case <-ticker.C:
            if err := p.relayPendingEvents(ctx); err != nil {
                log.Printf("ERROR: relay worker: %v", err)
            }
        }
    }
}

func (p *Publisher) relayPendingEvents(ctx context.Context) error {
    rows, err := p.db.QueryContext(ctx, `
        SELECT id, event_type, payload
        FROM event_outbox
        WHERE published = false
        ORDER BY occurred_at
        LIMIT 100
        FOR UPDATE SKIP LOCKED
    `)
    if err != nil {
        return err
    }
    defer rows.Close()

    for rows.Next() {
        var id, eventType string
        var payload []byte
        if err := rows.Scan(&id, &eventType, &payload); err != nil {
            continue
        }

        if err := p.redis.Publish(ctx, eventType, payload).Err(); err != nil {
            log.Printf("WARN: relay failed for event %s: %v", id, err)
            continue
        }

        // Отмечаем как опубликованное
        _, _ = p.db.ExecContext(ctx,
            `UPDATE event_outbox SET published = true, published_at = NOW() WHERE id = $1`,
            id,
        )
    }
    return rows.Err()
}

Реализация subscriber и обработчиков событий

Subscriber подписывается на каналы Redis и диспетчеризует события по зарегистрированным обработчикам. Используем паттерн registry для гибкой регистрации обработчиков:

package subscriber

import (
    "context"
    "encoding/json"
    "fmt"
    "log"

    "github.com/go-redis/redis/v9"

    "myapp/events"
)

// HandlerFunc — тип функции-обработчика события
type HandlerFunc func(ctx context.Context, event events.BaseEvent) error

// Subscriber управляет подписками на события
type Subscriber struct {
    redis    *redis.Client
    handlers map[events.EventType][]HandlerFunc
}

func NewSubscriber(redisClient *redis.Client) *Subscriber {
    return &Subscriber{
        redis:    redisClient,
        handlers: make(map[events.EventType][]HandlerFunc),
    }
}

// Register регистрирует обработчик для определённого типа события
func (s *Subscriber) Register(eventType events.EventType, handler HandlerFunc) {
    s.handlers[eventType] = append(s.handlers[eventType], handler)
}

// Listen запускает прослушивание каналов Redis
func (s *Subscriber) Listen(ctx context.Context) error {
    channels := make([]string, 0, len(s.handlers))
    for eventType := range s.handlers {
        channels = append(channels, string(eventType))
    }

    if len(channels) == 0 {
        return fmt.Errorf("no channels registered")
    }

    pubsub := s.redis.Subscribe(ctx, channels...)
    defer pubsub.Close()

    log.Printf("INFO: subscribed to channels: %v", channels)

    msgCh := pubsub.Channel()
    for {
        select {
        case <-ctx.Done():
            log.Println("INFO: subscriber shutting down")
            return nil
        case msg, ok := <-msgCh:
            if !ok {
                return fmt.Errorf("subscription channel closed")
            }
            go s.dispatch(ctx, msg.Channel, []byte(msg.Payload))
        }
    }
}

// dispatch десериализует событие и вызывает обработчики
func (s *Subscriber) dispatch(ctx context.Context, channel string, payload []byte) {
    var event events.BaseEvent
    if err := json.Unmarshal(payload, &event); err != nil {
        log.Printf("ERROR: unmarshal event on channel %s: %v", channel, err)
        return
    }

    handlers, ok := s.handlers[event.Type]
    if !ok {
        log.Printf("WARN: no handlers for event type %s", event.Type)
        return
    }

    for _, handler := range handlers {
        if err := handler(ctx, event); err != nil {
            log.Printf("ERROR: handler for %s failed: %v", event.Type, err)
        }
    }
}

Пример конкретного обработчика

package handlers

import (
    "context"
    "encoding/json"
    "log"

    "myapp/events"
)

// EmailHandler отправляет приветственное письмо при создании пользователя
type EmailHandler struct {
    emailService EmailService
}

func NewEmailHandler(es EmailService) *EmailHandler {
    return &EmailHandler{emailService: es}
}

func (h *EmailHandler) Handle(ctx context.Context, event events.BaseEvent) error {
    var payload events.UserCreatedPayload
    if err := remarshal(event.Payload, &payload); err != nil {
        return fmt.Errorf("parse UserCreatedPayload: %w", err)
    }

    log.Printf("INFO: sending welcome email to %s", payload.Email)
    return h.emailService.SendWelcome(ctx, payload.Email, payload.Name)
}

// remarshal преобразует interface{} payload в конкретный тип
func remarshal(src interface{}, dst interface{}) error {
    data, err := json.Marshal(src)
    if err != nil {
        return err
    }
    return json.Unmarshal(data, dst)
}

Гарантии доставки: at-least-once и идемпотентность

Redis Pub/Sub предоставляет гарантию at-most-once: сообщение доставляется ноль или один раз. Для достижения at-least-once мы используем комбинацию outbox pattern и relay worker, описанных выше. Однако при retry возможны дублирующиеся события, поэтому обработчики должны быть идемпотентными.

Паттерн обеспечения идемпотентности — таблица обработанных событий:

-- Таблица для дедупликации
CREATE TABLE IF NOT EXISTS processed_events (
    event_id    UUID PRIMARY KEY,
    handler     VARCHAR(255) NOT NULL,
    processed_at TIMESTAMPTZ DEFAULT NOW()
);

-- Индекс для быстрого поиска
CREATE INDEX IF NOT EXISTS idx_processed_events_handler
    ON processed_events (handler, event_id);
package handlers

import (
    "context"
    "database/sql"
    "errors"
    "fmt"

    "github.com/jmoiron/sqlx"
    "myapp/events"
)

// IdempotentHandler оборачивает обработчик для обеспечения идемпотентности
type IdempotentHandler struct {
    db          *sqlx.DB
    handlerName string
    inner       func(ctx context.Context, event events.BaseEvent) error
}

func NewIdempotentHandler(
    db *sqlx.DB,
    name string,
    inner func(ctx context.Context, event events.BaseEvent) error,
) *IdempotentHandler {
    return &IdempotentHandler{db: db, handlerName: name, inner: inner}
}

func (h *IdempotentHandler) Handle(ctx context.Context, event events.BaseEvent) error {
    // Проверяем, было ли событие уже обработано
    var exists bool
    err := h.db.QueryRowContext(ctx,
        `SELECT EXISTS(SELECT 1 FROM processed_events WHERE event_id = $1 AND handler = $2)`,
        event.ID, h.handlerName,
    ).Scan(&exists)
    if err != nil && !errors.Is(err, sql.ErrNoRows) {
        return fmt.Errorf("check idempotency: %w", err)
    }

    if exists {
        // Событие уже обработано — пропускаем
        return nil
    }

    // Выполняем бизнес-логику
    if err := h.inner(ctx, event); err != nil {
        return err
    }

    // Помечаем как обработанное
    _, err = h.db.ExecContext(ctx,
        `INSERT INTO processed_events (event_id, handler) VALUES ($1, $2) ON CONFLICT DO NOTHING`,
        event.ID, h.handlerName,
    )
    return err
}

Сохранение событий в PostgreSQL как audit log

Схема таблиц для outbox и audit log:

-- Таблица для Outbox Pattern
CREATE TABLE IF NOT EXISTS event_outbox (
    id           UUID PRIMARY KEY,
    event_type   VARCHAR(255) NOT NULL,
    version      INTEGER NOT NULL DEFAULT 1,
    source       VARCHAR(255) NOT NULL,
    occurred_at  TIMESTAMPTZ NOT NULL,
    payload      JSONB NOT NULL,
    published    BOOLEAN NOT NULL DEFAULT false,
    published_at TIMESTAMPTZ,
    created_at   TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX idx_event_outbox_unpublished
    ON event_outbox (occurred_at)
    WHERE published = false;

-- Audit log для аналитики и восстановления состояния
CREATE TABLE IF NOT EXISTS event_audit_log (
    id           BIGSERIAL PRIMARY KEY,
    event_id     UUID NOT NULL,
    event_type   VARCHAR(255) NOT NULL,
    version      INTEGER NOT NULL,
    source       VARCHAR(255) NOT NULL,
    occurred_at  TIMESTAMPTZ NOT NULL,
    payload      JSONB NOT NULL,
    created_at   TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX idx_event_audit_log_event_type
    ON event_audit_log (event_type, occurred_at DESC);

CREATE INDEX idx_event_audit_log_event_id
    ON event_audit_log (event_id);

Функция записи в audit log вызывается из subscriber после успешной обработки события:

func SaveToAuditLog(ctx context.Context, db *sqlx.DB, event events.BaseEvent) error {
    payload, err := json.Marshal(event.Payload)
    if err != nil {
        return err
    }
    _, err = db.ExecContext(ctx, `
        INSERT INTO event_audit_log
            (event_id, event_type, version, source, occurred_at, payload)
        VALUES ($1, $2, $3, $4, $5, $6)
        ON CONFLICT DO NOTHING
    `, event.ID, event.Type, event.Version, event.Source, event.OccurredAt, payload)
    return err
}

Запуск и оркестрация через Docker Compose

Docker Compose позволяет поднять всё окружение одной командой. Ниже — полноценный файл для разработки:

version: '3.9'

services:
  postgres:
    image: postgres:16-alpine
    container_name: eda_postgres
    environment:
      POSTGRES_USER: eda_user
      POSTGRES_PASSWORD: eda_secret
      POSTGRES_DB: eda_db
    ports:
      - "5432:5432"
    volumes:
      - postgres_data:/var/lib/postgresql/data
      - ./migrations:/docker-entrypoint-initdb.d
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U eda_user -d eda_db"]
      interval: 5s
      timeout: 5s
      retries: 5

  redis:
    image: redis:7-alpine
    container_name: eda_redis
    ports:
      - "6379:6379"
    command: redis-server --save 60 1 --loglevel warning
    volumes:
      - redis_data:/data
    healthcheck:
      test: ["CMD", "redis-cli", "ping"]
      interval: 5s
      timeout: 3s
      retries: 5

  publisher:
    build:
      context: .
      dockerfile: ./cmd/publisher/Dockerfile
    container_name: eda_publisher
    environment:
      DATABASE_URL: postgres://eda_user:eda_secret@postgres:5432/eda_db?sslmode=disable
      REDIS_URL: redis:6379
      SERVICE_NAME: publisher-service
    depends_on:
      postgres:
        condition: service_healthy
      redis:
        condition: service_healthy

  subscriber:
    build:
      context: .
      dockerfile: ./cmd/subscriber/Dockerfile
    container_name: eda_subscriber
    environment:
      DATABASE_URL: postgres://eda_user:eda_secret@postgres:5432/eda_db?sslmode=disable
      REDIS_URL: redis:6379
      SERVICE_NAME: subscriber-service
    depends_on:
      postgres:
        condition: service_healthy
      redis:
        condition: service_healthy
    deploy:
      replicas: 2

volumes:
  postgres_data:
  redis_data:

Тестирование event-driven системы

Тестирование EDA требует особого подхода. Разделяем тесты на три уровня: unit-тесты для обработчиков, интеграционные тесты с реальным Redis и PostgreSQL, и end-to-end тесты потоков событий.

Unit-тесты обработчиков

package handlers_test

import (
    "context"
    "testing"
    "time"

    "github.com/stretchr/testify/assert"
    "github.com/stretchr/testify/mock"

    "myapp/events"
    "myapp/handlers"
)

type MockEmailService struct {
    mock.Mock
}

func (m *MockEmailService) SendWelcome(ctx context.Context, email, name string) error {
    args := m.Called(ctx, email, name)
    return args.Error(0)
}

func TestEmailHandler_Handle(t *testing.T) {
    mockES := new(MockEmailService)
    mockES.On("SendWelcome", mock.Anything, "user@example.com", "Alice").Return(nil)

    handler := handlers.NewEmailHandler(mockES)

    payload := events.UserCreatedPayload{
        UserID: "123",
        Email:  "user@example.com",
        Name:   "Alice",
    }
    event := events.NewEvent(events.UserCreated, "user-service", 1, payload)

    err := handler.Handle(context.Background(), event)

    assert.NoError(t, err)
    mockES.AssertExpectations(t)
}

func TestIdempotentHandler_SkipsDuplicate(t *testing.T) {
    // Интеграционный тест с testcontainers-go или sqlmock
    // Первый вызов — обрабатывает событие
    // Второй вызов с тем же event.ID — пропускает
    // Подробная реализация зависит от тестовой инфраструктуры
    t.Log("See integration tests for full idempotency coverage")
}

Интеграционные тесты с testcontainers-go

package integration_test

import (
    "context"
    "testing"
    "time"

    "github.com/stretchr/testify/require"
    "github.com/testcontainers/testcontainers-go"
    "github.com/testcontainers/testcontainers-go/modules/redis"
    "github.com/testcontainers/testcontainers-go/modules/postgres"
)

func TestPublishSubscribeFlow(t *testing.T) {
    ctx := context.Background()

    // Поднимаем Redis контейнер
    redisContainer, err := redis.RunContainer(ctx,
        testcontainers.WithImage("redis:7-alpine"),
    )
    require.NoError(t, err)
    defer redisContainer.Terminate(ctx)

    // Поднимаем PostgreSQL контейнер
    pgContainer, err := postgres.RunContainer(ctx,
        testcontainers.WithImage("postgres:16-alpine"),
        postgres.WithDatabase("test_db"),
        postgres.WithUsername("test"),
        postgres.WithPassword("test"),
    )
    require.NoError(t, err)
    defer pgContainer.Terminate(ctx)

    // ... инициализация publisher и subscriber
    // ... публикация тестового события
    // ... проверка обработки через канал или WaitGroup с таймаутом

    received := make(chan events.BaseEvent, 1)
    // subscriber.Register(events.UserCreated, func(ctx context.Context, e events.BaseEvent) error {
    //     received <- e
    //     return nil
    // })

    select {
    case event := <-received:
        require.Equal(t, events.UserCreated, event.Type)
    case <-time.After(5 * time.Second):
        t.Fatal("timeout waiting for event")
    }
}

Ограничения Redis Pub/Sub и когда смотреть в сторону Kafka

Redis Pub/Sub — отличный инструмент для определённых сценариев, но у него есть принципиальные ограничения, которые необходимо понимать.

  • Нет персистентности сообщений: если подписчик не подключён в момент публикации, сообщение теряется безвозвратно. Именно поэтому мы используем outbox pattern.
  • Нет consumer groups: все подписчики на канал получают копию каждого сообщения. Невозможно распределить нагрузку между несколькими экземплярами одного сервиса без дополнительной логики.
  • Нет replay: нельзя переиграть события с определённого offset, как в Kafka.
  • Ограниченная пропускная способность: при очень высоких нагрузках (миллионы сообщений в секунду) Redis становится узким местом.

Используйте Redis Streams (XADD/XREADGROUP) если вам нужны consumer groups и базовая персистентность, оставаясь в экосистеме Redis. Переходите на Apache Kafka, когда требуется гарантированная доставка, replay событий, обработка миллионов сообщений в секунду или долгосрочное хранение event log.

В микросервисной архитектуре Redis Pub/Sub отлично подходит для: внутренних уведомлений с низкой латентностью, кэш-инвалидации, нотификаций в реальном времени и случаев, где потеря редкого сообщения некритична при наличии outbox-компенсации. Для финансовых транзакций, критичных бизнес-событий и систем с требованиями compliance — выбирайте Kafka или другие гарантированные брокеры.

Заключение

Мы построили полноценную event-driven систему на Go, объединив Redis Pub/Sub как легковесную шину событий с PostgreSQL для надёжной персистентности через outbox pattern. Ключевые паттерны, которые мы применили:

  • Outbox Pattern — сохраняем событие в БД до публикации в Redis, обеспечивая at-least-once доставку.
  • Idempotent Handlers — таблица обработанных событий защищает от дублирования при retry.
  • Event Registry — гибкая регистрация множества обработчиков на один тип события.
  • Relay Worker — фоновый процесс, который ретранслирует события при восстановлении Redis.

Асинхронная обработка событий на Go в 2026 году — это не экзотика, а практическая необходимость для масштабируемых микросервисных систем. Начните с Redis Pub/Sub + PostgreSQL и мигрируйте на Kafka только когда реально упрётесь в ограничения.

Технологии

Теги

Руслан Исмаилов

Senior Web / Backend разработчик. Senior web/backend разработчик с 9-летним опытом. Стек: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, микросервисы, CI/CD. Подробнее обо мне →