Эволюция системы хранения распределенных таблиц ByteDance

Архитектура
Эволюция системы хранения распределенных таблиц ByteDance

Эта статья выбрана из серии статей «Практика инфраструктуры Byte Beat».

Серия статей «Практика инфраструктуры ByteDance» представляет собой техническую галантерею, созданную техническими командами и экспертами отдела инфраструктуры ByteDance, в которой мы делимся с вами практическим опытом и уроками команды в процессе развития и эволюции инфраструктуры, а также всеми техническими студентами. общаться и расти вместе.

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

резюме

Распределенные системы хранения таблиц имеют широкий спектр сценариев применения в отрасли. Google выпустила два поколения распределенных систем хранения таблиц, Bigtable и Spanner, чтобы удовлетворить все потребности в хранении таблиц во внутренних и внешних облачных сервисах. Среди них реализация Bigtable с открытым исходным кодом, HBase, широко используемая отечественными и зарубежными компаниями, а также базовые слои графовой базы данных с открытым исходным кодом JanusGraph, базы данных временных рядов OpenTSDB, базы данных географической информации GeoMesa, реляционная база данных Phoenix. на HBase для хранения данных. ByteDance также использует HBase в качестве службы хранения таблиц.В практическом процессе ByteDance HBase имеет проблемы с высокой задержкой, длинным хвостом и недостаточной доступностью в сценарии чрезвычайно большого объема данных и чрезвычайно высокой нагрузки. С этой целью группа хранения таблиц ByteDance провела много исследований систем с открытым исходным кодом и, наконец, решила разработать набор с высокой доступностью, строгой согласованностью, малой задержкой, высокой пропускной способностью, глобальным порядком и совместимостью как с SSD, так и с жесткими дисками. которые совместимы с моделью данных и семантикой HBase Распределенная система хранения таблиц для решения технических проблем, вызванных быстрым развитием бизнеса ByteDance. В настоящее время первая репрезентативная грид-система хранения Bytable1.0 стабильно работает в качестве базового хранилища данных для поиска и рекомендаций, и наша команда разрабатывает на этой основе вторую репрезентативную грид-систему хранения Bytable2.0.

1. Предпосылки

С запуском и развитием проекта Toutiao по поиску всей сети бизнесу нужна система хранения таблиц с глобальным порядком, огромной емкостью и высокой производительностью для хранения всех ссылок и веб-страниц во всем Интернете, а также для обеспечения того, чтобы все изменения в Интернете могли быть записаны в режиме реального времени обновление в системе хранения таблиц. До этого Google, лидер в области поисковых технологий, создал для этой цели Bigtable и использовал его в нескольких своих проектах, включая Search, Earth и Finance. Среди них поисковый бизнес в основном обсуждался в опубликованной газете, поэтому многие компании, занимающиеся поисковым бизнесом, последовали их примеру. Наша команда первоначально обслуживала поисковый бизнес, используя HBase, реализацию Bigtable с открытым исходным кодом. В соответствии с требованиями обновления в реальном времени для всей связи сетевого канала необходимо обеспечить достаточно высокую доступность и достаточно низкую задержку.Из-за чрезвычайно большого объема данных будет создано большое количество фрагментов данных, а общий хвост задержка и доступность кластера будут увеличиваться с увеличением количества экземпляров сегментирования данных, что приводит к экспоненциальному ухудшению состояния, поэтому к задержке и доступности каждого экземпляра сегментирования предъявляются более высокие требования. Однако из-за проблем с высокой задержкой хвоста и низкой доступностью HBase не мог удовлетворить наши потребности, поэтому наша команда начала исследования и выбор технологии.

2. Bytable1.0 стала первой репрезентативной грид-системой хранения.

2.1 Технический отбор

