golang

System Design на практике: создаем микросервис генерации уникальных идентификаторов

  • понедельник, 24 августа 2026 г. в 00:00:15
https://habr.com/ru/articles/1068960/

Это вторая статья серии, посвященной проектированию системы сокращения ссылок. В предыдущей мы спроектировали архитектуру сервиса, разделил систему на микросервисы, определили зоны ответственности и описали взаимодействие. Сегодня перейдем от теории к практике и займемся основой сокращателя ссылок — микросервисом генерации уникальных идентификаторов. Почему генерацию идентификаторов стоит рассматривать в отдельной статье? На первый взгляд задача кажется тривиальной: взять строку из случайных символов или авто инкремент из базы данных. Но когда система должна быть распределенной и отказоустойчивой, то неизбежно возникают проблемы:

  • Использовать первичный ключ БД с атрибутом auto_increment проблематично в распределенном окружении.

  • Можно генерировать UUID, который имеет очень низкую, но всё же не нулевую вероятность повторения, и при этом занимают слишком много места (128 бит) и генерируют слишком длинные ссылки, что сводит на нет саму идею сокращателя.

  • Рандомные строки приводят к коллизиям и требуют постоянных проверок в хранилище идентификаторов, что тормозит систему в целом.

Нам нужен генератор, который выдает уникальные, короткие, последовательные ID с минимальной задержкой и не создает единой точки отказа.

Сегодня мы создадим такой генератор ID и разберем:

  1. Подходы к генерации уникальных идентификаторов.

  2. Генерацию на основе алгоритма Snowflake.

  3. Как выбрать минимальную длину ID чтобы ссылки оставались максимально короткими.

  4. Как спроектировать сервис чтобы он генерировал ID с предвыборкой ограниченной длины.

  5. Как обеспечить адаптацию к нагрузке и сделать систему отказоустойчивой, чтобы падение одного узла не остановило генерацию во всей системе.

  6. Как уберечь сервис от зависания запросов и падения при возникновении критических ошибок.

  7. Как дождаться завершения работы сервиса.

  8. Архитектурные вопросы и напишем тестируемый прототип микросервиса на языке Go.

Процесс проектирования должен начинаться со сбора требований, предъявляемых к системе. Рассмотрим основные.

Функциональные требования:

  1. Генерация уникального ID для исходной ссылки.

  2. Длина ID (количество занимаемых бит) должна быть минимальной для компактности итоговой ссылки.

Нефункциональные требования:

  1. Высокая доступность. Система должна продолжать работать при падении одного из серверов.

  2. Низкая задержка. Генерация ID должна занимать как можно меньше времени.

  3. Масштабируемость как возможность легко добавлять новые узлы генератора при росте трафика.

  4. Уникальность ID — исключение коллизий.

Рассмотрим 2 основных подхода решения задачи генерации ID, их плюсы и минусы.

Сервер тикетов

Это один из способов генерации уникальных ID. Давайте посмотрим, как они работают. Суть в том, что мы используем функцию auto_increment в отдельно взятом сервере баз данных. auto_increment — это встроенный инструмент БД, который автоматически создает новое последовательное число для каждой новой строки.

Создается таблица и для колонки с ID указывается свойство auto_increment. СУБД запускает внутренний счетчик для этой таблицы. По умолчанию он равен 1. Когда вы добавляете новую запись, СУБД берет текущее значение, отдает его новой строке, а сам счетчик увеличивает на единицу (шаг по умолчанию равен 1).

На базе этой функции и строится сервер тикетов. Схема выглядит так:

  1. Создается специальная таблица всего с двумя полями: id (свойство auto_increment) и, например, простая строка — заполнитель, не хранящая значимых данных.

  2. Микросервис посылает запрос на вставку строки в таблицу.

  3. База данных добавляет строку, увеличивает id на единицу и возвращает этот id микросервису, который использует это число как уникальный идентификатор.

У такого подхода возникают проблемы при масштабировании. Если использовать обычный auto_increment на одном сервере, то он станет единой точкой отказа. Для решения этой проблемы в распределенных системах делают так:

  • Используют несколько серверов тикетов (например, два).

  • На сервере № 1 настраивают шаг auto_increment = 2 и старт с 1. Он будет выдавать только нечетные числа: 1, 3, 5, 7 и так далее

  • На сервере № 2 настраивают шаг auto_increment = 2 и старт с 2. Он будет выдавать только четные числа: 2, 4, 6, 8 и так далее

Если один сервер № 1 упадет, то система продолжит работать, выдавая уникальные ID со второго сервера.

Основным преимуществом такого подхода является простота генерации числовых идентификаторов.

Использование сервера тикетов в микросервисах порождает следующие проблемы:

  1. Производительность. При генерации ID база данных использует внутренние блокировки счетчика. Когда микросервисы одновременно запрашивают новые ID, они встают в очередь. Скорость системы падает.

  2. Сетевые задержки. Микросервисам приходится делать лишний сетевой запрос к БД только ради получения ID перед тем, как сохранить данные у себя. Это сильно увеличивает общее время обработки запроса.

Snowflake ID

Это децентрализованный алгоритм генерации уникальных целочисленных идентификаторов для распределенных систем. Каждый ID содержит 64 бита и состоит из нескольких логических блоков:

  1. Знак (Sign 1 бит). Всегда равен 0 для обеспечения положительного значения числа.

  2. Временная метка (Timestamp, 41 бит). Количество миллисекунд, прошедших с некоторой заданной точки отсчета (запуска системы), например, 1 августа 2026 года. Позволяет сортировать ID по времени. Максимальная временная метка, которую можно представить с помощью 41 бита позволяет работать системе примерно 69 лет.

  3. Идентификатор дата — центра (Data center, 5 бит) и сервера в нём (Worker, 5 бит). Всего 10 бит, позволяющих различать до 1024 независимых узлов генерации в сети.

  4. Порядковый номер последовательности. Счетчик для ID, генерируемых в течение одной миллисекунды. Позволяет генерировать до 4096 уникальных ID на одном узле за одну миллисекунду.

Преимущества алгоритма:

  1. Уникальность без координации. Узлы генерируют ID автономно, не обращаясь к БД.

  2. Хронологическая сортировка. Идентификаторы «естественным» образом упорядочены по времени создания.

  3. Эффективность. Занимает 8 байт (64 бита), а не 16 как UUID.

