Эта статья выбрана из серии статей «Практика инфраструктуры Byte Beat».
Серия статей «Практика инфраструктуры ByteDance» представляет собой техническую галантерею, созданную техническими командами и экспертами отдела инфраструктуры ByteDance, в которой мы делимся с вами практическим опытом и уроками команды в процессе развития и эволюции инфраструктуры, а также всеми техническими студентами. общаться и расти вместе.
С тех пор, как Google выпустил документ Spanner, соответствующие продукты или услуги баз данных были запущены в стране и за рубежом для решения проблемы масштабируемости базы данных. ByteDance также приняла соответствующие технические решения для удовлетворения огромных потребностей в хранении данных. Этот обмен расскажет о проблемах, решениях и развитии технологий, с которыми мы столкнулись при создании таких систем.
Обзор прошлой ситуации
ключевая технология
Далее мы продолжаем обсуждать ключевые технологии распределенных транзакций, автоматического разделения и объединения разделов и балансировки нагрузки.
Распределенная транзакция
Представляя интерфейсную часть, мы упомянули атомарные WriteBatch и MultiGet ByteKV, которые удовлетворяют распределенному согласованному чтению моментальных снимков. WriteBatch означает, что все изменения в пакете либо выполняются успешно, либо завершаются неудачно, и не будет частичного успеха или частичного отказа. MultiGet означает, что части данных из других совершенных транзакций не будут считаны.
ByteKV примерно использует следующие технологии для реализации распределенных транзакций:
- Кластер обеспечивает глобально увеличивающиеся логические часы, и каждому запросу на чтение и запись присваивается отметка времени через этот модуль, тем самым присваивая глобальный порядок всем запросам.
- Каждое обновление ключа создает новую версию в системе, гарантируя, что новые операции записи не повлияют на моментальные снимки старых операций чтения.
- Двухэтапная фиксация введена в процесс запроса на запись, чтобы гарантировать, что записи могут быть отправлены упорядоченным и атомарным образом.
Служба глобального хронометража
Нет сомнений в том, что упорядочение всех событий упрощает многие проблемы в распределенных системах. Мы также всегда видим, как различные системы выбирают между различными физическими часами, логическими часами и смешанными логическими часами. С точки зрения производительности, стабильности и сложности реализации ByteKV реализует интерфейс, который обеспечивает глобальное инкрементное выделение временных меток в службе KVMaster, которая используется всеми модулями чтения и записи в кластере.Этот интерфейс гарантирует, что выдаваемые временные метки являются глобально уникальными. и Увеличение.
Эта структура принята, потому что мы считаем, что:
- Логика распределения часов очень проста, даже если она обеспечивается одним модулем, он может получить стабильную задержку и достаточную пропускную способность.
- Мы можем использовать протокол Raft для достижения высокой доступности модуля распределения часов, и выход из строя отдельной машины никогда не станет единой точкой системы.
Что касается конкретной реализации, чтобы обеспечить стабильность, эффективность и простоту использования часов, мы также предприняли некоторые инженерные усилия и оптимизации:
- Логика одного и того же клиента, принимающего часы, объединяется в пакеты, что может эффективно сократить количество RPC.
- При распределении часов следует использовать независимый сокет TCP, чтобы избежать помех от других запросов RPC.
- Распределение часов использует атомарные операции, полностью избегая использования блокировок.
- Часы должны быть максимально приближены к реальному физическому времени, что очень полезно для отладки некоторых проблем.
несколько версий
Почти все современные системы баз данных используют многоверсионный механизм как часть механизма контроля параллелизма транзакций, и ByteKV не является исключением. Преимущество нескольких версий в том, что операции чтения и записи не блокируют друг друга. Каждая запись в строку создает новую версию, а чтение обычно считывает существующую версию. Логическая организация данных выглядит следующим образом:
Несколько версий одного и того же ключа будут постоянно храниться вместе, чтобы облегчить поиск определенных версий, а версии расположены в порядке убывания, чтобы уменьшить накладные расходы на запросы.
Чтобы гарантировать, что закодированные данные могут быть отсортированы так, как мы ожидаем, мы используем сопоставимое с памятью кодирование для ключа RocksDB [2].Причина, по которой здесь нет пользовательской функции сравнения RocksDB, заключается в следующем:
- Ключевым размером сравнения является очень высокая частота чтения и записи движком, а memcmp по умолчанию очень дружелюбен к производительности.
- Уменьшите особую зависимость от RocksDB и улучшите гибкость архитектуры.
Чтобы избежать бесконечного расширения пространства, вызванного непрерывным накоплением нескольких версий одного и того же ключа, у ByteKV есть фоновая задача по периодической очистке старых версий и данных, помеченных для удаления. В последней части главы, посвященной механизму хранения, было сделано несколько вступлений.
двухэтапная фиксация
ByteKV использует двухэтапную фиксацию для реализации распределенных транзакций.Общая идея состоит в том, что весь процесс разделен на две фазы: первая фаза называется фазой подготовки.На этом этапе координатор отвечает за отправку запросов на подготовку участникам, а участники отвечают на запросы и распределяют Ресурсы, пре-коммит (мы называем пре-коммит данных Write Intent); после успешного выполнения всех участников первого этапа координатор запускает второй этап, Commit stage, в котором координатор совершает коммит транзакцию и дает все Участник отправляет команду отправки, а участник отвечает на запрос и конвертирует Write Intent в реальные данные. В ByteKV координатором является KVClient, а участниками являются все PartitionServers. Далее давайте рассмотрим некоторые детали реализации распределенных транзакций ByteKV с точки зрения атомарности и изоляции.
Во-первых, как обеспечить видимость атомарности транзакций для внешнего мира? Проблема, по сути, заключается в необходимости иметь постоянное состояние транзакций, которое можно модифицировать атомарно. В отрасли существует множество решений.Метод, принятый ByteKV, заключается в том, чтобы рассматривать состояние транзакции как обычные данные и хранить их во внутренней таблице отдельно. Мы называем эту таблицу таблицей состояний транзакций, и, как и другие бизнес-данные, она также распределяется по нескольким машинам. Таблица состояния транзакций включает следующую информацию:
- Статус транзакции: в том числе статус транзакции была начата, зафиксирована, отменена и так далее. Состояние транзакции само по себе является KV, и легко добиться атомарности.
- Номер версии транзакции: отметка времени, полученная от глобальных инкрементных часов при фиксации транзакции. Этот номер версии будет закодирован во всех ключах, измененных транзакцией.
- TTL транзакции: период ожидания транзакции, в основном для решения ситуации, когда транзакция мертва и все время занимает ресурсы. Когда другие транзакции получают доступ к ресурсам, измененным транзакцией, если они обнаруживают, что время ожидания транзакции истекло, они могут принудительно завершить транзакцию.
С помощью таблицы состояний транзакций на втором этапе координатору достаточно просто изменить состояние транзакции, чтобы завершить операции фиксации и отката транзакции. После завершения изменения состояния транзакции клиенту можно успешно ответить, а операции фиксации и очистки намерения записи выполняются асинхронно.
Второй вопрос: как обеспечить изоляцию и обработку конфликтов между транзакциями? ByteKV будет сортировать выполняемые транзакции в порядке поступления, а более поздняя транзакция будет ждать, пока будет прочитано намерение записи, пока предыдущая транзакция не завершится и намерение записи не будет очищено. Намерение записи невидимо для запроса на чтение.Если время подготовки транзакции, на которое указывает намерение записи, превышает время транзакции чтения, намерение записи будет проигнорировано, в противном случае запрос на чтение должен дождаться предыдущей транзакции. для завершения или отката, чтобы узнать, доступны ли данные для чтения. Ожидание фиксации транзакции может повлиять на задержку запроса на чтение.Простая оптимизация заключается в том, что запрос на чтение подталкивает отметку времени фиксации незафиксированной транзакции к отметке времени транзакции чтения. Ранее было сказано так много намерений записи, так как же кодируется намерение записи, чтобы данные, которые не были отправлены во время операции транзакции, не могли быть прочитаны другими транзакциями? Это также относительно просто, просто установите номер версии Write Intent на бесконечность.
Помимо вышеперечисленных проблем, распределенным транзакциям необходимо решить проблему отказоустойчивости. Здесь обсуждается только сценарий сбоя координатора.После сбоя координатора транзакция может находиться в зафиксированном состоянии или в незафиксированном состоянии, некоторые Write Intent в PartitionServer могут быть зафиксированы или очищены, или могут остаться там. Если транзакция была зафиксирована, когда последующая транзакция чтения-записи встретит оставшееся намерение записи, она поможет предыдущей транзакции зафиксировать или очистить намерение записи в соответствии с состоянием в таблице состояний транзакций; если транзакция не была зафиксирована, последующая транзакция будет в предыдущей транзакции.По истечении времени ожидания (TTL транзакции) состояние транзакции изменяется для отката, а намерение записи очищается асинхронно. Поскольку само намерение записи также содержит информацию, связанную с транзакцией, если мы также запишем список участников намерения записи, мы можем изменить флаг фиксации транзакции с атомарного изменения состояния транзакции на все намерения записи, завершающие сохраняемость, тем самым уменьшив задержку фиксация; и последующие операции могут восстановить состояние транзакции в соответствии со списком участников после обнаружения намерения записи.
Разделы автоматически разделяются и объединяются
Как упоминалось ранее, ByteKV использует разбиение по диапазонам для обеспечения масштабируемости.Одна из проблем, вызванных этим методом разбиения, заключается в том, что с развитием бизнеса исходная структура разделов больше не подходит для новой бизнес-модели. Например, служба записывает изменения хотспота, а хотспоты перемещаются с одного раздела на другой. Для того, чтобы решить эту проблему, ByteKV реализует функцию автоматического разделения: путем выборки пользовательской записи, когда объем данных превышает определенный порог, Range делится на два новых Range из середины. Функция разделения вместе с планированием обеспечивает возможность автоматического расширения.
В других сценариях, таких как большое количество TTL, большое количество операций записи перед удалением, большое количество разделов будет автоматически разбито. Когда TTL истекает и данные подвергаются сборке мусора, эти разделенные разделы образуют большое количество фрагментов данных: каждая группа Raft обслуживает только небольшой объем данных. Эти небольшие разделы вызывают бессмысленные накладные расходы, а сохранение их метаинформации также увеличивает нагрузку на KVMaster. В ответ на эту ситуацию ByteKV реализует функцию автоматического слияния, чтобы объединить некоторые меньшие интервалы со смежными интервалами.
Процесс слияния более сложен, чем расщепление.Мастер планирует объединить два соседних интервала в один блок, а затем инициирует операцию слияния. Как показано на рисунке выше, этот процесс разделен на два шага: сначала инициируется операция в левом интервале, получается точка синхронизации, затем инициируется операция слияния в правом интервале, правый интервал будет ждать, до тех пор, пока данные до точки синхронизации левого интервала на текущем сервере. Если все выполняется синхронно, метаинформация левого и правого интервалов может быть безопасно изменена для завершения операции слияния.
балансировки нагрузки
Балансировка нагрузки — одна из важных возможностей, необходимых для всех распределенных систем. Система, которая не может добиться балансировки нагрузки, не только не может полностью использовать вычислительные и складские ресурсы кластера, но и влияет на качество обслуживания из-за джиттера отдельных узлов из-за большой нагрузки. Разработка хорошей стратегии балансировки нагрузки столкнется с двумя трудностями: во-первых, существует множество аспектов ресурсов, которые необходимо сбалансировать, не только самое основное дисковое пространство, но и ЦП, ввод-вывод, пропускная способность сети, объем памяти и т. д.; второй — byte beating. Внутренние характеристики машин очень разнообразны, и разные узлы в одном кластере могут иметь разные процессоры, диски и память. Мы приняли пошаговый подход к разработке стратегии балансировки нагрузки, сначала рассматривая сценарий одномерной однородной модели, а затем расширяя его до многомерных гетерогенных моделей. Далее описывается эволюция стратегии.
Стратегия одномерного планирования
Возьмем в качестве примера одно измерение дискового пространства и предположим, что все узлы имеют одинаковую емкость диска. Использование дискового пространства каждым узлом равно сумме объемов данных всех реплик на этом узле. Выделение и размещение всех реплик по одной на определенном узле формирует схему размещения реплик. Должен быть план, значение дисперсии объема данных каждого узла самое низкое, это состояние называется «абсолютным равновесием». Поскольку данные продолжают записываться, объем данных на узле будет продолжать меняться.Чтобы поддерживать кластер в «абсолютно сбалансированном» состоянии, требуется непрерывное планирование, что приводит к большим накладным расходам на миграцию данных. Мало того, абсолютное равновесие в одном измерении сделает невозможным абсолютное равновесие в других измерениях. С точки зрения затрат и осуществимости мы определяем более слабое состояние равновесия, называемое «достаточным равновесием», которое ослабляет стандарт равновесия, с одной стороны, снижает чувствительность планирования, и небольшое количество данных не будет изменяться. вызывает частое планирование и, с другой стороны, позволяет нескольким измерениям достичь этого состояния слабого равновесия одновременно. Чтобы интуитивно выразить определение «достаточного равновесия», нарисуем такую схематическую диаграмму для иллюстрации:
- Каждый узел — это столбец, высота столбца — это его объем данных, и все узлы расположены в порядке от высокого к низкому.
- Рассчитайте средний объем данных Savg всех узлов и начертите горизонтальную линию, называемую средней линией.
- Добавьте или вычтите альфа-значение из среднего объема данных, чтобы получить значение максимального и минимального уровня воды. Альфа может составлять 10% или 20% от Savg, что определяет степень жесткости баланса. линия отметки уровня воды и линия отметки низкого уровня воды в соответствии со значением уровня воды.
- В соответствии с соотношением между объемом данных узла и тремя линиями они делятся на четыре области:
- Область высокой нагрузки/область активной миграции: объем данных узла выше максимального значения
- Область высокого баланса/зона пассивной миграции: объем данных узла ниже верхней отметки и выше среднего
- Зона низкого равновесия/зона пассивной миграции: объем данных узла выше минимальной отметки и ниже среднего
- Область низкой нагрузки/область активной миграции: объем данных узла ниже нижнего порогового значения.
- Когда узел находится в зоне высокой нагрузки, ему необходимо активно перемещать копию, а целевой узел находится в области перемещения; когда узел находится в зоне низкой нагрузки, ему необходимо активно перемещаться внутрь. копия, а исходный узел — область перемещения
- Когда все узлы находятся в двух зонах равновесия, кластер достигает «достаточно сбалансированного» состояния.На следующем рисунке показано «достаточно сбалансированное» состояние.
Многомерная стратегия планирования
Основываясь на предыдущем одномерном планировании, цель многомерного планирования состоит в том, чтобы заставить кластер достичь достаточно сбалансированного состояния в нескольких измерениях одновременно или в максимально возможном количестве.
Давайте сначала представим, что у каждого измерения есть вышеупомянутая схематическая диаграмма, представляющая его состояние равновесия, и есть N графиков в N измерениях. При переносе копии равновесное состояние всех измерений будет изменено одновременно, то есть все схематические диаграммы будут изменены. Если все измерения становятся более сбалансированными (больше узлов в зоне баланса) или если некоторые измерения становятся более сбалансированными, а другие остаются прежними (количество узлов в зоне баланса остается прежним), то миграция — это хорошее расписание; в любом случае , если все измерения становятся более несбалансированными (количество узлов в области баланса становится меньше) или если одни измерения становятся более несбалансированными, а другие остаются неизменными, то миграция является плохим графиком. Существует также третья ситуация, некоторые измерения более сбалансированы, а некоторые измерения более несбалансированы. Это нейтральное планирование, и это нейтральное планирование часто неизбежно. Например, в кластере есть только два узла A и B. трафик выше и объем данных B выше, перенос реплик из A в B сделает трафик более сбалансированным, а объем данных более несбалансированным, а перенос реплик из B в A будет противоположным.
Чтобы судить о том, допустимо ли такое нейтральное планирование, мы вводим понятие приоритета, присваивая уникальный приоритет каждому измерению и приемлемо жертвуя балансом низкокачественных измерений в обмен на более сбалансированные высококачественные измерения. , недопустимо жертвовать балансом качественных габаритов в обмен на более сбалансированные низкокачественные габариты.
Продолжая рассматривать предыдущий пример, поскольку слишком большой трафик повлияет на время отклика при чтении и записи и, таким образом, повлияет на качество обслуживания, мы считаем, что приоритет трафика выше, чем приоритет объема данных, поэтому миграция из A в B будет приемлемый. Но есть исключение: если предположить, что оставшееся дисковое пространство узла B близко к 0, и даже самая маленькая реплика в кластере не может его вместить, даже если приоритет трафика лучше, ни одна реплика не должна быть допущена к миграции на B. . Чтобы визуально выразить это состояние насыщения ресурсами, мы добавляем на схематическую диаграмму жесткую ограничительную линию:
На этой диаграмме стратегия многомерной балансировки нагрузки выглядит следующим образом:
- Отсортируйте несколько измерений в соответствии с их приоритетами и выполните одномерную стратегию планирования, описанную выше, в последовательности от высокооптимизированного измерения к низкооптимизированному измерению с небольшими изменениями в процессе:
- ближайший к исходному узлу
Sbestно меньше чемSbestКопия является кандидатом на миграцию, и если она вызывает какое-либо из следующих условий, она исключается, и в качестве кандидата выбирается следующая копия до тех пор, пока не будет найдена подходящая копия:- После миграции целевая машина будет выше максимальной отметки в более высоких оптимальных размерах.
- После миграции целевая машина будет выше жесткого предела в нижнем оптимальном измерении.
- Если для определенного целевого узла объект миграции не может быть выбран на исходном узле, узел, ранжированный перед целевым узлом, берется в качестве нового целевого узла, и описанный выше процесс повторяется.
- Если для всех целевых узлов объект миграции по-прежнему не может быть выбран на исходном узле, удалите исходный узел из отсортированного списка и повторите описанный выше процесс.
Стратегия планирования гетерогенной модели
Для гомогенных моделей единица нагрузки будет использовать одинаковую долю ресурсов на каждом узле, и мы можем планировать только на основе значения нагрузки, независимо от того, сколько машинных ресурсов используют эти нагрузки, но для гетерогенных моделей это недопустимо. Например, чтение 1 МБ данных с диска может занимать только 1 % пропускной способности ввода-вывода и 1 % циклов ЦП на высокопроизводительном сервере, в то время как на виртуальном сервере это может занимать 5 % пропускной способности ввода-вывода и 3 % циклов ЦП. машина % циклов процессора. На узлах с разной производительностью одна и та же нагрузка приведет к разному использованию ресурсов.
Чтобы применить предыдущую политику планирования к сценарию гетерогенных моделей, во-первых, необходимо изменить планирование на основе значения нагрузки на планирование на основе использования ресурсов. Для объема данных его следует изменить на использование дискового пространства, для трафика — на использование ЦП, использование ввода-вывода и т. д. Чтобы упростить стратегию, мы включаем использование памяти, дискового ввода-вывода, сетевого ввода-вывода и т. д. в использование ЦП. Объяснить, почему:
- Для памяти верхний предел использования памяти нашего процесса контролируется элементами конфигурации.Во время развертывания мы гарантируем, что использование памяти не будет превышать размер физической памяти, а оставшаяся физическая память будет использоваться для буфера/кеша операционной системы. , действительно могут быть использованы нами. Размер памяти будет влиять на производительность узла, влияя на размер MemTable и BlockCache, и этот эффект в конечном итоге отразится на использовании ЦП и ввода-вывода, поэтому мы можем учитывать использование памяти, изучая использование ЦП и ввода-вывода. . . .
- Для дискового ввода-вывода использование ввода-вывода в конечном итоге будет отражаться в использовании ЦП (синхронный ввод-вывод отражается в wa, а асинхронный ввод-вывод отражается в sys), поэтому мы можем учитывать использование дискового ввода-вывода, исследуя загрузку ЦП.
- В ЦП есть три уровня кеша, а также регистры, при рассмотрении загрузки ЦП он будет рассматриваться как единое целое, а использование кешей или регистров отдельно анализироваться не будет. Память и диск можно представить как кэши четвертого и пятого уровня ЦП. Чем меньше память, тем медленнее дисковый ввод-вывод и выше загрузка ЦП. Их можно рассматривать как единое целое.
Вторая проблема, которую необходимо решить с помощью гетерогенного планирования, — это отношение преобразования между использованием ресурсов и значением нагрузки. Например, загрузка ЦП узлов A и B составляет 50 % и 30 % соответственно, а также известны запросы на чтение и запись каждой реплики на узле.Как выбрать лучшую реплику с узла A для миграции на узел B, Чтобы свести к минимуму разрыв в использовании ЦП между A и B, мы должны рассчитать, сколько ресурсов ЦП будет генерировать каждая реплика на узлах A и B. Для этого мы собираем как можно больше информации о запросах на чтение и запись для каждой реплики, например:
- Размер ключа и значение запроса на чтение и запись
- скорость чтения кеша
- Степень рандомизации обновлений, доля удалений
На основе этой информации каждый запрос на чтение и запись преобразуется в N стандартных потоков. Например, запрос в пределах 1 КБ — это один стандартный трафик, а запрос размером 1–2 КБ — это два стандартных трафика; запрос, попадающий в кеш, — это один стандартный трафик, а запрос, не попадающий в кеш, — это два стандартных трафика. Зная общее значение стандартного трафика на узле, вы можете рассчитать загрузку ЦП, соответствующую стандартному трафику на этом узле, в соответствии с загрузкой ЦП, а затем вы можете рассчитать загрузку ЦП, соответствующую каждой реплике на каждом узле.
Таким образом, стратегия планирования для разнородных моделей требует лишь внесения следующих изменений на основе стратегии многомерного планирования:
- Узлы сортируются по использованию ресурсов, а не по значению нагрузки.
- Значение нагрузки каждой копии должно быть преобразовано в коэффициент использования ресурсов исходного узла и коэффициент использования ресурсов целевого узла соответственно.В гетерогенных моделях коэффициент использования ресурсов одной и той же копии будет сильно различаться.
Другая политика планирования
В KVMaster есть временная задача для выполнения вышеуказанной стратегии балансировки нагрузки, называемая «планировщик балансировки нагрузки», которая не будет здесь повторяться; в то же время есть еще одна временная задача, используемая для выполнения другого типа планирования, называемого « Планировщик размещения копий" ", в дополнение к базовым стратегиям, таким как уровень безопасности реплики (центр обработки данных/стойка/сервер) и обнаружение аномалий узлов, он также реализует следующие стратегии планирования:
- Стратегия изоляции бизнеса: разные пространства имен/таблицы могут храниться на разных узлах. Каждое пространство имен/таблица может указывать тег строкового типа, а каждый узел может указывать один или несколько тегов.Только когда тег пространства имен/таблицы, где находится копия, совпадает с тегом узла, его можно разместить на узел. Планировщик запланирует реплики, которые не соответствуют требованиям тега.
- Обнаружение горячих точек: когда объем данных фрагмента данных достигает определенного порога, происходит разделение.Кроме того, когда его трафик чтения и записи превышает определенное значение, кратное среднему значению, также происходит разделение. Когда происходит разделение, все копии одного из вновь созданных сегментов (левого или правого) будут перенесены на другие узлы, чтобы узлы не стали точками доступа.
- Обнаружение фрагмента: когда объем данных и трафик чтения и записи фрагмента данных меньше определенной доли среднего значения, он будет объединен с соседними фрагментами. Перед слиянием все копии маленьких осколков будут перенесены на узлы, где расположены соседние осколки.
слой таблицы
Как упоминалось ранее, модель данных KV слишком проста для удовлетворения потребностей некоторых сложных бизнес-сценариев. Например:
- Слишком много полей и типов
- Запрос со сложными условиями для разных размеров поля
- Поля или измерения запросов часто меняются вместе с требованиями.
Нам нужна более богатая модель данных для удовлетворения потребностей этих сценариев. Поверх слоя KV мы строим слой таблицы ByteSQL, реализованный вышеупомянутым SQLProxy. ByteSQL поддерживает запись и чтение с помощью языка структурированных запросов (SQL) и реализует интерактивные транзакции, которые поддерживают смешанные операции чтения-записи на основе интерфейсов WriteBatch и чтения моментальных снимков ByteKV.
Табличная модель
В модели табличного хранения данные организованы и хранятся в соответствии с двумя логическими уровнями базы данных и таблицы. В одном физическом кластере можно создать несколько баз данных, а в каждой базе данных можно создать несколько таблиц. Определение схемы таблицы содержит следующие элементы:
- Основные свойства таблицы, включая имя базы данных, имя таблицы, количество копий данных и т. д.
- Определение поля: содержит имя поля, тип, разрешать ли нулевые значения, значения по умолчанию и другие атрибуты. Таблица должна содержать хотя бы одно поле.
- Определение индекса: содержит имя индекса, список полей, включенных в индекс, и тип индекса (первичный ключ, уникальный ключ, ключ и т. д.). В таблице есть только один индекс первичного ключа (первичный ключ), и пользователи могут также добавить вторичный индекс (тип ключа или уникального ключа) для повышения производительности выполнения SQL. Каждый индекс может быть индексом с одним полем или совместным индексом с несколькими полями.
Каждая строка в таблице кодируется в несколько записей KV, хранящихся в ByteKV, в соответствии с индексом, и метод кодирования каждого типа индекса отличается. Строка первичного ключа содержит значения всех полей таблицы, а строка вторичного индекса содержит только поля, определяющие индекс и первичный ключ. Метод кодирования каждого индекса следующий:
Primary Key: pk_field1, pk_field2,... => non_pk_field1, non_pk_field2...
Unique Key: key_field1, key_field2,...=> pk_field1, pk_field2...
NonUnique Key: key_field1, key_field2,..., pk_field1, pk_field2...=> <null>
Где pk_field — это поле, определяющее первичный ключ, non_pk_field — это поле, которое не является первичным ключом в таблице, а key_field — это поле, определяющее вторичный индекс.=>Содержимое до и после соответствует частям Key и Value слоя KV соответственно. Кодирование ключевой части по-прежнему использует упомянутое выше кодирование, сравнимое с памятью, что гарантирует, что естественный порядок полей будет таким же, как порядок байтов после кодирования. Часть Value использует метод кодирования переменной длины, аналогичный protobuf, чтобы минимизировать размер закодированных данных. 1 байт используется при кодировании каждого поля, чтобы определить, является ли значение нулевым.
глобальный вторичный индекс
Пользователям часто приходится использовать поля непервичного ключа в качестве условий запроса, что требует создания вторичных индексов для этих полей. В традиционной архитектуре сегментирования (такой как кластер сегментов MySQL) выберите поле в таблице в качестве ключа сегментирования и разбейте всю таблицу на разные сегменты. Поскольку между разными сегментами не существует эффективного механизма распределенных транзакций, в каждом сегменте необходимо создавать вторичные индексы (т. е. локальные вторичные индексы). Проблема с этим решением заключается в том, что если условие запроса не содержит ключа сегментирования, необходимо просканировать все сегменты для объединения результатов, и глобальное ограничение уникальности не может быть реализовано.
Чтобы решить эту проблему, ByteSQL реализует глобальный вторичный индекс, который распределяет данные первичного ключа и данные вторичного индекса по разным сегментам ByteKV и может найти индекс индекса только в соответствии с условиями запроса на запись вторичного индекса, а затем найдите соответствующую запись первичного ключа. Этот метод позволяет избежать накладных расходов на сканирование всех сегментов для слияния результатов, а также может поддерживать глобальные ограничения уникальности за счет создания уникальных ключей, которые обладают сильной горизонтальной масштабируемостью.
интерактивная транзакция
Основываясь на многоверсионной функции ByteKV и атомарной записи (WriteBatch) нескольких записей, ByteSQL реализует транзакции чтения-записи, поддерживающие изоляцию моментальных снимков.Основные идеи реализации заключаются в следующем:
- Когда пользователь инициирует команду Start Transaction, ByteSQL получает глобально уникальную временную метку от ByteKV Master в качестве начальной временной метки транзакции (Start Timestamp).
- Все записи внутри транзакции кэшируются в локальном буфере записи ByteSQL, и каждая транзакция имеет свой собственный экземпляр буфера записи. Если это операция удаления, также запишите метку захоронения в буфере записи.
- Все операции чтения в транзакции сначала читают буфер записи и возвращаются напрямую, если в буфере записи есть запись (если возвращаемая запись Tombstone не существует в буфере записи); в противном случае попытайтесь прочитать запись в ByteKV, версия которой число меньше, чем Start Timestamp.
- Когда пользователь инициирует команду Commit Transaction, ByteSQL вызывает интерфейс WriteBatch ByteKV для отправки записей, кэшированных в буфере записи.В настоящее время отправка является условной: для каждого ключа в буфере записи необходимо убедиться, что есть не более чем Start Timestamp при отправке. Существуют более крупные версии. Если условие не выполняется, текущая транзакция должна быть прервана. Это условие реализуется через интерфейс CAS ByteKV.
Как видно из приведенного выше процесса, ByteSQL реализует обнаружение конфликтов транзакций в оптимистическом режиме. Этот режим очень эффективен в сценариях, где частота конфликтов при записи невелика. Если уровень конфликтов высок, транзакция будет часто прерываться.
Выполнить оптимизацию процесса
ByteSQL обеспечивает более богатую семантику SQL-запросов, но такие операции, как Get и Delete, добавляют дополнительные накладные расходы по сравнению с простыми операциями Put, Delete и другими операциями в модели KV. Операции Insert, Update и Delete в SQL на самом деле представляют собой процесс чтения перед записью. Возьмем, к примеру, Update. Сначала используйте операцию Get, чтобы прочитать старое значение из ByteKV, обновите некоторые поля старого значения в соответствии с предложением SQL Set, чтобы сгенерировать новые значения, а затем используйте операцию Put, чтобы записать новое значение в ByteKV. .
В некоторых сценариях значения некоторых полей могут быть автоматически сгенерированы в ByteSQL (например, автоматические первичные ключи и поля времени с атрибутами DEFAULT/ON UPDATE CURRENT_TIMESTAMP) или рассчитаны на основе зависимостей (например, SET a = a+1) , пользователю необходимо получить фактические измененные данные сразу после операции «Вставить», «Обновить» или «Удалить», а также выполнить операцию «Выбор» после записи. Всего требуется две операции Get и одна операция Put. Чтобы оптимизировать эффективность выполнения, ByteSQL реализует семантику возврата в синтаксисе PostgreSQL/Oracle: возвращает новое значение Insert/Update или старое значение Delete в том же запросе Query, экономя накладные расходы Get.
UPDATE table1 SET count = count + 1 WHERE id >= 10 RETURNING id, count;
Изменения онлайн-схемы
Непрерывное развитие и изменение бизнес-требований делают изменение схемы неизбежной задачей.Схема изменения схемы, встроенная в традиционные базы данных, обычно требует блокировки операций чтения и записи всей таблицы, что неприемлемо для онлайн-приложений. ByteSQL использует схему онлайн-изменения схемы Google F1 [3], которая не блокирует онлайн-запросы на чтение и запись в процессе изменения.
Метаданные схемы ByteSQL содержат определения библиотек и таблиц, которые хранятся в ByteKV. Экземпляры SQLProxy не имеют состояния, и каждый экземпляр периодически синхронизирует схему из ByteKV в локальную для анализа и выполнения запросов Query. В то же время в кластере есть специальный экземпляр Schema Change Worker, отвечающий за мониторинг и выполнение задач изменения схемы, отправленных пользователями. Как только рабочий процесс изменения схемы прослушивает запрос на изменение схемы, отправленный пользователем, он помещает его в очередь запросов и выполняет по порядку. В этом разделе объясняются причины введения промежуточного состояния схемы с точки зрения генерации и разрешения исключений согласованности данных. Подробное подтверждение правильности можно найти в оригинальной статье.
Поскольку разные экземпляры SQLProxy загружают схему в разное время, существует высокая вероятность того, что несколько версий схемы используются во всем кластере одновременно. Если процесс изменения схемы не обрабатывается должным образом, это приведет к несогласованности данных в таблице. Взяв в качестве примера создание вторичного индекса, рассмотрим следующий поток выполнения:
- Рабочий процесс изменения схемы выполняет задачу изменения создания индекса, включая заполнение ByteKV записями индекса и запись метаданных.
- Экземпляр SQLProxy 1 загружает метаданные схемы, содержащие новый индекс.
- Экземпляр SQLProxy 2 выполняет запрос на вставку. Поскольку экземпляр 2 не загрузил метаданные индекса, операция вставки не включает запись новых записей индекса.
- Экземпляр SQLProxy 2 выполняет запрос на удаление. Поскольку экземпляр 2 не загрузил метаданные индекса, операция Удалить не включает удаление новой записи индекса.
- Экземпляр SQLProxy 2 загружает метаданные схемы, содержащие новый индекс.
Оба шага 3 и 4 приводят к исключениям несоответствия между вторичным индексом и данными индекса первичного ключа: шаг 3 вызывает удаление записей вторичного индекса (потерянная запись), а шаг 4 вызывает устаревшие записи вторичного индекса (потерянное удаление). Причина этих исключений в том, что разные экземпляры SQLProxy загружают схему в разное время, в результате чего одни экземпляры думают, что индекс уже существует, а другие считают, что индекса не существует. В частности, исключение вставки на шаге 2 связано с тем, что индекс уже существует, но средство записи считает, что его не существует; исключение удаления на шаге 3 связано с тем, что средство записи воспринимает существование индекса, а средство удаления — нет. На самом деле операция обновления может вызвать оба вышеуказанных исключения.
Чтобы разрешить исключение «Потерянная запись», нам нужно убедиться, что для каждой вставленной строки данных экземпляр записи должен определить наличие индекса перед записью; а для исключения «Потерянное удаление» необходимо убедиться, что удаление экземпляр той же строки данных записывается до экземпляра записи. Воспринимает существование индекса (если экземпляр записывается, чтобы сначала воспринять индекс, а затем экземпляр удаляется, индекс может быть пропущен во время удаления, что приводит к потере Удалить).
Однако мы не можем напрямую контролировать воспринимаемый порядок различных экземпляров SQLProxy как экземпляров записи и экземпляров удаления, и вместо этого используем косвенный метод: определяем два промежуточных состояния для схемы для управления чтением и записью: состояние DeleteOnly и состояние WriteOnly, рабочий процесс изменения схемы. метаданные схемы состояния DeleteOnly, а затем запишите метаданные схемы состояния WriteOnly после синхронизации метаданных со всеми экземплярами. Экземпляры, воспринимающие состояние DeleteOnly, могут только удалять записи индекса, но не могут записывать записи индекса; экземпляры, воспринимающие состояние WriteOnly, могут как удалять, так и вставлять записи индекса. Это решает исключение Lost Delete.
Для исключения Lost Write мы не можем предотвратить запись данных экземплярами, которые еще не определили состояние Schema WriteOnly (поскольку весь процесс изменения схемы находится в режиме онлайн), но откладываем процесс заполнения записей индекса (называемый операцией Reorg в исходной статье). ) до тех пор, пока он не будет выполняться после этапа WriteOnly, таким образом заполняя записи индекса, соответствующие существующим данным в таблице, и заполняя записи индекса, отсутствующие из-за исключения Lost Write. После завершения операции заполнения метаданные схемы могут быть обновлены до общедоступного состояния, видимого для внешнего мира.
Мы решили проблему несогласованности данных при изменении схемы, введя два промежуточных состояния. Эти два промежуточных состояния являются внутренними для ByteSQL, и пользователь может видеть только индекс конечного общедоступного состояния. Здесь также возникает ключевой вопрос: как обеспечить синхронизацию состояния схемы во всех экземплярах SQLProxy в среде без глобальной информации об участниках? Решение состоит в том, чтобы поддерживать глобально фиксированное время аренды в схеме, и каждый SQLProxy должен перезагрузить схему из ByteKV, чтобы обновить контракт до истечения времени аренды. После того, как рабочий процесс изменения схемы каждый раз обновляет схему, ему необходимо дождаться успешной загрузки всех SQLProxy перед следующим обновлением. Для этого необходимо убедиться, что интервал между двумя обновлениями схемы превышает определенное время. Что касается того, насколько безопасен интервал, заинтересованные читатели могут подробно прочитать оригинальную статью [3], чтобы получить ответ. Если SQLProxy по какой-либо причине не может загрузить схему в течение срока аренды, текущий экземпляр ByteSQL устанавливается в недоступное состояние, и запросы на чтение и запись больше не обрабатываются.
будущее обсуждение
больше уровней согласованности
В сценарии межмашинного развертывания всегда есть запросы, требующие получения временных меток транзакций в машинном зале, что увеличивает задержку ответа; в то же время сетевая среда в машинном зале не так стабильна, как внутренняя компьютерная комната, а стабильность сети между машинными комнатами напрямую влияет на стабильность кластера. На самом деле, некоторые бизнес-сценарии не требуют надежных гарантий согласованности. В этих сценариях мы рассматриваем возможность введения гибридных логических часов HLC[4] для замены оригинальной глобальной службы синхронизации и преобразования ByteKV в систему, поддерживающую причинно-следственную согласованность. В то же время мы можем вернуть записанную метку времени в качестве пароля синхронизации клиенту, и клиент будет нести пароль синхронизации в последующих запросах для передачи события «произошло до» между событиями, которые имеют причинно-следственную связь в бизнесе, но не являются отношения, распознаваемые системой хранения, то есть непротиворечивость сеанса.
Кроме того, некоторые предприятия чрезвычайно чувствительны к задержкам и нуждаются в доступе к нескольким центрам обработки данных, а сценарии развертывания ByteKV с несколькими машинными помещениями не могут избежать задержек между машинными помещениями. Если этой части бизнеса нужно только поддерживать окончательную согласованность между компьютерными классами, мы можем синхронизировать данные в компьютерных залах, чтобы добиться эффекта окончательной согласованности класса.
Cloud Native
Поскольку CloudNative продолжает расти, он оказывает беспрецедентное влияние на существующую модель разработки и развертывания. ByteKV также продолжит изучение углубленной интеграции с CloudNative. Узнайте об автоматическом развертывании, автоматическом масштабировании и автоматическом восстановлении на основе Kubernetes. Дальнейшее улучшение использования ресурсов, снижение стоимости эксплуатации и обслуживания и повышение удобства использования услуг. Предоставьте ByteKV, более удобный для использования пользователями CloudNative.
использованная литература
[1] Ongaro D, Ousterhout J. In search of an understandable consensus algorithm (extended version)[J]. 2013.
[2] https://github.com/facebook/mysql-5.6/wiki/MyRocks-record-format#memcomparable-format
[3] Rae I, Rollins E, Shute J, et al. Online, asynchronous schema change in F1[J]. Proceedings of the VLDB Endowment, 2013, 6(11): 1045-1056.
[4] Kulkarni S, Demirbas M, Madeppa D, et al. Logical physical clocks and consistent snapshots in globally distributed databases[C]//The 18th International Conference on Principles of Distributed Systems. 2014.
[5] Peng D, Dabek F. Large-scale incremental processing using distributed transactions and notifications[J]. 2010.
[6] https://www.co ckroachlabs.com/blog/how-co ckroachdb-distributes-atomic-transactions/
[7] https://www.co ckroachlabs.com/blog/serializable-lockless-distributed-isolation-co ckroachdb/
[8] https://www.co ckroachlabs.com/blog/living-without-atomic-clocks/
[9] Shute J, Oancea M, Ellner S, et al. F1-the fault-tolerant distributed rdbms supporting google's ad business[J]. 2012.
[10] Корбетт Дж. К., Дин Дж., Эпштейн М. и др. Спаннер: глобально распределенная база данных Google[J], Транзакции ACM в компьютерных системах (TOCS), 2013, 31(3): 1-22.
[11] Ongaro D. Consensus: Bridging theory and practice[D]. Stanford University, 2014.
[12] Roohitavaf M, Ahn J S, Kang W H, et al. Session guarantees with raft and hybrid logical clocks[C]//Proceedings of the 20th International Conference on Distributed Computing and Networking. 2019: 100-109.
[13] Huang G, Cheng X, Wang J, et al. X-Engine: An optimized storage engine for large-scale E-commerce transaction processing[C]//Proceedings of the 2019 International Conference on Management of Data. 2019: 651-665.
поделиться больше
Чтобы получить адрес воспроизведения и ссылку на ppt последнего технического салона «Крупномасштабный проект гибридного развертывания в ByteDance», ответьте на фоне общедоступной учетной записи WeChat «Техническая команда ByteDance»: Салон архитектурных технологий.
Команда инфраструктуры ByteDance
Команда инфраструктуры ByteDance — важная команда, которая поддерживает бесперебойную работу множества пользовательских продуктов ByteDance, включая Douyin, Today's Toutiao, Xigua Video и Volcano Small Video, Стабильная разработка обеспечивает гарантии и импульс.
В компании команда инфраструктуры в основном отвечает за создание частного облака ByteDance, управление десятками тысяч кластеров серверного масштаба, отвечает за десятки тысяч гибридных развертываний вычислений/хранилищ и гибридных развертываний онлайн/офлайн, а также за поддержку стабильного хранилища. нескольких массивных данных ЭП.
В культурном плане команда активно использует открытый исходный код и инновационные аппаратные и программные архитектуры. Мы давно занимаемся набором студентов по направлению инфраструктура.Подробности на сайте job.bytedance.com.Если интересно, можете написать на guoxinyu.0372@bytedance.com.
Добро пожаловать в техническую команду ByteDance