По нашему требованию, сталкиваясь с массивными веб-ссылками и веб-данными, объем данных очень велик, и любое использование SSD-дисков значительно увеличит стоимость хранения, поэтому мы можем не только рассмотреть SSD-диски, но и предоставить HDD-диски.Эффективная поддержка . В настоящее время широко используемые глобально упорядоченные распределенные системы хранения таблиц с открытым исходным кодом включают HBase, TiDB, CockroachDB и т. д. Было подтверждено, что HBase не может удовлетворить наши потребности из-за длинного хвоста проблем с задержкой и доступностью. TiDB и CockroachDB существуют под наши нужды: (1) требуется многоканальное слияние и сканирование данных при переносе данных; (2) при записи данных требуются два лога; эти два вызывают недружественные проблемы с поворотом головы на HDD-диски. Поскольку ни одна из этих систем не может решить наши проблемы, наша команда, наконец, решила самостоятельно разработать набор систем высокой доступности, строгой согласованности, низкой задержки и высокой пропускной способности, совместимых с моделью данных и семантикой HBase, на основе текущего уровня программных и аппаратных ресурсов. и бизнес-сценариев, предоставляемых компанией. , Распределенная система хранения таблиц Bytable1.0, которая глобально упорядочена и совместима как с SSD, так и с HDD дисками. При разработке схемы принцип, которого придерживается наша команда, заключается в том, чтобы максимально упростить дизайн с учетом потребностей бизнеса. Распределенные системы хранения таблиц в отрасли часто используют общую архитектуру хранения, но распределенная система хранения файлов внутри компании на тот момент не могла обеспечить достаточно высокую доступность и SLA. Чтобы удовлетворить требования бизнеса к высокой доступности, мы решили напрямую управлять локальными дисками, не полагаясь на распределенную систему хранения файлов.

2.2 Архитектура системы

Давайте представим системную архитектуру Bytable1.0, как показано на рисунке ниже, она в основном состоит из трех модулей: Master, PlacementDriver и TabletServer. Клиент свяжется с Мастером, чтобы получить метаинформацию планшета (аналогично Bigtable, используя планшет для представления фрагмента данных в таблице, представляющего часть данных, соответствующих KeyRange в таблице), чтобы получить адрес TabletServer, где находится соответствующий TabletServer, а затем клиент связывается с соответствующим TabletServer для чтения и записи данных. Подобно NameNode и DataNode в HDFS, которые имеют дело с задачами плоскости управления и плоскости данных соответственно, мы делим весь Bytable1.0 на плоскость управления (главный), плоскость принятия решений (драйвер размещения) и плоскость данных. (TabletServer).Мы подробно познакомим вас с каждым из них.Плоский дизайн интерьера.

Bytable1.0 Global架构图
Диаграмма глобальной архитектуры Bytable1.0

2.2.1 Мастер плоскости управления

Мастер управляет определенным процессом активного переключения планшета, выбора пассивного мастера, миграции, разделения и слияния.Кроме того, Мастер предоставляет информацию о местонахождении планшета Клиенту. Здесь общий процесс описан на примере переключения пассивного мастера.Если Мастер не получает сообщение пульса от определенного TabletServer в течение заданного времени, TabletServer будет определен как Offline, а затем Tablet, который является лидером на TabletServer инициирует раунд выбора мастера. Процесс выбора мастера аналогичен Рафту, за исключением того, что эта часть логически перемещается в Мастер для унифицированного управления.Сначала собираются последние лог-номера каждой копии Планшета, а затем достаточно новая копия выбирается в качестве главной цели выбора, и голосование выдается каждой копии Планшета.При голосовании После успеха отправьте основной запрос на выбранную основную цель. Поскольку мастер берет на себя всю работу уровня управления, логика и взаимодействие RPC будут очень сложными, а поскольку нет особенно высоких требований к производительности, язык Go выбран для разработки, чтобы снизить затраты на НИОКР и повысить эффективность НИОКР.

2.2.2 Драйвер размещения плоскости принятия решений

Мастер контролирует конкретный процесс указанной операции, но не принимает решения о том, следует ли инициировать указанную операцию (за исключением пассивного выбора мастера), а задачу этого решения берет на себя PlacementDriver. Он сканирует Мастер, чтобы получить информацию о загрузке раздела, генерирует информацию о решении путем расчета различных стратегий и отправляет ее Мастеру для фактического выполнения. Поскольку обновление стратегии чаще происходит чаще, если стратегия разбита на два модуля, обновление стратегии не требует непрерывного перезапуска Мастера, который может быть совершенно незаметен для пользователя. Чтобы удешевить написание сложных стратегий, PD также выбрала для разработки язык Go.

2.2.3 Плоскость данных TabletServer

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