Минусы алгоритма:

  1. Проблема перевода назад часов на сервере из — за чего ID могут дублироваться.

  2. Ограничение на количество серверов.

По моему мнению, плюсы алгоритма Snowflake ID перевешивают его минусы. Поэтому наш микросервис будет использовать этот алгоритм.

Для начала давайте вспомним требования, предъявляемые к генератору, которые описывались в первой статье:

  1. Количество операции записи в день. Генерируется 100 миллионов URL адресов.

  2. Количество операций записи в секунду: 100 миллионов / 24 / 3600 = 1160.

  3. Пусть сервис для сокращения URL адресов проработает 8 лет, тогда мы должны поддерживать хранение 100 миллионов ID в день * 365 дней * 8 лет = 292 миллиардов записей.

Чтобы сократить исходный URL адрес, нужно реализовать механизм, который преобразует его в строку, которая должна быть как можно короче. В литературе такой механизм называют хеш — функцией.

Пусть строковое представление ID состоит из символов [0 — 9, a — z, A — Z]. Общее число различных симовлов равно 10 + 26 + 26 = 62. Т.е. мы используем систему счисления с основанием 62. Чтобы определить длину ID, нужно найти наименьшее n, при котором 62^n ≥ 292 миллиардов.

Рассмотрим последовательные значения n = 6,7,8:

62^6 = 56 800 235 584

62^7 = 3 521 614 606 208 ~3,5 триллиона

62^8 = 218 340 105 584 896

3,5 триллиона при n = 62 ^ 7 более чем достаточно для хранения 292 миллиардов URL адресов. Поэтому длина ID будет равна 7. Если взять общеизвестные функции хеширования вроде CRC32, MD5 или SHA-1, то даже самое короткое значение хеша, получаемое при использовании CRC32, будет больше 7 символов. Как его сократить? В качестве одного из решений можно взять первые 7 символов хеша, но это увеличивает вероятность конфликтов. Чтобы значения хеша не дублировались, мы можем рекурсивно добавлять к ним заданную строку, пока они не станут уникальными. Но это влечет удлинение строки.

Поэтому подход, основанный на хешировании, нам не подходит.

Другим распространенным методом сокращения URL адресов является преобразование значения в другую систему счисления. Поскольку для строкового представления ID используется набор из 62 символов, то выберем алгоритм base62. Чтобы понять, как происходит это преобразование, переведем десятичное число 11 157 в систему счисления с основанием 62. Как понятно из названия, base62 — это способ кодирования с использованием 62 символов. Они соотносятся как 0 — 0,..., 9 — 9, 10 — a, 11 — b,..., 35 — z, 36 — A, …, 61 — Z, где символ «a» соответствует 10, а «Z» соответствует 61. 11 157 = 2 x 62^2 + 55 x 62^1 + 59 x 62^0 = [2, 55, 59] → [2, T, X] в представлении base62.

Сокращенный URL адрес имеет переменную длину, которая увеличивается вместе с ID. Этот подход использует генератор уникальных ID, исключающий генерацию дубликатов. Таким образом, длина ID будет составлять не более 7 символов.

Перед реализацией функции получения уникального ID необходимо внести корректировки в алгоритм Snowflake. Для представления ID он использует 64 бита. Исходя из наших требований к системе, такая длина будет избыточной. В нашем случае для представления максимального ID достаточно 42 бит. Это дает возможность хранить примерно 4,39 триллиона ID, что гораздо больше 3,5 триллионов.

Распределение бит будет следующим:

  1. Временная метка — 28 бит. Предполагаемая продолжительность работы системы — порядка 8 лет, генерация ID идет со скоростью 1160 в секунду. Поэтому нам не нужна миллисекундная точность, достаточно секундной. 2^28 секунд — это примерно 8.5 лет.

  2. Идентификатор дата центра. Достаточно двух бит (4 значения). Для ID сервера достаточно одного бита (2 сервера на дата центр).

  3. Порядковый номер последовательности. Под это поле выделим оставшиеся 11 бит, что позволяет хранить 2048 значений, и превосходит оговоренные требования к системе.

Создадим проект на языке Go. Дерево каталогов будет расширяться по мере добавления функционала. Сейчас оно выглядит следующим образом:

shortener
├── client
│   └── idclient      Клиент для проверки работы микросервиса генератора
├── cmd               Директория с кодом серверов для запуска микросервисов
│   └── idgenservice  Кодом микросервиса генератора ID
├── config            Код для задания конфигов  
├── internal          Внутренности проекта, реализация микросервисов
│   ├── base62        Код преобразования ID в base62 строку и обратно
│   └── idgenerator   Бизнес-логика генератора ID и тесты к нему
│       ├── client    Клиентская библиотека для работы с генератором
│       └── handler   Обработчик gRPC
└── proto             Хранит proto файлы и сгенерированный с их помощью код
└── utils             Вспомогательные функции

В директории client будем хранить код программ клиентов для микросервисов. Клиент показывает основной сценарий работы с сервисом, позволяет отправлять ему запросы и получать ответы.

В папке cmd будем хранить main функции, отвечающие за работу серверной части микросервисов. В директории internal для каждого сервиса будет отдельная папка с функционалом и тестами.

Так как для взаимодействия микросервисов между собой используем gRPC, то описания сервисов и сгенерированный код будем хранить в proto.

Папка utils хранит вспомогательный функционал, который могут разделять все сервисы системы.

Реализация алгоритма Snowflake

Приступим к написанию генератора уникальных ID. Будем следовать принятому в Go способу кодирования, когда метод возвращает некоторое значение в случае успешного завершения и ошибку в случае неудачи. Поэтому для начала введем ошибки, которые могут возникнуть при работе генератора.

var (
	ErrCounterExhausted = errors.New("idgen: counter space exhausted")
	ErrDataCenterID     = errors.New("datacenterID must be between 0 and 3")
	ErrMachineID        = errors.New("machineID must be between 0 and 1")
	ErrClockBackward    = errors.New("clock moved backwards, refusing to generate ID")
    ErrInvalidBatchSize = errors.New("Invalid batch size specified")
)

Так как алгоритм Snowflake кодирует данные в виде битовых полей разной длины, мы тоже введем такие поля с указанием размеров, максимальных значений каждой группы битов, а также смещений. ID будем хранить, используя тип данных int64, в котором значащими будут только младшие 42 бита.


