Привет, Хабр! Продолжаю серию статей по проектированию и разработке отказоустойчивой микросервисной системы сокращения ссылок.

Если вы только присоединились, рекомендую ознакомиться с предыдущими частями:

  • Часть 1. Определили требования к системе и декомпозировали её на независимые сервисы.

  • Часть 2. Создали микросервис генерации уникальных идентификаторов.

Наш микросервис генерации ID полностью готов, но пока живет только в среде разработки. Сегодня упакуем его в контейнер Docker и развернем в локальном Kubernetes-кластере Minikube. При этом решим несколько важных проблем:

  1. Алгоритм генерации ID требует, чтобы каждый экземпляр микросервиса обладал уникальным числовым идентификатором.

  2. Для производительной и отказоустойчивой работы микросервисной архитектуры требуется поддержка одновременной работы нескольких экземпляров микросервиса и балансировка нагрузки между ними. Для этого нужно решить задачи обнаружения сервисов и мониторинга их готовности к обработке запросов. Разберем как решить эти задачи при помощи Kubernetes.

  3. Получение уникальных ID клиентским сервисом внутри Kubernetes.

Эта статья ориентирована на начинающих разработчиков, поэтому вначале рассмотрим внутреннее устройство Kubernetes (сокращенно K8s).

Кластер K8s состоит из узлов (Node или нод). Узел — это виртуальная машина (ВМ), на которой запускаются приложения, именуемые подами (Pods или поды ), представляющими из себя набор контейнеров. Различают два вида узлов: Master node (мастер узел) и Worker node.

Master node включает компоненты, необходимые для работы кластера и управления им:

  • API Server — основной интерфейс для работы с кластером, обрабатывающий запросы API.

  • etcd — хранилище данных о состоянии кластера.

  • Scheduler определяет узлы на которых будут запускаться новые Pods.

  • Controller Manager отвечает за поддержание желаемое состояние кластера, управляя различными контроллерами.

Контроллеры в Kubernetes — управляющие циклы, которые отслеживают состояние вашего кластера, затем вносят или запрашивают изменения там, где это необходимо. Каждый контроллер пытается привести текущее состояние кластера ближе к желаемому состоянию.

Worker node — рабочий узел, на котором запускаются пользовательские Pods. Они включают:

  • Kubelet агент, который отвечает за запуск Pods на узле и управление их состоянием.

  • Kube—proxy обеспечивает сетевое взаимодействие и балансировку нагрузки между Pods.

  • Пользовательский Pod — базовая единица развёртывания в K8s. Каждый под — это один или несколько контейнеров, имеющих общие сетевые ресурсы и хранилище. Это своеобразная «обертка» вокруг одного или нескольких Docker—контейнеров.

В K8s Docker образы не запускаются напрямую. Вместо этого вы оперируете подами. Поды могут постоянно создаваться, уничтожаться и получать новые IP адреса после перезапуска. Контейнеры в Pod разделяют один IP адрес.

Pods использует тома для хранения данных. Они могут использоваться всеми контейнерами, входящими в Pod.

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

В Kubernetes используется декларативный подход к управлению инфраструктурой. Вместо того чтобы задавать команды вроде "запустить контейнер, выделить ему 2 ГБ памяти и открыть порт", мы описываем в специальном файле желаемое конечное состояние системы. Мы как бы говорим: "хочу чтобы в системе работал такой сервис", а плоскость управления сама решает, как этого достичь. Этот файл конфигурации пишется в формате yaml и называется манифестом.

Каждый манифест Kubernetes состоит из четырех обязательных базовых блоков:

# 1. Версия API, которую использует ресурс
apiVersion: apps/v1

# 2. Тип создаваемого ресурса (сущности)
kind: StatefulSet

# 3. Метаданные, идентифицирующие этот конкретный объект
metadata:
  name: idgen-service
  labels:
    app: idgen-service