Путь данных является критическим путем этой системы, которая требует более высокой производительности для повышения пропускной способности и уменьшения задержки.На уровне языка мы выбрали C++, чтобы обеспечить эффективную и предсказуемую производительность.Кроме того, мы также выполнили механизм ведения журналов и механизм обработки данных. Много оптимизаций. В настоящее время большинство систем на рынке, таких как MySQL, MongoDB, TiDB, CockroachDB и т. д., при записи журналов записывают два журнала: журнал репликации и журнал двигателя.Трафик удваивается, и голова часто качается, ограничивая производительность записи. Поэтому в нашей реализации журнал движка избегается и используется только журнал репликации, чтобы избежать проблемы качания головы. Кроме того, в решениях TiDB и CockroachDB RocksDB используется в качестве комбинированного решения для хранения нескольких наборов реплицированных журналов, но это приведет к ненужной вставке MemTable и накладным расходам SST Flush & Compaction, что недостаточно эффективно. Поэтому наша команда разработала механизм хранения WAL для объединения нескольких наборов реплицированных журналов и реализации записи журнала только один раз, не вызывая раскачивания головки жесткого диска, а полоса пропускания записи жесткого диска может быть заполнена даже при сжатии. не выполняется. .

Что касается механизма хранения данных, наша команда использует RocksDB, который популярен и широко проверен в отрасли, в качестве механизма хранения. В отличие от TiDB и CockroachDB, каждый планшет Bytable1.0 соответствует экземпляру RocksDB.Во время процесса переноса данных чтение, передача и запись файлов RocksDB могут выполняться напрямую без дополнительного сканирования и создания файлов, а также внедрения файлов, избегая движения головы. проблема, вызванная слиянием нескольких файлов во время сканирования, потенциальная проблема с остановкой записи, вызванная ошибочным входом операции внедрения файла RocksDB в L0, и операция внедрения файла RocksDB не разрешает проблему одновременной обычной записи.Легкая миграция данных. Кроме того, размер LSMTree можно разумно контролировать, чтобы избежать проблемы, связанной с тем, что размер LSMTree не соответствует параметрам RocksDB.

Эффективная реализация операций разделения и слияния — это проблема, возникающая из-за того, что планшет соответствует экземпляру RocksDB. Для операции Split необходимо разделить исходный экземпляр RocksDB на два экземпляра RocksDB.В нашей реализации мы используем характеристики LSMTree.Сначала мы выполним раунд жесткой цепочки полного файла движка, затем перестанем писать, а затем выполнить еще один раунд. Файлы инкрементного движка жестко привязаны и открыты для записи, а время недоступности — это время, необходимое для небольшого количества инкрементных операций жесткой цепочки. Для операции слияния мы модифицировали RocksDB, чтобы во время внедрения файлов в манифест сохранялись введенный файл и GlobalSequence, что позволяет напрямую вводить файлы в другой экземпляр RocksDB. Подобно процессу Split, процесс Merge также использует добавочные стратегии для внедрения, чтобы сократить время простоя.

По сравнению с реализацией нескольких планшетов, совместно использующих экземпляр RocksDB, решение Bytable 1.0 имеет больше преимуществ в стабильности и эффективности миграции и контроле размера одного LSMTree, но есть определенные недостатки в дрожании, вызванном разделением и слиянием.

На следующем рисунке показано сравнение задержек между Bytable1.0 TabletServer и HBase RegionServer в одной и той же среде компьютера в смешанном сценарии чтения-записи (50% чтения, 50% записи). Видно, что Bytable1.0 значительно уменьшил среднюю задержку и задержку p99.

2.3 Детали оптимизации

2.3.1 Оптимизация считывания точек

Поскольку мы инкапсулировали модель данных Table в модель KV RocksDB, невозможно использовать ее встроенный фильтр BloomFilter для оптимизации производительности чтения точек.Мы вручную генерируем фильтр BloomFilter для каждого SST в процессе Flush и Compaction и сохраняем его в свойство TableProperty и отключил встроенный фильтр BloomFilter для экономии места. При запросе используйте функцию TableFilter, чтобы найти созданный вручную BloomFilter, соответствующий SST, для оценки, чтобы оптимизировать производительность чтения точки.

2.3.2 Оптимизация записи точки доступа

Поскольку мы используем Raft в качестве протокола репликации реплики, а процесс Apply в Raft является последовательным, когда возникает точка доступа для записи, нам необходимо разделить планшет точки доступа для повышения производительности. Поскольку наши операции разделения и слияния относительно тяжелы, решение о разделении точек доступа не подходит для слишком реального времени и чувствительности.Чтобы увеличить производительность записи одного планшета, мы внедрили новый тип RocksDB MemTable, который реализуется только в очереди записи при записи.При сложности онлайн-записи O(1) данные одновременно записываются в фактический SkipList пулом фоновых потоков, а механизм ReadIndex используется для ожидания при чтении. Увеличьте максимальную грузоподъемность одного потока до максимальной грузоподъемности одной машины.