const (
	//размеры полей в битах
	timestampBits  = 28
	datacenterBits = 2
	machineBits    = 1
	sequenceBits   = 11
	//максимальные значения полей
	MaxDatacenterID        = (1 << datacenterBits) - 1 // 3
	MaxMachineID           = (1 << machineBits) - 1    // 1
	MaxSequence            = (1 << sequenceBits) - 1   // 2047
	MaxTimestampBits int64 = (1 << timestampBits) - 1  // 2^28-1
	//смещения полей
	machineShift    = sequenceBits
	datacenterShift = sequenceBits + machineBits
	timestampShift  = sequenceBits + machineBits + datacenterBits
	//временная метка начал работы с сервисом
	ShortenerEpoch int64 = 1777939200
)

Также будем стараться соответствовать SOLID. В частности, принципу инверсии зависимостей (буква D), а также разделения интерфейсов (буква I). Первый говорит о том, что зависимости должны строиться на абстракциях (интерфейсах), а не на конкретных реализациях. Согласно второму, лучше создать несколько маленьких специализированных интерфейсов, чем один большой на все случаи жизни. Микросервисы, для минимизации ошибок, должны покрываться тестами, разработка которых значительно упростится при использовании интерфейсов.

// BatchGenerator — интерфейс генератора. Позволяет подменять реализацию в тестах.
type BatchGenerator interface {
	NextBatch(n int) (IDBatch, error)
}

// Интерфейс определения текущего времени
type DBTime interface {
	Now() time.Time
}

Введем необходимые типы данных.

type IDBatch []int64 

// Generator для потокобезопасной реализация Snowflake ID.
type Generator struct {
	mu            sync.Mutex
	epoch         int64 
	datacenterID  int64
	machineID     int64
	sequence      int64
	lastTimestamp int64 // для диагностирования перевода часов сервера назад
	timeEngine    pool.DBTime
}

// Конфиг для микросервиса
type Config struct {
	DatacenterID int64 // номер дата центра
	MachineID    int64 // номер сервера в этом центре
}

Generator хранит данные для создания уникальных ID. К нему могут обращаться для записи несколько разных горутин одновременно. Поэтому для предотвращения гонки используем мьютекс.

Теперь можно написать функцию создания генератора. Вначале валидируем значения входных параметров и возвращаем ошибку в случае несоответствия значений.

func NewIDGenerator(cfg *Config, timeEngine pool.DBTime) (*Generator, error) {
	if cfg.DatacenterID < 0 || cfg.DatacenterID > MaxDatacenterID {
		return nil, ErrDataCenterID
	}
	if cfg.MachineID < 0 || cfg.MachineID > MaxMachineID {
		return nil, ErrMachineID
	}

	return &Generator{
		epoch:         ShortenerEpoch,
		datacenterID:  cfg.DatacenterID,
		machineID:     cfg.MachineID,
		lastTimestamp: -1,
		timeEngine:    timeEngine,
	}, nil
}

Перейдем к генерации ID. Запрашивать их каждый раз по одному слишком расточительно для производительности системы. Поэтому напишем метод, возвращающий срез ID заданной длины n (который, для краткости, назовем батчем) при каждом вызове.

func (g *Generator) NextBatch(n int) (IDBatch, error) {
	if n <= 0 {
		return nil, ErrInvalidBatchSize
	}

	// 1. Выделяем память до блокировки мьютекса,
	// чтобы не держать Lock во время выделения памяти рантаймом Go
	batch := make(IDBatch, n)

	g.mu.Lock()
	defer g.mu.Unlock()

	now := g.timeEngine.Now().Unix()
    // если часы перевели назад, то возвращаем ошибку
	if now < g.lastTimestamp {
		return nil, ErrClockBackward
	}

	// 2. Если мы находимся в рамках той же секунды, что и lastTimestamp, то
	//  проверяем, влезает ли весь батч в лимит
	if now == g.lastTimestamp {
		// Проверяем, не превысит ли g.sequence + n максимальный лимит 2047
		if g.sequence+int64(n) > MaxSequence {
			// Если батч не влезает целиком, ждем следующую секунду
			now = g.waitNextTime(now, g.lastTimestamp)
			g.sequence = 0
		}
	} else {
		g.sequence = 0
	}

	g.lastTimestamp = now
	timeStamp := now - g.epoch
	if timeStamp > MaxTimestampBits {
		return nil, ErrCounterExhausted
	}

	// Вычисляем базовую часть ID (время, датацентр, машина) один раз
	baseID := (timeStamp << timestampShift) |
		(g.datacenterID << datacenterShift) |
		(g.machineID << machineShift)

	// 3. Заполняем батч.
	for i := range n {
		batch[i] = baseID | g.sequence
		g.sequence++
	}

	return batch, nil
}

func (g *Generator) waitNextTime(now int64, last int64) int64 {

	for now <= last {
		exactNow := g.timeEngine.Now()
		// Считаем, сколько миллисекунд осталось до конца текущей секунды
		msPassed := exactNow.Nanosecond() / int(time.Millisecond)
		msToWait := 1000 - msPassed
		time.Sleep(time.Duration(msToWait+1) * time.Millisecond)

		// Обновляем Unix-время для проверки условия цикла
		now = g.timeEngine.Now().Unix()
	}

	return now
}

Теперь оценим скорость генерации батчей. Так как в методе NextBatch присутствует ожидание следующей секунды, то скорость генерации ID в основном будет ограничена знрачением аргумента метода g.waitNextTime.

Измерим скорость вычисления батча при помощи теста. Основная идея — выставить время так, чтобы генатор не ждал целую секунду и создать требуемые 1160 ID. Для этого создадим мок для управления временем в тесте.

// mockTimeEngine позволяет вручную управлять временем в тестах
type mockTimeEngine struct {
	fixedTime time.Time
}

func (m *mockTimeEngine) Now() time.Time {
	return m.fixedTime
}

И тест замера скорости генерации батча из 1160 элементов.

