Hazelcast: Хороший, плохой, злой

от автора

Всем привет! Меня зовут Петрович и я работаю техлидом java команды в небольшом международном финтехе. В этой статье я хотел бы поделиться с вами опытом миграции наших сервисов с Redis на Hazelcast, а также о том какие уроки мы усвоили. Статья будет полезна тем кто рассматривает отказоустойчивые и более простые в поддержке альтернативы Redis.

Наша команда отвечает за подсистему оценки платежеспособности клиентов (в народе — скоринг). Стек — классический джентльменский набор для java команды — последняя LTS версия Java, Spring Boot (с Hibernate), база PostgreSQL, брокер сообщений RabbitMQ, авторизация через Keycloak, кэш Redis, собираем в docker контейнеры с помощью jib и деплоим в k8s через gitlab, arttifactory используем для хранения стартеров и библиотек.

В прод окружении наши немногочисленные сервисы работают в кластерах k8s в единственном экземпляре, будучи, как правило, распределенными по 3 нодам с приличным запасом по ресурсам. Сделано так специально чтобы в случае отказа одной из нод другие могли автоматически принять нагрузку с выбывшей из строя напарницы без участия сотрудников разработки или поддержки. Ну и rollout update делать удобно и быстро — у нас распределенный монолит и при деплое необходимо деплоить сразу все сервисы одновременно.

С точки зрения кластеров k8s наша подсистема полностью изолирована и не зависит от других подсистем компании. Другими словами, все необходимые для работы подсистемы ресурсы мы хостим и поддерживаем сами. Redis, Keycloak, S3 — все это работает внутри каждого из наших кластеров, обеспечивая достаточное для бесперебойного функционирования резервирование. Все, кроме Redis.

Redis — это исторически сложившееся решение (да, я знаю, именно с этой фразы начинается объяснение самых паршивых решений, но что поделать, такова реальность). И я немного слукавил сказав что это кэш. Дело в том что у нас используется Camunda для настройки роутинга процесса скоринга между сервисами. Для Camunda требуются данные, много данных. Эти данные хранятся во внутренней структуре, назовем ее моделью. Так вот, эта модель может спокойно превышать 5-10 мегабайт, так как хранит внутри себя примерно все что требуется для расчета. Просто и удобно — все собрано в одном месте в менеджере процессов и Camunda в любой момент имеет доступ к любым данным. Хранить такую структуру в памяти процесса Camunda не очень практично, так как все запущенные процессы периодически сбрасываются в базу данных и база очень быстро становится бутылочным горлышком (плюс к этому еще и чудовищно разбухает, так как при каждом коммите в базу обновляется вся запись с переменной модели в bpmn-процессе). Поэтому несколько лет назад было принято решение вынести хранение модели скоринга в память, а в Camunda оставить только ссылку на ключ в кэше. Таким образом у нас в команде появился in-memory data grid, замаскировавшийся под обычный кэш.

До определенного момента такая архитектура работала и работала хорошо. До тех пор пока не начала отказывать нода, где был развернут Redis. Дело в том что, ввиду ограниченности ресурсов и времени, кэш был развернут в единственном экземпляре. При выходе из строя одной ноды выходил из строя полностью весь процесс скоринга. Когда Redis внедряли, то делали быстро, без проектирования и оглядки на общепринятые стандарты. Нужно было сделать хоть что-то, потому что одна из метрик SLA в нашей команде это общее время скоринга, а Redis болячку закрывал.

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

  1. Система должна быть бесплатной (оплачивать лицензию не хочется, ровно как и подписку на SaaS)

  2. Система должна уметь автоматически горизонтально масштабироваться (минимум 3 инстанса и чтобы для этого не нужно было ничего дополнительно делать)

  3. Система должна переживать падение одного из инстансов (без деградации общей производительности, а значит все данные должны храниться на всех инстансах)

  4. Система должна иметь автоматическую репликацию, в идеале иметь толерантность к split brain

  5. Система не должна жрать ресурсы как не в себя (максимум 2-3 ГБ RAM на инстанс и пару ядер)

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

  7. Также система должна быть дружелюбна к Java (мы стараемся держать стек гомогенным и не привносить без особой надобности другие языки программирования)

  8. Система должна поддерживать распределенные блокировки

  9. Система должна поддерживать персистентность кэшей (в идеале в PostgreSQL)