2.3.3 Запись оптимизации противодавления

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

2.4 Распределенные транзакции

В настоящее время в некоторых бизнес-сценариях существует потребность в распределенных транзакциях.В соответствии с моделью транзакций Percolator наша команда реализовала набор решений для распределенных транзакций на верхнем уровне Bytable.Соответствующее содержание будет подробно описано в последующих специальные статьи.

3. Bytable2.0 разрабатывает идеальную грид-систему хранения второго порядка

С запуском и стабильностью Bytable 1.0 бизнес предъявляет более высокие требования к стоимости, стабильности и производительности системы хранения таблиц. С точки зрения затрат, в настоящее время использование трех копий для хранения данных требует почти в два раза больше места для хранения, чем решение EC (Erasure Code) для распределенной файловой системы, в разы больше ЦП. С точки зрения использования ресурсов в настоящее время существует несоответствие между запросами на чтение и запись и занятостью дискового пространства, и будут ситуации, когда ресурсы ЦП остаются, но дисковое пространство исчерпано, или ресурсы ЦП исчерпаны, но дисковое пространство остается. С точки зрения подключения к автономной экосистеме Hadoop создание резервных копий, экспорт и загрузка данных теперь немного хлопотны. С точки зрения стабильности, в процессе разделения и слияния все еще существует определенное дрожание задержки, и мы также надеемся максимально оптимизировать задержку хвоста. С точки зрения сложности эксплуатации и обслуживания, основанной на распределенной файловой системе, наши сервисы могут быть сделаны без сохранения состояния.В сочетании с технологией планирования K8S нам больше не нужно рассматривать отчет о сбоях и перезапускать машины и диски, которые будут значительно сократить наши расходы на эксплуатацию и техническое обслуживание. Поэтому необходимо спроектировать и внедрить распределенную систему хранения таблиц на основе общего хранилища.

3.1 Технический отбор

В отличие от цели разработки Bytable 1.0, которая заключается только в удовлетворении текущих потребностей бизнеса, мы надеемся найти оптимальное решение в более широком масштабе и создать ведущую в мире распределенную систему хранения таблиц. Поскольку распределенная система хранения файлов ByteDance постепенно переходит в онлайн-режим, у нас есть распределенная система хранения файлов ByteStore, которую можно использовать онлайн. Кроме того, ByteDance также разработала ByteJournal, распределенную систему хранения журналов, использующую протокол Quorum, которая имеет функции неупорядоченной репликации и отправки журналов, аналогичные Paxos, для максимального сокращения задержки. Кроме того, ByteJournal также может обеспечить в общей сложности 5 копий записи и 2 копии могут быть успешно отправлены, но требует расширенных функций 4 копий Recover для обеспечения более экстремальной оптимизации хвостовой задержки, чем протокол Paxos. На этот раз мы решили встать на плечи гигантов, использовать ByteStore для хранения файлов движка, использовать ByteJournal для управления журналами репликации и полагаться на эти две распределенные системы хранения для реализации Bytable2.0, грид-системы хранения второго поколения, основанной на общее хранилище. Кроме того, у международных компаний также есть глобально согласованные требования к хранению таблиц, и они надеются установить и изменить конфигурацию репликации некоторых данных на мелкозернистой основе, тем самым контролируя задержку чтения и записи соответствующих данных во время доступа (через расстояние, когда доступ к данным).

3.2 Архитектура системы

В системном дизайне Bytable2.0 для справки используется идея дизайна Spanner, распределенной системы хранения таблиц Google второго поколения. В этом разделе мы представим системную архитектуру Bytable2.0.Как показано на рисунке ниже, мы разделяем Bytable2.0 на уровень чтения/записи данных (TabletServer), уровень глобального управления (GlobalMaster), уровень планирования вычислений (GroupMaster) и слой планирования данных (PlacementDriver).