func BenchmarkNextBatch_1160(b *testing.B) {
	timeMock := &mockTimeEngine{fixedTime: time.Unix(ShortenerEpoch+10, 0)}
	cfg := &Config{DatacenterID: 1, MachineID: 1}
	gen, _ := NewIDGenerator(cfg, timeMock)

	// Сбрасываем таймер перед запуском цикла, чтобы не учитывать время инициализации
	b.ResetTimer()

	for i := 0; i < b.N; i++ {
		// Искусственно сдвигаем время вперед в моке, чтобы генератор
		// не ждал в цикле наступления новой секунды при частых вызовах
		timeMock.fixedTime = timeMock.fixedTime.Add(time.Second)

		_, _ = gen.NextBatch(1160)
	}
}

В результате прогона теста у меня получился следующий результат.

BenchmarkNextBatch_1160-6 113 046 11 404 ns/op 9472 B/op 1 allocs/op PASS

То есть время генерации батча из 1160 элементов составляет 11 404 наносекунды что подтверждает сделанный ранее вывод о том, что скорость функции генерации ограничена временем ожидания до наступления следующей секунды. Для батча из 2048 элементов время составило 21 478 наносекунд.

Такие результаты меня устраивают. Если в будущем скорость генерации ID вдруг станет узким местом, то метод NextBatch можно будет ускорить «перематывая» время вперед, но пока в этом нет необходимости.

Генерация идентификаторов для клиента

Теперь рассмотрим вопросы организации обмена батчами между микросервисом генерации и клиентами.

Первый вопрос — как минимизировать время ожидания клиентов, запрашивающих батч? В идеале это время должно стремиться к нулю. Второй — как минимизировать объем памяти, который тратит сервис для хранения батчей? Мы хотим чтобы в распоряжении микросервиса был определенный запас батчей, размер которого не превосходит заданную величину, чтобы не расходовать память сервера впустую, если нагрузка на сервис мала.

Для решения этих вопросов можно использовать классическую схему производитель — потребитель с ограниченным буфером.

  1. Микросервис генерации ID запускает горутину, которая генерирует батчи и заполняет буфер, размер которого задается в конфиге, например 100 элементов. Сервис, работая параллельно извлекает батчи из буфера и отдает их клиентам по запросу.

  2. Если записывающая горутина заполняет весь буфер, то она останавливается и ожидает пока не освободится место.

  3. Если буфер пуст, то все клиенты блокируются, ожидая пока микросервис не запишет в него хотя бы один батч.

Используем буферизованный канал типа chan []int64, где ёмкость канала задается сервисом. Каналы в Go спроектированы как безопасные для конкурентного доступа многих читателей и писателей одновременно. Рантайм сам синхронизирует доступ к внутренней очереди канала. Поэтому:

  • Если два клиента одновременно вызовут чтение канала рантайм гарантирует, что каждый получит свой уникальный батч — дублирования не будет.

  • Продюсер безопасно записывает в канал батчи, не заботясь о том, сколько сейчас читателей ждёт.

Для того чтобы можно было легко тестировать взаимодействие производитель — потребитель используем следующие сущности:

  1. Generator. Создает один батч вызовом NextBatch(n int).

  2. Buffer. Хранить готовые батчи. Блокирует производителя, если буфер полон. Заставляет ждать потребителя, если буфер пуст.

  3. Producer организует взаимодействие. Вызывает Generator, помещает батч в Buffer. Единственный, кто знает про две другие сущности.

Написание кода Buffer начнем с интерфейса IDBuffer.

// IDBuffer описывает методы, которые нужны обработчику
type IDBuffer interface {
	Push(ctx context.Context, batch IDBatch) error
	TakeBatch(ctx context.Context) (IDBatch, error)
}

type Buffer struct {
	ch chan IDBatch
}

// Создаем буфер заданной емкости 
func NewBuffer(batches int) IDBuffer {
	if batches < 1 {
		batches = 1
	}

	return &Buffer{ch: make(chan IDBatch, batches)}
}

// Push помещает готовый батч в буфер. Блокируется, если все слоты заняты,
// пока потребитель не освободит место (TakeBatch) либо не отменится контекст ctx.
func (b *Buffer) Push(ctx context.Context, batch IDBatch) error {
	select {
	case b.ch <- batch:
		return nil
	case <-ctx.Done():
		return ctx.Err()
	}
}

// TakeBatch извлекает батч из буфера. Вызов блокируется, если буфер пуст,
// до тех пор пока producer не добавит батч (Push) либо не отменится ctx.
func (b *Buffer) TakeBatch(ctx context.Context) (IDBatch, error) {
	select {
	case batch := <-b.ch:
		return batch, nil
	case <-ctx.Done():
		return nil, ctx.Err()
	}
}

Обратите внимание на то, что нам нужно иметь возможность прерывать работу методов Buffer из вне. Поэтому в методы Push и TakeBatch передается контекст.

Осталось написать функционал производителя (Producer), который координирует работу генератора и буфера.

type Producer struct {
	gen BatchGenerator // генратор
	buf       IDBuffer // буфер
  	batchSize int      // размер батча 
}

// Создаёт Producer, который будет генерировать батчи размером
// batchSize через gen и складывать их в buf.
func NewProducer(gen BatchGenerator, buf IDBuffer, batchSize int) *Producer {
	p := &Producer{
		gen:       gen,
		buf:       buf,
		batchSize: batchSize,
	}

	return p
}

// Run создает цикл работы производителя: сгенерировать батч -> положить в буфер.
// Предназначен для запуска в отдельной горутине: go p.Run(ctx).
func (p *Producer) Run(ctx context.Context) (err error) {

	defer func() {
		if r := recover(); r != nil {
			log.Printf("[CRITICAL] Producer run fall with panic: %v.", r)
			err = fmt.Errorf("producer panic: %v", r)
		}
	}()

	for {
		// Проверяем, не завершен ли контекст перед генерацией нового батча
		if err = ctx.Err(); err != nil {
			log.Printf("Shutting down Run with error: %v", err)
			return err
		}

		batch, batchErr := p.gen.NextBatch(p.batchSize)
		if batchErr != nil {
			return batchErr

		}
		// Push сам блокируется, если буфер полон, и сам же реагирует на
        // отмену ctx. Producer.Run просто транслирует эту отмену в выход
        // из цикла.
		if pushErr := p.buf.Push(ctx, batch); pushErr != nil {
			// ctx отменён/просрочен, пока Push ждал место в буфере —
			// корректно завершаем горутину, ничего не "теряя" молча.
			return fmt.Errorf("buffer push failed: %w", pushErr)
		}
	}
}