# 4. Спецификация — самое главное тело манифеста
spec:
  replicas: 2
  # ... (описание контейнеров, ресурсов, и сетевых портов)
  1. apiVersion указывает версию API сервера Kubernetes, с которой будет взаимодействовать этот объект. Например, для Pod или Service нужна v1. Более сложные инфраструктурные сущности, такие как Deployment или StatefulSet, требуют расширенной ветки apps/v1.

  2. kind определяет тип ресурса, который вы хотите создать. Именно эта строчка говорит Kubernetes, какой именно контроллер нужно активировать для обработки вашего файла.

  3. metadata — это "паспорт" объекта. Здесь указывается уникальное имя (name), по которому вы будете искать ресурс через консоль. Также здесь прописываются метки — ключевой механизм связывания сущностей. Например, Service находит нужные ему поды именно по совпадению меток в графе labels.

  4. spec описывает из чего состоит ваш ресурс. Для контроллеров в этой секции указывается:

    1. replicas — сколько копий нужно запустить.

    2. template — шаблон, по которому будут генерироваться новые поды. Внутри него описывается образ контейнера (image), порты (ports), переменные окружения (env), требования к памяти и процессору (resources), а также пробы работоспособности (probes). Если в image указать название образа, то будет использован образ из публичного реестра Docker.

Реестр Docker — это специальное хранилище и сервис для управления, публикации и скачивания образов контейнеров (Docker образов).

Чтобы создать Pod из примера выше, необходимо передать его манифест кластеру. Это можно сделать при помощи kubectl консольной утилиты для взаимодействия с объектами кластера.

Если под — это элемент инфраструктуры, то контроллер можно сравнить с управляющим, который следит за жизненным циклом этих элементов. Он непрерывно сравнивает реальное состояние кластера с тем состоянием, которое описано в манифесте.

В основе его работы лежит так называемый цикл сверки:

  1. Контроллер смотрит в манифест (желаемое состояние): например, "хочу 3 реплики приложения".

  2. Опрашивает кластер (текущее состояние): "в живых осталось только 2 пода, один упал".

  3. Устраняет разницу: отдает команду плоскости управления поднять новый под взамен упавшего.

В зависимости от задач в K8s применяются следующие виды контроллеров:

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

  • Deployment используют для управления обновлениями Pods и ReplicaSets. Deployment предоставляет управление обновлениями приложений. Это делает его основным инструментом для развёртывания и управления в Kubernetes. Ему не важно, в каком порядке запускаются поды и какие у них имена так как они абсолютно взаимозаменяемы.

  • StatefulSet — специализированный контроллер для приложений с состоянием (stateful). StatefulSet гарантирует, что каждому поду присваивается неизменяемый индекс, начиная с нуля. Поды запускаются по очереди, а при перезапуске полностью сохраняют свое имя.

  • DaemonSet гарантирует, что копия конкретного пода будет запущена на каждом рабочем узле кластера.

  • Job — контроллер для разовых задач. Он запускает под, даёт выполнить ему задачу, после чего завершает.

Наш генератор ID зависит от уникальных значений идентификаторов дата центра и ВМ внутри него. Deployment нам не подходит так как не гарантирует уникальности при перезапуске пода. Поэтому мы используем StatefulSet — специализированный контроллер для приложений с состоянием, который присваивает каждому поду неизменяемый индекс, начиная с нуля.

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

apiVersion: v1
kind: Service
metadata:
  name: idgen-headless
  labels:
    app: idgen
spec:
  # clusterIP: None делает сервис Headless. Обязательно для StatefulSet
  # и для правильного разрешения адресов в gRPC клиентах.
  clusterIP: None 
  selector:
    app: idgen
  ports:
    - name: grpc
      port: 50051
      targetPort: 50051
---
apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: idgen
spec:
  serviceName: "idgen-headless"
  replicas: 3 # Запустит 3 пода: idgen-0, idgen-1, idgen-2
  selector:
    matchLabels:
      app: idgen
  template:
    metadata:
      labels:
        app: idgen
    spec:
      containers:
        - name: idgen-service
          # Замените на имя вашего собранного образа
          image: practicalprogrammer123/idgen:v1
          imagePullPolicy: IfNotPresent # Для локального тестирования в Minikube/Kind
          ports:
            - containerPort: 50051
              name: grpc
          
          # Переменные окружения, которые K8s передает внутрь Go-кода
          env:
            # Назначаем порт
            - name: PORT
              value: "50051"
         

          # Проверки здоровья (Probes). Kubernetes использует gRPC-проверку,
          # чтобы знать, готов ли ваш Snowflake выдавать ID.
          # Для этого в Go должен быть подключен "google.golang.org/grpc/health"
          readinessProbe:
            grpc:
              port: 50051
              service: "readiness" # проверка готовности
            initialDelaySeconds: 2 # Подождать 2 секунды после старта перед первой проверкой
            periodSeconds: 5       # Проверять каждые 5 секунд 
          livenessProbe:
            grpc:
              port: 50051
              service: "liveness"  # проверка живучести
            initialDelaySeconds: 5
            periodSeconds: 10