Основные различия между Bytable2.0 и Bytable1.0 заключаются в следующем:

  1. В отличие от Bytable1.0, который хранит информацию о распространении Tablet в Master, Bytable2.0 создаст соответствующую MetaTable для каждой таблицы и сохранит ее в TabletServer. Одним из преимуществ этого является то, что он может изолировать доступ к информации о распространении Tablet для каждой таблицы.Другое преимущество заключается в том, что если MetaTable становится узким местом, RootTable можно добавить в более поздних реализациях, чтобы сделать ее разделяемой.
  2. В отличие от Bytable1.0, где Tablet соответствует группе репликации, группа репликации в Bytable2.0 будет содержать несколько Tablet. Окончательная производительность заключается в том, что TabletServer имеет реплику нескольких групп репликации (реплика — это копия группы репликации), а группа репликации имеет несколько планшетов. Преимущество этого заключается в том, что он позволяет помещать данные нескольких прерывистых диапазонов ключей в одну и ту же группу репликации для облегчения чрезвычайно тонкого разделения и позволяет обрабатывать данные прерывистых диапазонов ключей локально в одной группе репликации. Разделение возможности репликации ReplicationGroup и разделения KeyRange может поддерживать конвергентность числа ReplicationGroups даже при наличии множества разделений KeyRange.
  3. В отличие от Bytable1.0, у которого был только один мастер, Bytable2.0 представил два компонента: GlobalMaster и GroupMaster. GlobalMaster отвечает за запись метаинформации всех таблиц и групп репликации. GroupMaster постоянно синхронизирует метаинформацию, управляемую GlobalMaster, и сохраняет ее локально с избыточностью.Он управляет только активацией и размещением реплик экземпляров TabletServer в соответствующем регионе, что реализует автономию внутри региона и устраняет влияние состояния сети. между регионами на внутреннюю доступность региона.

Основные различия между Bytable2.0 и HBase заключаются в следующем:

  1. HBase поддерживает только синхронную или асинхронную репликацию master-slave, а при чтении гарантирует только внутрикластерную согласованность, но не межкластерную. Bytable 2.0 может поддерживать репликацию данных Quorum между регионами для обеспечения глобальной согласованности.
  2. Отношение сопоставления между Region и KeyRange в HBase является сопоставлением «один к одному», поэтому, если вы хотите установить отношение репликации KeyRange на мелкозернистой основе, будет создано большое количество небольших регионов, что не является дружественным. к HBase. А Bytable 2.0 может поддерживать детализированные настройки отношений репликации данных Tablet.

Когда клиент читает и записывает данные, он сначала связывается с GroupMaster, чтобы получить информацию о местоположении каждой реплики ReplicationGroup, соответствующей MetaTable соответствующей таблицы данных, а затем запрашивает соответствующую MetaTable, чтобы найти ReplicationGroupId соответствующего Tablet, а затем переходит к GroupMaster для запроса информации о местоположении соответствующей ReplicationGroup.Replica и, наконец, находит соответствующий TabletServer для чтения и записи на основе этой информации о местоположении.

Bytable2.0 Global架构图
Диаграмма глобальной архитектуры Bytable2.0

В сценарии развертывания в одной комнате приведенная выше архитектура слишком раздута.Поскольку мы сделали достаточно хорошую реализацию модуляризации, мы можем скомпилировать три модуля GlobalMaster, GroupMaster и PlacementDriver в один и тот же двоичный файл Master, чтобы сократить развертывание и обслуживание. сложности, как показано на рисунке ниже, кластер может иметь только две службы, Master и TabletServer, во время развертывания.

Bytable2.0 精简部署
Оптимизированное развертывание Bytable2.0

3.2.1 Уровень чтения и записи данных TabletServer

В отличие от однозначного соответствия между группами репликации и планшетами в Bytable1.0, TabletServer в Bytable2.0 содержит реплики нескольких групп репликации, и каждая группа репликации содержит несколько планшетов. Этот дизайн может сделать разделение и слияние планшета очень легким, просто измените метаданные без паузы. В соответствии с этой схемой миграция больше не является вопросом переноса планшета с одного TabletServer на другой TabletServer, а заключается в импорте соответствующих данных Tablet из ReplicationGroup в другую ReplicationGroup.Мы используем метод извлечения диапазона файлов SST, чтобы легко завершить соответствие. Работа.