Нужно обратить внимание на несколько важных моментов:

  1. Так Run стартует из горутины и работает на сервере, то нам нужно позаботиться о перехвате паники, которая может возникнуть в одном из вызываемых методов. Для этого используется блок, содержащий вызов recover.

  2. Все ошибки, возникающие во время работы сервера, необходимо логировать.

  3. Сервер должен иметь возможность корректно завершить работу Run при помощи передаваемого в неё контекста. Это необходимо для реализации так называемого «корректного завершения работы» (graceful shutdown).

  4. Если Generator.NextBatch вернул ошибку, Producer не передаёт частично заполненный или пустой срез в Buffer.Push, а возвращает ошибку.

Явная ошибка — правильный выбор. Пусть Generator либо гарантированно отдаёт полный батч, либо явно сигнализирует о невозможности это сделать — а как реагировать на это решает Producer, а не Generator.

Разработка микросервиса генерации ID

Теперь приступим к созданию сервиса, соединяющего воедино весь функционал. Как говорилось в предыдущей статье, микросервисы будут взаимодействовать между собой по протоколу gRPC. Поэтому начнем с написания proto файла, хранящего контракт API, описывающий структуры передаваемых данных и доступные методы. Этот файл использует язык Protocol Buffers (Protobuf). На основе файла будет создан код для клиент — серверного взаимодействия.

Хранить proto файлы будем в директории proto, а сгенерированный код во вложенных папках (в данном случае idservice).

syntax = "proto3";

package generator;

option go_package = "./proto/idservice";

// IDService раздаёт клиентам батчи уникальных идентификаторов
// из предварительно заполненного буфера.
service IDService {
  rpc GetIDBatch (GetBatchRequest) returns (GetBatchResponse);
}

message GetBatchRequest {}

message GetBatchResponse {
  repeated int64 ids = 1;
}

Написание proto файла начинают с указания версии языка (обычно syntax = “proto3”). Далее идет пакет — пространство имен для предотвращения конфликтов между разными сервисами (package generator;) и опция go_package = “./proto/idservice”, указывающая компилятору protoc, куда именно сохранить созданный код.

Далее идет описание сервисов. Сервис — это набор RPC — методов, которые сервер может выполнять, с указанием типов входящих и исходящих сообщений. Поддерживаются четыре типа вызовов:

  1. Унарные (обычный запрос — ответ).

  2. Потоковые со стороны сервера (сервер отправляет поток данных).

  3. Потоковые со стороны клиента (клиент отправляет поток данных).

  4. Двунаправленный стриминг.

В нашем случае у сервера всего один метод GetIDBatch, принимающий входной параметр — сообщение типа GetBatchRequest и возвращающий GetBatchResponse. Сообщение представляет структуру данных, содержащую поля с указанием типов (например, int32, string, bool) и их уникальных порядковых номеров (тегов). У метода обязательно должен быть аргумент (в данном случае GetBatchRequest), даже если он пустой. В дальнейшем в этот параметр можно добавить поля с информацией для сервера.

Сообщение GetBatchResponse содержит одно поле с именем ids. Ключевое слово repeated указывает на то, что поле является массивом. Оно может содержать любое количество элементов. При компиляции.proto файла для Go repeated превращается в слайс []int64. Единица после знака равно задает тег поля. Менять его после запуска сервиса в продакшн нельзя.

После компиляции появляются два файла с кодом для сервиса. В файле idservice.pb.go размещен код сообщений. Среди прочего в нем содержатся описания структур GetBatchRequest и GetBatchResponse. Поле repeated int64 ids превращено в слайс []int64. Для каждого поля созданы методы получения значения, например func (x *GetBatchResponse) GetIds() []int64.

Файл idservice_grpc.pb.go содержит всё необходимое для создания gRPC — клиента и сервера. В том числе, код клиента (IDServiceClient), описывающий методы, которые он может вызывать удаленно. В нашем случае там будет объявлен интерфейс клиента,

type IDServiceClient interface {
	GetIDBatch(ctx context.Context, in *GetBatchRequest, opts ...grpc.CallOption) (*GetBatchResponse, error)
}

а также структура iDServiceClient, которая берет на себя работу по отправке запроса на сервер. С её помощью создание клиента происходит в одну строчку.

func NewIDServiceClient(cc grpc.ClientConnInterface) IDServiceClient {
	return &iDServiceClient{cc}
}

Также в этом файле объявлен интерфейс сервера (IDServiceServer).

type IDServiceServer interface {
	GetIDBatch(context.Context, *GetBatchRequest) (*GetBatchResponse, error)
	mustEmbedUnimplementedIDServiceServer()
}

Его необходимо написать в своем коде. Также здесь реализован регистратор сервера в виде метода RegisterIDServiceServer(s grpc.ServiceRegistrar, srv IDServiceServer). Она связывает нашу бизнес — логику с запущенным gRPC — сервером.

Кроме того, в файле есть структура UnimplementedIDServiceServer, которую встраивают в сервер для обеспечения обратной совместимости, чтобы код компилировался, даже если вы добавили новый метод в proto файл, но еще не написали для него код.

Напишем серверный обработчик gRPC запроса от клиента на получение батча.

type GRPCHandler struct {
	// Встраиваем обязательную заглушку для обратной совместимости
	pb.UnimplementedIDServiceServer

	// Внедряем буфер как зависимость
	buffer idgenerator.IDBuffer
}

// Создание gRPC обработчика
func NewHandler(buf idgenerator.IDBuffer) *GRPCHandler {
	return &GRPCHandler{
		buffer: buf,
	}
}

// Реализация метода GetIDBatch сервера
func (h *GRPCHandler) GetIDBatch(ctx context.Context, req *pb.GetBatchRequest) (*pb.GetBatchResponse, error){

	batch, err := h.buffer.TakeBatch(ctx) // забрать готовый батч из общего буфера
	if err != nil {
		return nil, err
	}
	return &pb.GetBatchResponse{Ids: batch}, nil
}

Обработчик запроса принимает буфер в качестве зависимости. Это сделано для удобства тестирования. Логика метода GetIDBatch оказалась очень простой. Обработчик забирает один батч из буфера, который наполняет Producer.Run.

Создаем gRPC сервер

Теперь можно объединить описанные механики для написания сервера. Мы пишем тестируемый код, поэтому вынесем логику в отдельную структуру App вместо того, чтобы разместить её в функции main.