Вначале создается Headless Service со специальным флагом clusterIP: None. Это встроенный в K8s механизм обнаружения сервисов, который мы обсуждали в первой части. Он позволяет подам StatefulSet иметь стабильные сетевые имена и общаться по gRPC напрямую. Headless Service позволяет реализовать балансировку gRPC запросов на стороне клиента в Kubernetes.

Я не являюсь специалистом в K8s, но насколько я понимаю, при использовании обычного сервиса K8s дает один виртуальный IP. Если gRPC-клиент обратится по нему, то установит одно постоянное TCP-соединение с этим IP. Из-за того что gRPC мультиплексирует запросы в одном HTTP/2 потоке, все последующие запросы пойдут на один и тот же под. Балансировка не будет работать.

Headless Service с clusterIP: None не создает виртуальный IP. Вместо этого он настраивает внутренний DNS-сервер Kubernetes так, что имя idgen-headless начинает возвращать список IP адресов всех живых подов этого StatefulSet. gRPC клиент на Go делает запрос к idgen-headless, получает этот массив IP адресов и самостоятельно распределяет gRPC запросы между ними, например, по алгоритму round robin.

port: 50051 — это порт, к которому будут обращаться другие микросервисы (клиенты) внутри кластера, когда захотят вызвать генератор ID.

targetPort: 50051 — это порт, который слушает микросервис внутри контейнера.

StatefulSet запускает 3 пода, используя образы docker, заданные при помощи тега image. В описании сервиса стоит обратить внимание на проверки здоровья (probes). Используются два механизма проверки здоровья в Kubernetes, которые решают совершенно разные задачи.

livenessProbe (проба живучести) проверяет, может ли контейнер принимать запросы (не завис ли он). Если проверка не пройдена, K8s перезапускает контейнер, чтобы заставить его работать. Это нужно в случае если приложение намертво зависло и не может продолжать работу.

readinessProbe (проба готовности) проверяет, готов ли контейнер принимать трафик от пользователей. Если проверка не пройдена, то убирает под из списка обслуживания балансировщика. Под не перезапускается, на него просто временно не идет трафик. Это нужно в случаях когда приложение запускается или перегружено запросами и нужно дать ему время.

Теперь посмотрим какие изменения нужны в коде для корректной работы в K8s. Первое изменение касается логирования. Нужно писать логи не в файл, как мы делали в предыдущей части, а в os.Stdout. Кроме того, при запуске микросервиса нужно получать идентификаторы дата центра и машины в нем из переменных окружения, куда их пробрасывает K8s.

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

	if nodeId, err := utils.ExtractNodeID(); err != nil {
		log.Printf("failed to extract nodeId: %v", err)
		return
	} else {
		datacenterID, machineID := utils.ExtractMachineAndDataCenterID(nodeId)
		log.Printf("starting with data center Id: %d and machine Id: %d", datacenterID, machineID)
		cfg.SetDataCenterAndMachineID(datacenterID, machineID)
	}

	port := os.Getenv("PORT")
	if len(port) != 0 {
		cfg.SetPort(port)
	}

	timeEngine := pool.UnixTimeReal{}
	// 2. Создаем генератор, буфер и продюсер
	gen, err := idgenerator.NewIDGenerator(&idgenerator.Config{
		DatacenterID: cfg.DatacenterID,
		MachineID:    cfg.MachineID,
	}, timeEngine)

Код получения идентификаторов также крайне прост.

// Считывает имя хоста пода и извлекает его порядковый номер
func ExtractNodeID() (int64, error) {
	// В Kubernetes StatefulSet hostname всегда равен имени пода (например, "idgen-0")
	hostName, err := os.Hostname()
	if err != nil {
		return 0, fmt.Errorf("failed to get hostname: %w", err)
	}

	return ExtractPodNumber(hostName)
}

// Пытается получить номер пода
func ExtractPodNumber(hostName string) (int64, error) {
	// Находим последний дефис в имени хоста
	lastDash := strings.LastIndex(hostName, "-")
	if lastDash == -1 || lastDash == len(hostName)-1 {
		// Если дефиса нет, значит код запущен локально (например, на Windows/Mac)
		// Возвращаем ошибку или дефолтный ID для локальной разработки
		return 0, fmt.Errorf("hostname %s does not contain a valid pod index", hostName)
	}

	// Вырезаем подстроку после дефиса (например, из "idgen-1" получим "1")
	indexStr := hostName[lastDash+1:]

	nodeID, err := strconv.ParseInt(indexStr, 10, 64)
	if err != nil {
		return 0, fmt.Errorf("failed to parse node ID from substring %s: %w", indexStr, err)
	}

	return nodeID, nil
}

