Построение event-driven системы на Go с Redis Pub/Sub и PostgreSQL
Введение в 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. Подробнее обо мне →