В процессе записи Leader Replica сначала записывает в WAL через ByteJournal и записывает в MemTable, когда журнал успешно записан. Follower Replica читает WAL из ByteJournal и записывает в MemTable. Реплика лидера регулярно выполняет операции Checkpoint и Compaction для записи файлов SST в ByteStore. Когда ведущая реплика и ведомая реплика совместно используют распределенную файловую систему, только ведущая реплика выполняет операции проверки и уплотнения. Если ведущая реплика и подчиненная реплика не используют общую распределенную файловую систему, в каждой распределенной файловой системе будет выбран лидер сжатия для выполнения операций контрольной точки и сжатия соответственно. Таким образом, мы можем решить проблему дополнительных потерь ЦП, вызванных предыдущим уплотнением нескольких копий.В то же время мы также можем использовать реплику Follwer для решения проблемы отсутствия избыточности master-slave в более ранних версиях HBase.

В процессе чтения и ведущая реплика, и ведомая реплика считывают данные SST из соответствующего ByteStore, объединяют их с MemTable в памяти и, наконец, возвращают пользователю.Вспомогательная реплика может указать версию для чтения или используйте метод ReadIndex.Получите последний номер журнала из ведущей реплики и подождите, пока он будет прочитан локально после применения.

3.2.2 Глобальное управление GlobalMaster

GlobalMaster — это глобальный модуль управления всем кластером, который отвечает за доступ к рабочему интерфейсу для обслуживающего и обслуживающего персонала, а также за создание и удаление таблиц и ColumnFamily. Он хранит метаинформацию обо всех таблицах и ColumnFamily, а также сохраняет метаинформацию обо всех ReplicationGroups, содержащихся в каждой таблице. GroupMaster синхронизирует эту метаинформацию с GlobalMaster, чтобы получить реплику, которая нуждается в планировании для управления размещением реплики ReplicationGroup в соответствующем регионе.

3.2.3 Уровень планирования вычислений GroupMaster

GroupMaster — это модуль управления в регионе, который управляет несколькими экземплярами TabletServer в регионе и репликами, выделенными каждой ReplicationGroup в регионе. В дополнение к пульсу и работе поддержания активности нескольких TabletServers, зарегистрированных в этом GroupMaster, он также получает информацию о репликах, которая должна быть запланирована, из метаинформации, синхронизированной GlobalMaster, и планирует ее на нескольких TabletServers, зарегистрированных в этом GroupMaster. Укажите TabletServer открыть или закрыть соответствующую реплику.

3.2.4 Уровень планирования данных PlacementDriver

PlacementDriver – это модуль принятия решений для планирования данных. Он связывается с TabletServer, собирает информацию о нагрузке каждого TabletServer, сканирует MetaTablet, соответствующий каждой таблице, выполняет расчет балансировки нагрузки и, наконец, решает, следует ли разделить, объединить или перенести определенные планшеты. работать.

3.3 Перспективы будущего

3.3.1 Автономное уплотнение

По сравнению с чтением и записью потоковых данных, операция сжатия на самом деле является автономной операцией пакетного типа, и это автономное задание выполняется в онлайн-службе, что часто снижает эффективное использование ресурсов и качество обслуживания онлайн-службы. Поэтому мы планируем передать выполнение этого процесса уплотнения Yarn или другим автономным платформам, выполняющим задания, чтобы снизить потребление ЦП для онлайн-сервисов. Решение FPGA, используемое командой Aliyun X-Engine, использует аналогичный метод для разгрузки нагрузки ЦП локального уплотнения на FPGA. Однако он в значительной степени зависит от конкретных моделей и аппаратного обеспечения.Вы можете рассмотреть возможность использования платформы планирования Yarn для более разумного планирования ресурсов FPGA и полного использования его ресурсов.

3.3.2 Механизм хранения аналитических колонок

Как упоминалось в документе Spanner: Становление системой SQL, Spanner в настоящее время предоставляет набор аналитических механизмов с хранилищем столбцов в блоках, что дало хорошие результаты в некоторых сценариях. TiDB также предоставляет решение для аналитического запроса, подвешивая подчиненный узел ClickHouse (аналитические данные с открытым исходным кодом Яндекса) для Raft. Bytable 2.0 также имеет несколько аналитических сценариев использования.Мы планируем предоставить набор механизмов хранения аналитических столбцов в будущем и постепенно превратить его в систему HTAP.

3.3.3 Мультимодальная база данных

В настоящее время на рынке баз данных существуют различные API-интерфейсы, такие как Redis, MongoDB, MySQL и т. д. В то же время на рынке с открытым исходным кодом нижние уровни баз данных, таких как графовая база данных JanusGraph, база данных временных рядов OpenTSDB, географическая информационная база данных GeoMesa и реляционная база данных Phoenix, все основаны на HBase для хранения данных. Это вдохновило нас на то, что в будущем Bytable 2.0 можно будет использовать в качестве общей базовой службы хранения, реализуя различные уровни API на верхнем уровне и предоставляя недорогие решения для этих служб API с большой емкостью.