type App struct {
	grpcServer *grpc.Server          // Сервер
	producer   *idgenerator.Producer // Продюсер, реализующий метод Run(ctx)
	wg         sync.WaitGroup
	port       string
}

// Создаем App с использованием конфига, буфера и генератора
func NewApp(cfg config.Config, buf idgenerator.IDBuffer, 
            gen idgenerator.BatchGenerator, batchsize int) *App {
	// 1. Создаем продюсера
	producer := idgenerator.NewProducer(gen, buf, batchsize)

	// 2. Задаем опции логирования и восстановления при сбое
	loggerOpts := []logging.Option{
		logging.WithLogOnEvents(logging.StartCall, logging.FinishCall),
	}
  	grpcLogger := logging.LoggerFunc(func(ctx context.Context, lvl logging.Level, msg string, fields ...any) {
		log.Printf("[%s] %s %v", utils.CreateDebugLevelString(lvl), msg, fields)
	})

	recoveryOpts := []recovery.Option{
		recovery.WithRecoveryHandler(func(p any) (err error) {
			stackTrace := debug.Stack()
			log.Printf("Captured critical error (panic):: %v\n Call stack:\n%s", p, string(stackTrace))
			return status.Errorf(codes.Internal, "Internal server error")
		}),
	}

	// 3. Создание gRPC сервера
	gRPCServer := grpc.NewServer(
		grpc.ChainUnaryInterceptor(
			recovery.UnaryServerInterceptor(recoveryOpts...),
			utils.EnforceDeadlineInterceptor(),
			logging.UnaryServerInterceptor(grpcLogger, loggerOpts...),
		),
	)

	// 4. Регистрация обработчика запроса от клиента
	grpcHandler := NewHandler(buf)
	pb.RegisterIDServiceServer(gRPCServer, grpcHandler)

	return &App{
		grpcServer: gRPCServer,
		producer:   producer,
		port:       cfg.Port,
	}
}

При создании сервера нужно учитывать следующее:

  1. Серверу может понадобиться логирование запросов. Как минимум на этапе тестирования и отладки. Для этого задается метод logging.LoggerFunc. Он срабатывает в ответ на определенные события, задаваемые в loggerOpts.

  2. Необходимо уметь реагировать на паники. Для этого вызовом recovery.WithRecoveryHandler задается обработчик, вызываемый в случае возникновения критических ошибок. Если внутри любого обработчика gRPC запроса произойдет panic (например, обращение к nil‑указателю), сервер её перехватит. Функция — перехватчик выполняет следующие действия. Вызывает debug.Stack(), чтобы получить детальный стек вызовов (traceback) позволяющий понять, где именно произошла ошибка. Записывает ошибку и стек в лог для разработчиков. Возвращает клиенту понятную gRPC ошибку («Internal server error») вместо обрыва соединения.

  3. Очень важно защитить сервер от запросов, у которых не задан крайний срок выполнения (deadline). Для отслеживания этих запросов используем обработчик (интерцептор в терминах gRPC) EnforceDeadlineInterceptor.

func EnforceDeadlineInterceptor() grpc.UnaryServerInterceptor {
	return func(
		ctx context.Context, req any, info *grpc.UnaryServerInfo,
		handler grpc.UnaryHandler, ) (any, error) {

		// Проверяем, задан ли уже дедлайн клиентом во входящем контексте
		_, hasDeadline := ctx.Deadline()

		if !hasDeadline {
			// Отклоняем запрос клиента, влзвращая ошибку
			return nil, status.Error(codes.InvalidArgument, "gRPC deadline must be specified by the client")
		}

		// Передаем контекст дальше в обработчик сервиса
		return handler(ctx, req)
	}
}

Вызов grpc.NewServer создает и настраивает gRPC — сервер, подключая к нему цепочку из трех промежуточных обработчиков для одиночных запросов. Благодаря методу grpc.ChainUnaryInterceptor, каждый входящий запрос будет проходить сквозь эти обработчики поочередно перед тем как попасть в наш обработчик запроса генерации батча.

Порядок указания обработчиков в ChainUnaryInterceptor имеет важное значение. Если в интерцепторе дедлайнов, логирования или самом обработчике клиентского запроса случится паника, она будет перехвачена, стек вызовов запишется в лог и клиенту вернется ошибка в качестве ответа. Если таймаут запроса не задан, обработчик прервет выполнение запроса и вернет ошибку. Если ошибок нет, то обработчик логирования запишет информацию о начале запроса. Затем управление передается в наш обработчик GetIDBatch. Как только он завершит работу (успешно или с ошибкой), обработчик логирования запишет FinishCall с результатом выполнения.

Метод Run запускает создание батчей продюсером (метод producer.Run), а также gRPC сервер.

// Run запускает все асинхронные компоненты приложения
func (a *App) Run(ctx context.Context, lis net.Listener) error {
	// Создаем буферизированный канал на 2 элемента
	errChan := make(chan error, 2)
	// Запускаем продюсера
	a.wg.Add(1)
	go func() {
		defer a.wg.Done()
		errChan <- a.producer.Run(ctx)
	}()

	// Запускаем gRPC сервер
	go func() {
		log.Printf("gRPC server listening on %s", lis.Addr().String())
		if err := a.grpcServer.Serve(lis); err != nil && 
           !errors.Is(err, grpc.ErrServerStopped) {
			errChan <- err
		}
	}()

	// Ожидаем сигнала отмены контекста (Ctrl+C / SIGTERM) или ошибки сервера
	select {
	case <-ctx.Done():
		log.Printf("selected case case <-ctx.Done():, stopping App")
		return a.Stop()
	case err := <-errChan:
		log.Printf("selected case err := <-errChan:")
		return err
	}
}

Также этот метод позволяет серверу корректно остановить работу, если продюсер завершился с ошибкой или получен сигнал отмены контекста.

Корректное завершение работы происходит в методе Stop.