Если программа стартует в K8s, то имя хоста будет состоять из имени пода, к которому добавлено число, например idgen-0. Функция извлекает это число и возвращает. Если извлечь не удалось, то программа запускается без Kubernetes. Поэтому берем значения по умолчанию.

После получения числа выделяем два младших бита, считая что нулевой хранит ID машины, а первый — идентификатор дата центра.

func ExtractMachineAndDataCenterID(nodeId int64) (datacenterID int64, machineID int64) {
	datacenterID = nodeId >> 1 // сдвиг на 1 бит вправо даст ID датацентра
	machineID = nodeId & 1     // остаток даст ID машины
	return datacenterID, machineID
}

Для того чтобы клиентский пул на Go автоматически обнаруживал реплики нашего микросервиса через Headless Service и балансировал запросы между ними, необходимо немного изменить метод setupClientgRPCConnection, написанный в предыдущей части.

func (a *IDConsumerApp) setupClientgRPCConnection(ctx context.Context) error {
	dialCtx, dialCancel := context.WithTimeout(ctx, a.cfg.DialTimeout)
	defer dialCancel()

	// для kubernetes настроим конфигурацию сервиса для включения round robin балансировки
	serviceConfig := `{"loadBalancingConfig": [{"round_robin":{}}]}`

	conn, err := grpc.NewClient(a.cfg.Addr, grpc.WithTransportCredentials(insecure.NewCredentials()),
		grpc.WithDefaultServiceConfig(serviceConfig))
	if err != nil {
		return err
	}

	a.conn = conn
	return nil
}

Здесь в метод NewClient добавляется параметр serviceConfig, задающий политику балансировки round robin.

Осталось добавить проверки самочувствия сервиса. Используем пакет grpc_health_v1.

Добавим сервер отслеживания состояния здоровья в метод NewApp.

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. Добавляем функции проверки здоровья
	healthServer := health.NewServer()
	healthServer.SetServingStatus("liveness", healthgrpc.HealthCheckResponse_SERVING)

	// Но приложение пока НЕ ГОТОВО принимать трафик клиентов (Readiness = NOT_SERVING)
	// так как мы, например, ещё не заполнили буфер ID или не проверили сеть
	healthServer.SetServingStatus("readiness", healthgrpc.HealthCheckResponse_NOT_SERVING)

	// 5. Регистрация обработчиков
	grpcHandler := NewHandler(buf)
	pb.RegisterIDServiceServer(gRPCServer, grpcHandler)
	healthgrpc.RegisterHealthServer(gRPCServer, healthServer)

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

При создании приложения выставляем liveness статус в HealthCheckResponse_SERVING (обслуживается), а readiness статус в HealthCheckResponse_NOT_SERVING (не обслуживается).

Далее в методе Run запускаем монитор состояния буфера в который заносятся готове батчи ID.

func (a *App) monitorBufferReadiness(ctx context.Context, lowWatermark int) {
	ticker := time.NewTicker(500 * time.Millisecond)
	defer ticker.Stop()

	currServing := true

	for {
		select {
		case <-ctx.Done():
			return
		case <-ticker.C:
			curLength := a.producer.GetNumBatches()
			if curLength < lowWatermark {
				if currServing {
					a.healthServer.SetServingStatus("readiness", healthgrpc.HealthCheckResponse_NOT_SERVING)
					log.Printf("Buffer is running low: (%d) batches", curLength)
					currServing = false
				}
			} else {
				if !currServing {
					a.healthServer.SetServingStatus("readiness", healthgrpc.HealthCheckResponse_SERVING)
					log.Printf("Buffer recovered: (%d) batches", curLength)
					currServing = true
				}
			}
		}
	}
}

В нем мониторим количество заполненных элементов буфера и если оно меньше заданной величины lowWatermark, то выставляем readiness в HealthCheckResponse_NOT_SERVING и сбрасываем флаг обслуживания новых запросов (currServing). В противном случае если флаг currServing был сброшен, то сервис готов к обработке новых запросов и можно выставить readiness в значение healthgrpc.HealthCheckResponse_SERVING.

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

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

Комментарии (0)