Redis не подходил сразу по 3 параметрам. Нет, конечно на сегодня уже существуют Sentinel и кластерный Redis, но все они не подходят по критериям или простоты поддержки или потребления ресурсов или дружелюбности к экосистеме Java — все-таки Redis написан на C и для его настройки и тюнинга требуется отдельная (причем недешевая) компетенция.

Замечательный Apache Ignite подходил почти по всем параметрам кроме ресурсов и простоты поддержки. Классные CouchDB и CockroachDB отсекались примерно по той же причине, к тому же были не особо java friendly. Монгу ставить никто не хотел — это вообще не кэш, ровно как и всякие Aerospike и DynamoDB. Да еще и платные.

Когда-то давно, на этом же проекте я хотел внедрять Hazelcast, так как на тот момент он выбивал страйк — закрывал 9 из 9 требований (внедрить не успел — ушел работать в другую компанию). К тому же Hazelcast имел крутую фичу, так называемый NearCache — это когда горячие данные складываются прямо в хип сервиса и поход на сервер Hazelcast для чтения не осуществляется. Хранение модели подходило под этот сценарий. Сэкономить 5-10 мс на каждой операции чтения/записи дает прирост в 5-10% на потоке. Это много. Горизонтальное масштабирование вшито в архитектуру этого решения, есть настраиваемая репликация, от потери данных из-за падения ноды защищает репликация и алгоритм консенсуса Raft (как это было раньше в kafka), а встроенная CPSubsystem гарантирует корректную работу распределенных блокировок (CP это как раз consistency и partition tolerance из CAP теоремы). К тому же Hazelcast написан на Java и поддерживает все LTS версии вплоть до 25. Добавляя его в инфраструктуру мы сразу же можем распространить на него наши внутренние практики — тюнинга GC, мониторинга, алертинга. Взвесив все за и против, посоветовавшись с коллегами, освоив некоторое количество гайдов по интеграции, было принято решение внедрять Hazelcast.

Итак, чтобы развернуть Hazelcast в k8s, помимо специфичных для вашей системы вещей, вам необходимо всего лишь развернуть DaemonSet, как, например, приведенный ниже:

apiVersion: apps/v1kind: DaemonSetmetadata:  name: hazelcast  namespace: hazelcastspec:  selector:    matchLabels:      app: hazelcast  template:    metadata:      labels:        app: hazelcast    spec:      serviceAccountName: hazelcast      containers:        - name: hazelcast          image: hazelcast/hazelcast:5.7.0          ports:            - containerPort: 5701          env:            - name: JAVA_OPTS              value: "-Xms512m -Xmx1g -Dhazelcast.config=/opt/hazelcast/config/hazelcast.yaml"          volumeMounts:            - name: hazelcast-config              mountPath: /opt/hazelcast/config/hazelcast.yaml              subPath: hazelcast.yaml          livenessProbe:            httpGet:              path: /hazelcast/health/ready              port: 5701            initialDelaySeconds: 30            periodSeconds: 10          readinessProbe:            httpGet:              path: /hazelcast/health/ready              port: 5701            initialDelaySeconds: 20            periodSeconds: 5          resources:            requests:              memory: "1Gi"              cpu: "500m"            limits:              memory: "1500Mi"              cpu: "2000m"      volumes:        - name: hazelcast-config          configMap:            name: hazelcast-config

Cобственно, все. Настройка самого hazelcast сводится к тому что вы говорите ему в файле hazelcast.yaml что он работает в k8s и может сам себя обнаружить в таком-то неймспейсе:

hazelcast:  cluster-name: your-cluster  network:    join:      multicast:        enabled: false # мультикаст выключаем, так как в k8s окружении он работает плохо и добавляет ощутимые задержки      kubernetes:        enabled: true # вместо мультикаста говорим что hazelcast работает в k8s        namespace: your-namespace  # и указываем конкретный неймспейс        service-name: hazelcast-service # и даже сервис, чтобы не блукать по кластерным дебрям в поисках других инстансов  properties:    hazelcast.discovery.enabled: false # можно не включать, но если hazelcast.kubernetes.enabled = false, то нужно включать обязательно

Настройки кэша могут задаваться как на сервере, так и на клиенте. В нашем случае мы использовали серверный вариант, так как количество кэшей одинаковое и меняется редко. Все взаимодействие с Hazelcast мы вынесли в стартер и уже там при создании кэша указывали:

hazelcast:  map:    map-name:      # 3 инстанса, 1 мастер, 2 реплики, падение 1 ноды не выводит из строя весь кластер, даже если это мастер нода      backup-count: 2       # асинхронная репликация менее надежна, но в случае если потеря данных допустима или серверов много - можно сделать часть копий асинхронными      async-backup-count: 0      # разрешаем чтение из бэкапа, CQRS как он есть      read-backup-data: true

Первую версию стартера мы написали и отдалили за 2 дня. Подключили в проект, настроили запись модели, пока что без персистентности, развернули в прод кластерах Hazelcast, задеплоили новую версию, переключили и начали измерять. И ничего не изменилось. Метрики показывали точно такую же производительность. Деградации не было ни в первый час, ни к концу дня, ни через неделю. Улучшения производительности не было, но и деградации не было также. Но чутье мне говорило — не может быть все так просто. Где-то должен быть подвох. Не бывает так что с первого раза все работает. Чутье меня не подвело — проблемы начались сразу же как мы пошли дальше.

Первым делом нужно было прикрутить персистентность. Hazelcast в этом плане достаточно демократичен. Если вы хотите персистентность, то вы можете написать собственный класс, который будет заниматься сохранением в и извлечением из базы записей. Структуру таблицы, пути извлечения и записи вы определяете сами. При старте приложения ваш класс передается на сервер и там же выполняется. Звучит круто, но тут и начинаются большие проблемы. Главная проблема заключается в том что кэши на стороне сервера иммутабельные. Их невозможно перезаписать ни при каких условиях. Версионность кэшей не поддерживается. А это означает что если вы выпустили новую версию вашего класса для персистентности и хотите обновить его на сервере Hazelcast — сделать этого не получится, придется переименовывать кэш и создавать заново. Либо полностью стирать кэш на сервере и создавать заново.

Далее, так как ваш класс загружается в другой classpath, в нем может не оказаться необходимых библиотек. Hazelcast уже содержит в себе некоторые библиотеки — например драйвер для PostgreSQL, но все остальное попросту отсутствуют. Если вы хотите использовать JOOQ или Spring Data с Hibernate — то вы не хотите использовать JOOQ или Spring Data с Hibernate. У вас нет стандартных средств загрузить внутрь Hazelcast сторонние библиотеки. Да и даже если вы придумаете такой способ, то это уже будет не совсем коробочное решение. Это будет то от чего как раз мы хотели уйти — пострадают простота поддержки и разработки. То есть как бы вы не хотели писать по-своему — разгуляться не получится. В распоряжении стандартное JDBC API и вперед маслать. Да, если вы хотите загружать в какое-то экзотическое хранилище и драйвера к этому хранилищу нет на сервере, то у вас есть возможность инициализировать подключение самостоятельно реализовав всю логику в методе init класса реализующего MapLoaderlifeCycleSupport. Но это все не то.