// Stop выполняет graceful shutdown сервера
func (a *App) Stop() error {
	log.Println("Inside App Stop() shutting down gRPC server...")
    // контекст ожидания завершения сервера
	shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer cancel()

	grpcDone := make(chan struct{})
	go func() {
		a.grpcServer.GracefulStop() // ожидаем завершение всех соединений
		close(grpcDone)
	}()

	producerDone := make(chan struct{})
	go func() {
		a.wg.Wait()
		close(producerDone)
	}()

	select {
	case <-shutdownCtx.Done():
		log.Println("Shutdown timeout exceeded! Forcing stop...")
		a.grpcServer.Stop()
		return errors.New("shutdown timeout exceeded")
	case <-grpcDone:
		select {
		case <-producerDone:
			log.Println("All systems cleanly exited")
		case <-shutdownCtx.Done():
			log.Println("Warning: Producer shutdown timed out")
			return errors.New("shutdown timeout exceeded: producer hung")
		}
	}

	log.Println("gRPC server stopped")
	return nil
}

Для этого создается контекст shutdownCtx с таймаутом 5 секунд. После этого в отдельных горутинах вызывается блокирующий метод grpcServer.GracefulStop(), а также ожидание завершения продюсера. Сервер прекращает принимать новые входящие соединения и ждет, пока все активные запросы завершат свою работу.

Ожидание выполняется в отдельных горутинах, так как вызовы блокируют текущий поток выполнения. Сервер не должен ждать бесконечно. Поэтому ожидание завершения происходит при помощи чтения из каналов grpcDone, producerDone, а также shutdownCtx.Done() контекста завершения.

Используя экземпляр структуры App, можно написать функцию main нашего микросервиса.

const (
	buffersize = 3 
	batchsize  = 3
)

func main() {
	// 1. Загружаем конфигурацию
	cfg := config.Load()

	// 2. Настройка gRPC-слоя
	// Обязательно закрываем файл при завершении работы всего приложения
	logFile := utils.CreateLogFile("grpc_server" + cfg.Port + ".log")
	defer logFile.Close()
	log.SetOutput(logFile)

	timeEngine := pool.UnixTimeReal{}
	// 3. Создаем генератор, буфер и продюсер
	gen, err := idgenerator.NewIDGenerator(&idgenerator.Config{
		DatacenterID: cfg.DatacenterID,
		MachineID:    cfg.MachineID,}, timeEngine)
	if err != nil {
		log.Printf("failed to create generator: %v", err)
		return
	}

	buf := idgenerator.NewBuffer(buffersize)
	app := handler.NewApp(cfg, buf, gen, batchsize)

	// 4. Слушаем порт
	lis, err := net.Listen("tcp", ":"+cfg.Port)
	if err != nil {
		log.Printf("failed to listen port %s: %v", cfg.Port, err)
		return
	}

	// 5. Отслеживаем системные сигналы ОС
	ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
	defer stop()

	if err := app.Run(ctx, lis); err != nil {
		log.Printf("App execution failed: %v", err)
	}
	log.Println("gRPC server stopped")
}

Создание клиента для получения уникальных ID от сервера

Для полноты картины напишем простой клиент, циклический посылающий запросы на получение ID к gRPC серверу и выводящий результаты в лог.

В основе клиента будет лежать тот же механизм взаимодействия производитель — потребитель с ограниченным буфером. Теперь Generator — не создает уникальные ID по алгоритму Snowflake, а делает gRPC вызов GetBatch к микросервису. Producer в фоне поддерживает буфер заполненным на заданное количество (lookahead) батчей. Поэтому в типичном сценарии использования вызов для получения очередного ID (NextID) почти всегда отдаёт его из уже полученного от сервера среза, не дожидаясь обращения к сервису. Producer.Run запрашивает по сети следующий батч до того, как текущий слот буфера освобождён — то есть запрос следующего батча стартует, пока клиент ещё использует текущий, а не после того, как он его исчерпал.

Генератор теперь выглядит следующим образом.

// BatchFetcher — интерфейс, нужный пулу ID.
// Тесты подставляют фейковую реализацию, не поднимая настоящее gRPC-соединение.
type BatchFetcher interface {
	GetIDBatch(ctx context.Context, in *pb.GetBatchRequest, opts ...grpc.CallOption) (*pb.GetBatchResponse, error)
}

// rpcGenerator адаптирует BatchFetcher под интерфейс producer.Generator,
// чтобы переиспользовать существующий Producer для предвыборки по сети.
type rpcGenerator struct {
	client  BatchFetcher
	ctx     context.Context
	timeout time.Duration
}

// NextBatch запрашивает у сервиса очередной батч по gRPC.
// Параметр n намеренно игнорируется: реальный размер батча определяет
// сервис, клиент на него не влияет.
func (g *rpcGenerator) NextBatch(_ int) (idgenerator.IDBatch, error) {
	ctx, cancel := context.WithTimeout(g.ctx, g.timeout)
	defer cancel()

	resp, err := g.client.GetIDBatch(ctx, &pb.GetBatchRequest{})
	if err != nil {
		return nil, err
	}
	return resp.Ids, nil
}

Теперь реализуем пул для хранения ID на клиенте с предварительной выборкой.

// Клиентский локальный пул ID с предвыборкой.
type Pool struct {
	buf idgenerator.IDBuffer
	mu      sync.Mutex
	current idgenerator.IDBatch
	idx     int
	cancel context.CancelFunc
}

// Создаёт пул и сразу запускает фоновую предвыборку.
// lookaheadBatches задает глубину локальной очереди в батчах. 
// lookaheadBatches = 1 уже даёт скрытие сетевой задержки за счёт предвыборки
// rpcTimeout задает таймаут на один вызов GetBatch к сервису.
func NewPool(ctx context.Context, client BatchFetcher, lookaheadBatches int, rpcTimeout time.Duration, wg *sync.WaitGroup) *Pool {
	if lookaheadBatches < 1 {
		lookaheadBatches = 1
	}

	buf := idgenerator.NewBuffer(lookaheadBatches)

	cancelCtx, cancel := context.WithCancel(ctx)
	gen := &rpcGenerator{client: client, ctx: cancelCtx, timeout: rpcTimeout}
	prod := idgenerator.NewProducer(gen, buf, 0) // batchSize не используется rpcGenerator'ом — реальный размер решает сервис

	wg.Add(1)
	go func() {
		defer wg.Done() // Сообщаем серверу, что горутина продюсера полностью завершилась
		prod.Run(cancelCtx)
	}()

	return &Pool{buf: buf, cancel: cancel}
}

// Останавливает фоновую предвыборку. Уже полученные (в том числе
// предвыбранные, но ещё не розданные) id по-прежнему можно забрать через
// NextID, пока локальный запас не иссякнет. Новых батчей запрошено не будет.
func (p *Pool) Close() {
	p.cancel()
}