3.3.4 Исследование нового механизма распределенных транзакций

При реализации распределенных транзакций в отрасли часто приходится искать компромисс между пропускной способностью и задержкой. В решении Bytable для распределенных транзакций используется модель транзакций Percolator, которая опирается на одноточечный сервис синхронизации TSO и обеспечивает определенный компромисс между пропускной способностью и задержкой. В последующей схеме распределенных транзакций Bytable2.0 мы надеемся дополнительно оптимизировать проблему пропускной способности и проблему задержки между регионами модели Percolator в схеме распределенных транзакций Bytable1.0 и изучить новые решения.

4. Резюме

Благодаря быстрому развитию сценариев поиска и рекомендаций ByteDance, а также постоянному расширению потребностей бизнеса мы быстро построили распределенную систему хранения таблиц первого поколения Bytable1.0 на ранней стадии, решив большое количество бизнес-задач. Но в то же время из-за упрощения и компромисса системы первого этапа все еще есть некоторые проблемы, такие как стоимость и масштабируемость, которые необходимо оптимизировать.Мы решаем эти проблемы в Spanner-подобной таблице второго поколения. система хранения Bytable2.0. Заинтересованные студенты могут присоединиться к команде по хранению таблиц, чтобы создать идеальное решение для распределенного хранения таблиц и вместе построить надежную инфраструктуру ByteDance, чтобы сопровождать быстрое развитие компании.

5. Ссылки

[1] Chang, Fay, et al. "Bigtable: A distributed storage system for structured data." ACM Transactions on Computer Systems (TOCS) 26.2 (2008): 1-26.

[2] Apache HBase: An open-source, distributed, versioned, column-oriented store modeled after Google' Bigtable.

[3] TiDB: An open source distributed HTAP database compatible with the MySQL protocol.

[4] Taft, Rebecca, et al. "Co ckroachDB: The Resilient Geo-Distributed SQL Database." Proceedings of the 2020 ACM SIGMOD International Conference on Management of Data. 2020.

[5] Dong, Siying, et al. "Optimizing Space Amplification in RocksDB." CIDR. Vol. 3. 2017.

[6] Корбетт, Джеймс К. и др. «Spanner: глобально распределенная база данных Google.» ACM Transactions on Computer Systems (TOCS) 31.3 (2013): 1-22.

[7] Bacon, David F., et al. "Spanner: Becoming a SQL system." Proceedings of the 2017 ACM International Conference on Management of Data. 2017.

[8] Huang, Gui, et al. "X-Engine: An optimized storage engine for large-scale E-commerce transaction processing." Proceedings of the 2019 International Conference on Management of Data. 2019.

поделиться больше

ByteDance самостоятельно разработала надежную и согласованную онлайн-практику хранения KV и таблиц — часть 1

ByteDance самостоятельно разработала строгую согласованную онлайн-практику хранения KV и таблиц — часть 2

Эксклюзивное интервью InfoQ с Toutiao Search: от рекомендации к поиску, как построить еще одну возможность поисковой технологии?

Практика ByteDance в сетевой библиотеке Go

Команда инфраструктуры ByteDance

Команда инфраструктуры ByteDance — важная команда, которая поддерживает бесперебойную работу множества пользовательских продуктов ByteDance, включая Douyin, Today's Toutiao, Xigua Video и Volcano Small Video, Стабильная разработка обеспечивает гарантии и импульс.

В компании команда инфраструктуры в основном отвечает за создание частного облака ByteDance, управление десятками тысяч кластеров серверного масштаба, отвечает за десятки тысяч гибридных развертываний вычислений/хранилищ и гибридных развертываний онлайн/офлайн, а также за поддержку стабильного хранилища. нескольких массивных данных ЭП.

В культурном плане команда активно использует открытый исходный код и инновационные аппаратные и программные архитектуры. Мы давно набираем студентов по направлению инфраструктура.Подробности смотрите на job.bytedance.com ("Читать исходный текст" в конце статьи).Если вам интересно, вы можете связаться с guoxinyu .0372@bytedance.com.

Добро пожаловать в техническую команду ByteDance