Мы написали реализацию на базе JDBC API и закрыли вопрос. Но также не забывайте, что проектирование таблицы для хранения кэша тоже теперь лежит на ваших плечах. Если вы спроектируете таблицу, в которой не будет индексов, на каждый запрос будет производиться фулскан, то и результат моментально скажется на производительности всего кэша. В идеале вам нужна хэш-таблица, но только в базе. У Hazelcast есть возможность настраивать момент когда запись из памяти попадет в бд через write-delay-seconds. Если выставлено 0 — вы получаете синхронную репликацию из памяти в бд. В остальных случаях репликация асинхронная. Возможность группировки изменений в пакеты также есть и настраивается с помощью write-batch-size. Еще из полезного можно при старте приложения загрузить сразу весь кэш в память меняя initial-mode на EAGER — полезно если вам нужно вынести загрузку кэшей на этап запуска приложения, главное не забыть вызвать на старте метод IMap.getMap.

Следующим камнем преткновения стала CPSubsystem. Начиная с версии 5.4.0 это платный функционал. Да, можно использовать устаревшую версию 5.3.4, но это явно не наш вариант — мы стараемся использовать только последние стабильные версии библиотек и фреймворков. В последней релизной версии (5.7.0 на момент написания статьи) CPSubsystem попросту недоступна и нам пришлось использовать для блокировок стандартный IMap.tryLock, у которого есть проблемы с транзакциями, дедлоками, фризом на длительный промежуток времени при малых значениях таймаута. Чуть менее надежно, чем ожидалось, но на практике пока что ни разу не столкнулись с каким-то проблемами.

Весь процесс перехода с Redis на Hazelcast занял 2 недели. Разумеется, по закону подлости, аккурат во время перехода вновь случился отказ оборудования — одна из нод кластера ушла в офлайн. С единственным отличием — в этот раз аварии не получилось. После того как нода отвалилась мы получили алерт и в этот же момент поды, запущенные на умершей ноде уже мигрировали на оставшиеся в строю ноды, запустились и продолжили работать как ни в чем не бывало! Ноль секунд простоя, ни одного зависшего процесса, ни одной ошибки. Спустя 2 часа ноду оживили и вернули в строй, автоматическая балансировка разогнала поды обратно и все это происходило под нагрузкой — нам не пришлось предпринимать никаких дополнительных усилий. Мы не поднимали по звонку дежурного разработчика. Бизнес даже не узнал о том что произошел отказ. И я хочу сказать — это круто. Такой классный результат, о котором мы даже не мечтали. Спустя 3 месяца эксплуатации мы научились выявлять и исправлять проблемы с памятью, подкрутили настройки GC, перешли на jdk 25 на сервере hazelcast и забыли про существование такого вопроса как кэш. В новых сервисах мы просто подключаем зависимость, пишем пару классов и пользуемся благами цивилизации.

В случае масштабирования — при добавлении нод в кластер никаких дополнительных приседаний не потребуется — DaemonSet сам развернет копию Hazelcast на этой ноде. Если понадобится еще более избыточное реплицирование — мы поменяем 1 цифру в ConfigMap и сделаем rollout restart. Проверено на практике — даже под нагрузкой это делать безопасно.

Если статья понравилась — подписывайтесь. Если есть чем поделиться — приходите в комментарии, буду рад пообщаться и почитать про ваш опыт работы с Hazelcast. В скором времени я расскажу как мы готовим переход на GitOps и настраиваем автоматическое масштабирование сервисов, если пришло много траффика и есть риск OOMKill. Разумеется, все сами, дешево и сердито, но надежно.

P.S.: Вы, возможно, могли заметить что я не упомянул Infinispan — братишку Hazelcast, но только полностью open source, некогда бывший под эгидой RedHat, а нынe Commonhaus Foundation. И да, нас буквально сбил с толку GitHub проекта, где Infinispan, дословно In-Memory Distributed Database. Мы ошибочно записали его в один ряд с мастодонтами вида Apache Ignite, но на поверку оказалось что его можно и нужно использовать у нас. Infinispan из коробки умеет в распределенные блокировки (ClusteredLock с настройкой Reliability.CONSISTENT), а также за персистентность отвечает сам — вы не сможете загрузить в него свой код (максимум что вам доступно — настроить типы и названия некоторых колонок). Именно поэтому мы уже запланировали переход c Hazelcast на Infinispan.

ссылка на оригинал статьи https://habr.com/ru/articles/1065852/