// Отдаёт очередной уникальный id.
// Если в текущем локальном батче есть запас, то возвращает id мгновенно, без
// сетевого вызова. Если локальный батч исчерпан, блокируется на
// buf.TakeBatch: в типичном случае там уже лежит предвыбранный батч
// Безопасен для конкурентного вызова из нескольких горутин: mutex
// удерживается на всё время получени батча намеренно.
// Это гарантирует, что при одновременном исчерпании батч запрашивается
// только один раз, а не по разу на каждую заблокированную горутину.
func (p *Pool) NextID(ctx context.Context) (int64, error) {
	p.mu.Lock()
	defer p.mu.Unlock()

	for p.idx >= len(p.current) {
		batch, err := p.buf.TakeBatch(ctx)
		if err != nil {
			return 0, err
		}
		p.current = batch
		p.idx = 0
	}

	id := p.current[p.idx]
	p.idx++
	return id, nil
}

Осталось написать код клиента. Для хранения состояния клиента создадим структуру IDConsumerApp.

// Конфигурация вынесена в отдельную структуру для гибкости
type Config struct {
	Addr             string
	Count            int
	LookaheadBatches int
	RPCTimeout       time.Duration
	DialTimeout      time.Duration
}

// IDConsumerApp объединяет зависимости и состояние клиентского приложения
type IDConsumerApp struct {
	cfg  Config
	wg   sync.WaitGroup
	conn *grpc.ClientConn
	pool *client.Pool // Предполагается, что структура клиентского пула из вашего пакета
}

func NewIDConsumerApp(cfg Config) *IDConsumerApp {
	return &IDConsumerApp{
		cfg: cfg,
	}
}

Основной метод IDConsumerApp — это Run.

// Точка входа, управляющая жизненным циклом клиента
func (a *IDConsumerApp) Run(ctx context.Context) error {
	// 1. Подключаемся к gRPC серверу
	if err := a.setupClientgRPCConnection(ctx); err != nil {
		return fmt.Errorf("unable to establish gRPC connection: %w", err)
	}
	defer func() {
		if a.conn != nil {
			a.conn.Close()
		}
	}()

	// 2. Инициализируем пул предвыборки ID
	grpcClient := pb.NewIDServiceClient(a.conn)
	a.pool = client.NewPool(ctx, grpcClient, a.cfg.LookaheadBatches, a.cfg.RPCTimeout, &a.wg)
	
	log.Printf("Connected to %s, prefetching=%d batches. ID request...", a.cfg.Addr, a.cfg.LookaheadBatches)

	// 3. Запускаем основной цикл получения ID
	a.runIDConsumerLoop(ctx)
	// и закрываем пул после завершения цикла, тем самым завершая горутину продюсера
	a.pool.Close()
	// 4. Выполняем Graceful Shutdown для фоновых горутин пула
	a.waitForShutdown()
	return nil
}

Вначале мы создаем подключение к серверу и используем его для создания клиента вызовом NewIDServiceClient, написанным для нас gRPC. Далее в цикле получаем Count уникальных ID в методе runIDConsumerLoop.

// запрашивает пачки ID из пула согласно лимиту count
func (a *IDConsumerApp) runIDConsumerLoop(ctx context.Context) {
	got := 0
	for a.cfg.Count == 0 || got < a.cfg.Count {
		id, err := a.pool.NextID(ctx)
		if err != nil {
			if ctx.Err() != nil {
				log.Printf("The loop was stopped by the context signal: %v", ctx.Err())
				break
			}
			log.Printf("NextID method error: %v", err)
			break
		}

		log.Printf("Successfully received id=%d", id)
		got++
		time.Sleep(1 * time.Second)
	}
	log.Printf("Cycle completed: %d identifiers received in total", got)
}

После этого корректно завершаем работу клиента.

func (a *IDConsumerApp) waitForShutdown() {
	shutdownCtx, cancel := context.WithTimeout(context.Background(), a.cfg.DialTimeout)
	defer cancel()

	done := make(chan struct{})
	go func() {
		a.wg.Wait()
		close(done)
	}()

	select {
	case <-done:
		log.Println("Goroutines finished")
	case <-shutdownCtx.Done():
		log.Println("Waiting timeout expired. Force finish")
	}
}

Функция main просто создает клиента и вызывает его метод Run.

func main() {
	// 1. Инициализируем логирование
	logFile := utils.CreateLogFile("grpc_client.log")
	defer logFile.Close()
	log.SetOutput(logFile)

	// 2. Настраиваем системный контекст для Ctrl+C / SIGTERM
	ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
	defer stop()

	// 3. Собираем конфигурацию
	cfg := Config{
		Addr:             "localhost:50051",
		Count:            10,
		LookaheadBatches: 2,
		RPCTimeout:       2 * time.Second,
		DialTimeout:      5 * time.Second,
	}

	// 4. Запускаем приложение
	app := NewIDConsumerApp(cfg)
	if err := app.Run(ctx); err != nil {
		log.Printf("Application critical error: %v", err)
	}
}

Выводы

Сегодня мы разработали микросервис генерации уникальных ID на основе алгоритма Snowflake. Длина ID составляет 42 бита и является минимально возможной для представления их требуемого количества. Сервис может генерировать 2048 ID в секунду. Всего может быть запущено восемь экземпляров (4 дата центра по 2 сервера в каждом). Сервисы работают независимо и не требуют обращений к базе данных. Падение одного экземпляра не останавливает генерацию ID в системе. Чтобы снизить время ожидания клиента микросервис отдаёт ему срез ID. Сервер логирует запросы клиентов, корректно обрабатывает панику и защищен от запросов без установленного крайнего срока завершения. При получении соответствующего сигнала сервер выполняет graceful shutdown с ожиданием завершения текущих запросов.

Мы рассмотрели не все вопросы, касающиеся создания gRPC сервера. Остаются, как минимум, три темы:

  1. Сбор метрик функционирования сервера.

  2. Балансировка нагрузки при помощи gRPC.

  3. Обнаружение сервисов.

Их мы рассмотрим в следующий раз. Эта статья и так получилась слишком объемной. В ней возможно присутствуют некоторые ошибки и неточности о которых прошу сообщать в комментариях. Продолжение следует...