Эта страница переведена машинным переводом. Читайте английский оригинал. English

Библиотека IBSurgeon

Параллельное чтение данных в Firebird

D.Simonov, V.Horsun

версия 1.0.5 от 05.12.2023

Этот материал спонсируется и создан при спонсорской поддержке IBSurgeon www.ib-aid.com, поставщика HQbird (расширенной дистрибуции Firebird) и поставщика услуг по оптимизации производительности, миграции и технической поддержке Firebird.

Материал распространяется по лицензии Public Documentation License https://www.firebirdsql.org/file/documentation/html/en/licenses/pdl/public-documentation-license.html

Связанные материалы:

Предисловие

Firebird 5.0 представил возможность использования параллелизма при создании резервной копии с помощью утилиты gbak, среди прочих параллельных функций. Изначально эта функция появилась в HQbird 2.5, затем в HQbird 3.0 и 4.0, а затем была перенесена в Firebird 5.0.

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

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

Пример утилиты FBCSVExport

Для демонстрации параллельного чтения данных из СУБД Firebird была написана примерная утилита, которая экспортирует данные из одной или нескольких таблиц в формат CSV.

Её описание и открытый исходный код находятся здесь: https://github.com/IBSurgeon/FBCSVExport.git

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

По любым вопросам обращайтесь по адресу [email protected].

Параллельное чтение

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

Создадим пул рабочих потоков. Основной поток приложения также является рабочим потоком, поэтому количество дополнительных рабочих потоков должно быть N - 1, где N - общее количество параллельных рабочих потоков. Каждый рабочий поток будет использовать своё собственное соединение и транзакцию.

Первая проблема: как обеспечить согласованность читаемых данных?

Согласованное чтение данных

Поскольку каждый рабочий поток использует своё собственное соединение и свою собственную транзакцию, возникает проблема несогласованного чтения - если таблица одновременно изменяется другими пользователями, то читаемые данные могут быть несогласованными. В однопоточном режиме gbak использует транзакцию с режимом изоляции SNAPSHOT, что позволяет читать согласованную информацию на момент начала SNAPSHOT-транзакции. Но здесь у нас несколько транзакций, и необходимо, чтобы они видели один и тот же “снимок”, чтобы читать одни и те же неизменяемые данные.

Механизм создания общего снимка для разных транзакций с режимом изоляции SNAPSHOT был введён в Firebird 4.0 (изначально в HQBird 2.5, но в Firebird 4.0/HQBird 4.0 он проще и эффективнее). Есть два способа создать общий снимок:

  1. С помощью SQL
  • получить номер снимка из основной транзакции (которая запущена в основном рабочем потоке).
sql
SELECT RDB$GET_CONTEXT('SYSTEM', 'SNAPSHOT_NUMBER') FROM RDB$DATABASE
  • запустить другие транзакции с помощью следующего SQL:
sql
SET TRANSACTION SNAPSHOT AT NUMBER snapshot_number

где snapshot_number - номер, полученный предыдущим запросом.

  1. С помощью API
  • получить номер снимка из основной транзакции (которая запущена в основном рабочем потоке) с помощью функции isc_transaction_info или ITransaction.getInfo с тегом fb_info_tra_snapshot_number;

  • запустить другие транзакции с тегом isc_tpb_at_snapshot_number с полученным номером снимка.

Примерная утилита FBCSVExport, а также gbak, использует второй подход. Эти подходы можно смешивать - например, получить номер снимка с помощью SQL, а затем использовать полученный номер снимка для запуска других транзакций с помощью API, или наоборот.

В FBCSVExport мы получаем номер снимка с помощью следующего кода:

cpp
ISC_INT64 getSnapshotNumber(Firebird::ThrowStatusWrapper* status, Firebird::ITransaction* tra)
{
    ISC_INT64 ret = 0;
    unsigned char in_buf[] = { fb_info_tra_snapshot_number, isc_info_end };
    unsigned char out_buf[16] = { 0 };

    tra->getInfo(status, sizeof(in_buf), in_buf, sizeof(out_buf), out_buf);

    unsigned char* p = out_buf, * e = out_buf + sizeof(out_buf);
    while (p < e)
    {
        short len = 0;
        switch (*p++)
        {
        case isc_info_error:
        case isc_info_end:
            p = e;
            break;

        case fb_info_tra_snapshot_number:
            len = static_cast(isc_vax_integer(reinterpret_cast<char*>(p), 2));
            p += 2;
            ret = isc_portable_integer(p, len);
            p += len;
            break;
        }
    }
    return ret;
}

Для запуска транзакции с номером снимка мы используем следующий код:

sql
Firebird::AutoDispose tpbWorkerBuilder(fbUtil->getXpbBuilder(&status, Firebird::IXpbBuilder::TPB, nullptr, 0));
tpbWorkerBuilder->insertTag(&status, isc_tpb_concurrency);
tpbWorkerBuilder->insertBigInt(&status, isc_tpb_at_snapshot_number, snapshotNumber);

Firebird::AutoRelease workerTra(
    workerAtt->startTransaction(
        &status,
        tpbWorkerBuilder->getBufferLength(&status),
        tpbWorkerBuilder->getBuffer(&status)
    )
);

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

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

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

Разделение большой таблицы на части

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

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

  • На сколько частей обработки следует разделить таблицу?

  • Как лучше всего разделить таблицу с точки зрения доступа к данным?

Ответим на эти вопросы по порядку.

На сколько частей обработки следует разделить таблицу?

Предположим идеальный сценарий - сервер и клиент выделены под Firebird, то есть все процессоры полностью в нашем распоряжении. Тогда рекомендуется:

а) Использовать в качестве максимального количества параллельных частей удвоенное количество ядер процессора на сервере. Почему 2x ядер? Мы точно знаем, что будут задержки, связанные с вводом-выводом, поэтому мы можем позволить некоторое дополнительное использование процессора. Однако это число следует рассматривать как начальную настройку, на практике оно зависит от данных.

б) Учитывать количество ядер на клиенте: если на сервере их значительно больше (обычная ситуация), то может иметь смысл дополнительно ограничить количество частей разбиения, чтобы не перегружать клиент (он всё равно не сможет обработать больше, а затраты на переключение потоков никуда не денутся). Более точно можно будет решить, наблюдая за загрузкой процессора на клиенте и сервере - если на клиенте она 100%, а на сервере заметно меньше, то имеет смысл уменьшить количество частей.

c) если клиент и сервер находятся на одном хосте, то см. (a).

Если клиент и/или сервер заняты чем-то другим, возможно, придется уменьшить количество частей. На это также может влиять способность дисков на сервере обрабатывать множество запросов ввода-вывода одновременно (следите за размером очереди и временем отклика).

Как лучше всего разделить таблицу с точки зрения доступа к данным?

Для эффективной параллельной обработки важно обеспечить равномерное распределение задач между обработчиками и минимизировать их взаимную синхронизацию. Кроме того, нужно помнить, что синхронизация обработчиков может происходить как на стороне сервера, так и на стороне клиента. Например, несколько обработчиков не должны использовать одно и то же соединение с базой данных. Менее очевидный пример: плохо, если разные обработчики читают записи с одних и тех же страниц базы данных. Например, когда два обработчика читают четные и нечетные записи - это неэффективно. Синхронизация на клиенте может происходить при распределении задач, при обработке полученных данных (выделение памяти под результаты) и так далее.

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

Давайте посмотрим, как gbak решает эту проблему.

Для gbak единицей работы является набор записей со страниц данных (DP), принадлежащих одной странице указателей (PP). С одной стороны, это достаточно большое количество записей, чтобы держать обработчик занятым без необходимости часто запрашивать новую порцию данных (синхронизация). С другой стороны, даже если такие наборы записей не имеют точно одинаковый размер, это позволит относительно равномерно загрузить обработчиков. То есть вполне возможно, что один обработчик прочитает N записей из одного PP, а другой - M записей, и M будет сильно отличаться от N. Такой подход не идеален, но его довольно просто реализовать, и он обычно достаточно эффективен, по крайней мере, в больших масштабах (с десятками или сотнями (или более) PP).

Как получить количество PP (страниц указателей) для заданной таблицы? Это довольно просто и, что важно, быстро вычисляется из таблицы RDB$PAGES:

sql
SELECT RDB$PAGE_SEQUENCE
FROM RDB$PAGES
WHERE RDB$RELATION_ID = ? AND RDB$PAGE_TYPE = 4
ORDER BY RDB$PAGE_SEQUENCE DESC ROWS 1

Далее мы могли бы просто разделить количество PP на количество обработчиков и выдать каждому обработчику свой фрагмент. Это подходит для сценария, когда параллельная обработка выполняется разработчиком, который знает распределение данных. Но для более распространенного сценария нет гарантии, что такие «большие» фрагменты будут означать одинаковый объем работы. Нам не интересно видеть ситуацию, когда 15 обработчиков закончили свою работу и простаивают, а 16-й долго читает свои 10 миллионов записей.

Вот почему gbak делает это иначе. Существует координатор работ, который выдает каждому процессору по одному PP за раз. Координатор знает, сколько всего PP и сколько уже выдано в работу. Когда обработчик завершает чтение своих записей, он обращается к координатору за новым номером PP. Это продолжается до тех пор, пока PP не закончатся (или есть активные обработчики). Конечно, такое взаимодействие обработчиков с координатором требует синхронизации. Опыт показывает, что объем работы, выдаваемый на один PP, позволяет не синхронизироваться слишком часто. Такой подход позволяет практически равномерно загрузить всех обработчиков (а значит, и ядра процессора) работой, независимо от фактического количества записей, принадлежащих каждому PP.

Как обработчик читает записи из своего PP? Для этого, начиная с Firebird 4.0 (впервые появилась в HQBird 2.5), существует встроенная функция MAKE_DBKEY(). С ее помощью можно получить RDB$DB_KEY (физический номер записи) для первой записи на указанном PP.

А с помощью этих RDB$DB_KEY выбираются нужные записи:

sql
SELECT *
FROM relation
WHERE RDB$DB_KEY >= MAKE_DBKEY(:rel_id, 0, 0, :loPP)
    AND RDB$DB_KEY < MAKE_DBKEY(:rel_id, 0, 0, :hiPP)

Например, если задать loPP = 0 и hiPP = 1, то будут прочитаны все записи с PP = 0, и только с него.

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

Реализация утилиты FBCSVExport

Утилита FBCSVExport предназначена для экспорта данных из таблиц базы данных Firebird в формат CSV.

Каждая таблица экспортируется в файл с именем .csv. В обычном (однопоточном) режиме данные из таблиц экспортируются последовательно в алфавитном порядке имен таблиц.

В параллельном режиме таблицы экспортируются параллельно, каждая таблица в отдельном потоке. Если таблица очень большая, она разбивается на части, и каждая часть экспортируется в отдельном потоке. Для каждой части большой таблицы создается отдельный файл с именем .csv.partN, где N - номер части.

Когда все части большой таблицы экспортированы, файлы частей объединяются в файл с именем .csv.

Для указания того, какие таблицы будут экспортированы, используется регулярное выражение. Экспортироваться могут только обычные таблицы (системные таблицы, GTT, представления, внешние таблицы не поддерживаются). Регулярные выражения должны быть в синтаксисе SQL, то есть те, которые используются в предикате SIMILAR TO.

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

sql
SELECT
    R.RDB$RELATION_ID AS RELATION_ID,
    TRIM(R.RDB$RELATION_NAME) AS RELATION_NAME,
    P.RDB$PAGE_SEQUENCE AS PAGE_SEQUENCE,
    COUNT(P.RDB$PAGE_SEQUENCE) OVER(PARTITION BY R.RDB$RELATION_NAME) AS PP_CNT
FROM RDB$RELATIONS R
JOIN RDB$PAGES P ON P.RDB$RELATION_ID = R.RDB$RELATION_ID
WHERE R.RDB$SYSTEM_FLAG = 0 AND
      R.RDB$RELATION_TYPE = 0 AND
      P.RDB$PAGE_TYPE = 4 AND
      TRIM(R.RDB$RELATION_NAME) SIMILAR TO CAST(? AS VARCHAR(8191))
ORDER BY R.RDB$RELATION_NAME, P.RDB$PAGE_SEQUENCE

В однопоточном режиме этот запрос можно упростить до

sql
SELECT
    R.RDB$RELATION_ID AS RELATION_ID,
    TRIM(R.RDB$RELATION_NAME) AS RELATION_NAME,
    0 AS PAGE_SEQUENCE,
    1 AS PP_CNT
FROM RDB$RELATIONS R
WHERE R.RDB$SYSTEM_FLAG = 0 AND
      R.RDB$RELATION_TYPE = 0 AND
      TRIM(R.RDB$RELATION_NAME) SIMILAR TO CAST(? AS VARCHAR(8191))
ORDER BY R.RDB$RELATION_NAME

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

Результат этого запроса формируется в вектор структур:

cpp
struct TableDesc
{
    TableDesc() = default;
    TableDesc(const OutputRecord& rec)
        : releation_id(rec->releation_id)
        , relation_name(rec->relation_name.str, rec->relation_name.length)
        , page_sequence(rec->page_sequence)
        , pp_cnt(rec->pp_cnt)
    {}

    short releation_id;
    std::string relation_name;
    int32_t page_sequence;
    int64_t pp_cnt;
};

Этот вектор заполняется с помощью функции, объявленной как:

cpp
std::vector getTablesDesc(
    Firebird::ThrowStatusWrapper* status,
    Firebird::IAttachment* att,
    Firebird::ITransaction* tra,
    unsigned int sqlDialect,
    const std::string& tableIncludeFilter,
    bool singleWorker = true);

Последний параметр singleWorker переключает режим заполнения std::vector: если singleWorker = true, используется запрос для однопоточного режима; если singleWorker = false, используется более затратный и сложный запрос для многопоточного режима. Я не буду приводить саму реализацию - она довольно проста, и вы можете увидеть ее в исходном коде проекта.

Для экспорта таблицы в формат CSV был разработан класс CSVExportTable, который содержит следующие методы:

cpp
    void prepare(Firebird::ThrowStatusWrapper* status, const std::string& tableName,
                 unsigned int sqlDialect, bool withDbkeyFilter = false);

    void printHeader(Firebird::ThrowStatusWrapper* status, csv::CSVFile& csv);

void printData(Firebird::ThrowStatusWrapper* status, csv::CSVFile& csv, int64_t ppNum = 0);

Code

Метод `prepare` предназначен для построения и подготовки запроса, который используется для экспорта таблицы в формат CSV. Внутренний запрос строится по-разному в зависимости от параметра `withDbkeyFilter`. Если `withDbkeyFilter = true`, то запрос строится с фильтрацией по диапазону `RDB$DB_KEY`:

```sql hljs
SELECT *
FROM tableName
WHERE RDB$DB_KEY >= MAKE_DBKEY('tableName', 0, 0, ?)
  AND RDB$DB_KEY < MAKE_DBKEY('tableName', 0, 0, ?)

в противном случае используется упрощённый запрос:

sql
SELECT *
FROM tableName

Значение параметра withDbkeyFilter устанавливается в true, если используется многопоточный режим и таблица большая. Мы считаем таблицу большой, если pp_cnt > 1.

Метод printHeader предназначен для вывода заголовка CSV-файла (имена столбцов таблицы).

Метод printData выводит данные таблицы в CSV-файл, начиная с номера страницы PP ppNum, если запрос был подготовлен с использованием фильтра по диапазону RDB$DB_KEY, и все данные таблицы в противном случае.

Теперь рассмотрим код для однопоточного режима

cpp
...

// Открытие основного соединения
Firebird::AutoRelease att(
    provider->attachDatabase(
        &status,
        m_database.c_str(),
        dbpLength,
        dpb
    )
);

// Запуск основной транзакции в режиме изоляции SNAPSHOT
Firebird::AutoDispose tpbBuilder(fbUtil->getXpbBuilder(&status, Firebird::IXpbBuilder::TPB, nullptr, 0));
tpbBuilder->insertTag(&status, isc_tpb_concurrency);

Firebird::AutoRelease tra(
    att->startTransaction(
        &status,
        tpbBuilder->getBufferLength(&status),
        tpbBuilder->getBuffer(&status)
    )
);
// Получение списка таблиц с использованием регулярного выражения в m_filter.
// m_parallel задаёт количество параллельных потоков; когда оно равно 1,
// используется упрощённый запрос для получения списка таблиц,
// в противном случае для каждой таблицы генерируется список PP и их количество.
auto tables = getTablesDesc(&status, att, tra, m_sqlDialect, m_filter, m_parallel == 1);

if (m_parallel == 1) {
    FBExport::CSVExportTable csvExport(att, tra, fb_master);
    for (const auto& tableDesc : tables) {
        // здесь нет смысла использовать фильтр по диапазону RDB$DB_KEY
        csvExport.prepare(&status, tableDesc.relation_name, m_sqlDialect, false);
        const std::string fileName = tableDesc.relation_name + ".csv";
        csv::CSVFile csv(m_outputDir / fileName);
        if (m_printHeader) {
            csvExport.printHeader(&status, csv);
        }
        csvExport.printData(&status, csv);
    }
}

Здесь всё довольно просто и не требует дополнительных пояснений, поэтому перейдём к многопоточной части.

Для того чтобы экспорт выполнялся в многопоточном режиме, необходимо создать дополнительные m_parallel - 1 рабочих потоков. Почему количество дополнительных потоков на 1 меньше? Да, потому что основной поток также будет экспортировать данные, и он равен дополнительным потокам. Вынесем общую часть основного и дополнительного потоков в отдельную функцию:

cpp
void ExportApp::exportByTableDesc(Firebird::ThrowStatusWrapper* status, FBExport::CSVExportTable& csvExport, const TableDesc& tableDesc)
{
    // Если tableDesc имеет pp_cnt > 1, то он описывает только часть таблицы, и необходимо построить
    // запрос с использованием фильтра по диапазону RDB$DB_KEY.

    bool withDbKeyFilter = tableDesc.pp_cnt > 1;
    csvExport.prepare(status, tableDesc.relation_name, m_sqlDialect, withDbKeyFilter);
    std::string fileName = tableDesc.relation_name + ".csv";

    // Если это не первая часть таблицы, то записываем эту часть в файл .csv.part, где
    // N - номер PP. Позже части таблицы будут объединены в единый файл .csv
    if (tableDesc.page_sequence > 0) {
        fileName += ".part_" + std::to_string(tableDesc.page_sequence);
    }
    csv::CSVFile csv(m_outputDir / fileName);
    // Заголовок CSV-файла должен выводиться только в первой части таблицы.
    if (tableDesc.page_sequence == 0 && m_printHeader) {
        csvExport.printHeader(status, csv);
    }
    csvExport.printData(status, csv, tableDesc.page_sequence);
}

Описания таблиц или их частей находятся в общем векторе со структурами TableDesc. Из этого вектора каждый рабочий поток берёт таблицу или следующую часть. Чтобы предотвратить гонки данных, необходимо синхронизировать доступ к общему ресурсу. Но сам std::vector не изменяется, поэтому можно синхронизировать только общую переменную, которая является индексом в этом векторе. Это легко сделать с помощью std::atomic в качестве такой переменной.

cpp
if (m_parallel == 1) {
    ...
}
else {
    // Определение количества дополнительных рабочих потоков
    const auto workerCount = m_parallel - 1;

    // Получение номера снимка из основной транзакции
    auto snapshotNumber = getSnapshotNumber(&status, tra);
    // переменная для хранения исключения внутри потока
    std::exception_ptr exceptionPointer = nullptr;
    std::mutex m;
	// атомарный счётчик
    // это индекс следующей таблицы или её части
    std::atomic<size_t> counter = 0;
    // пул рабочих потоков
    std::vector<std::thread> thread_pool;
    thread_pool.reserve(workerCount);
    for (int i = 0; i < workerCount; i++) {
		// для каждого потока создаём собственное соединение
        Firebird::AutoRelease workerAtt(
            provider->attachDatabase(
                &status,
                m_database.c_str(),
                dbpLength,
                dpb
            )
        );
		// и собственную транзакцию, которой передаём номер снимка
        // для создания общего снимка
        Firebird::AutoDispose tpbWorkerBuilder(fbUtil->getXpbBuilder(&status, Firebird::IXpbBuilder::TPB, nullptr, 0));
        tpbWorkerBuilder->insertTag(&status, isc_tpb_concurrency);
        tpbWorkerBuilder->insertBigInt(&status, isc_tpb_at_snapshot_number, snapshotNumber);

        Firebird::AutoRelease workerTra(
            workerAtt->startTransaction(
                &status,
                tpbWorkerBuilder->getBufferLength(&status),
                tpbWorkerBuilder->getBuffer(&status)
            )
        );
		// создаём поток
        std::thread t([att = std::move(workerAtt), tra = std::move(workerTra), this,
                       &m, &tables, &counter, &exceptionPointer]() mutable {

            Firebird::ThrowStatusWrapper status(fb_master->getStatus());
            try {
                FBExport::CSVExportTable csvExport(att, tra, fb_master);
                while (true) {
                    // инкрементируем атомарный счётчик
                    size_t localCounter = counter++;
					// если таблицы или их части закончились, выходим
                    // из бесконечного цикла и завершаем поток
                    if (localCounter >= tables.size())
                        break;
                    // получаем описание таблицы или её части
                    const auto& tableDesc = tables[localCounter];
                    // и выполняем экспорт
                    exportByTableDesc(&status, csvExport, tableDesc);
                }
                if (tra) {
                    tra->commit(&status);
                    tra.release();
                }

                if (att) {
                    att->detach(&status);
                    att.release();
                }
            }
            catch (...) {
				// если возникло исключение, сохраняем его для
                // последующего освобождения в основном потоке
                std::unique_lock<std::mutex> lock(m);
                exceptionPointer = std::current_exception();
            }
        });
        thread_pool.push_back(std::move(t));
    }

    // экспорт в основном потоке
    FBExport::CSVExportTable csvExport(att, tra, fb_master);
    while (true) {
        // инкрементируем атомарный счётчик
        size_t localCounter = counter++;
        if (localCounter >= tables.size())
            break;
        // если таблицы или их части закончились, выходим

```cpp
// из бесконечного цикла
        const auto& tableDesc = tables[localCounter];
        exportByTableDesc(&status, csvExport, tableDesc);
    }
    // ожидание завершения рабочих потоков
    for (auto& th : thread_pool) {
        th.join();
    }
    // если в рабочих потоках было исключение, выбрасываем его снова
    if (exceptionPointer) {
        std::rethrow_exception(exceptionPointer);
    }
    ...

Осталось только объединить файлы, созданные для частей таблиц, в единый файл для каждой из этих таблиц.

cpp
for (size_t i = 0; i < tables.size(); i++) {
    const auto& tableDesc = tables[i];
    // если количество PP больше 1,
    // то таблица большая и для неё было несколько частей
    if (tableDesc.pp_cnt > 1) {
        // основной файл для таблицы
        std::string fileName = tableDesc.relation_name + ".csv";
        std::ofstream ofile(m_outputDir / fileName, std::ios::out | std::ios::app);
        i++;
        for (int64_t j = 1; j < tableDesc.pp_cnt; j++, i++) {
            // файлы частей таблицы
            std::string partFileName = fileName + ".part_" + std::to_string(j);
            auto partFilePath = m_outputDir / partFileName;
            std::ifstream ifile(partFilePath, std::ios::in);
            ofile << ifile.rdbuf();
            ifile.close();
            fs::remove(partFilePath);
        }
        ofile.close();
    }
}

Давайте измерим производительность инструмента в однопоточном и многопоточном режимах.

Бенчмарк инструмента FBCSVExport

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

  • Операционная система: Windows 10 x64.

  • Процессор: Intel Core i3 8100, 4 ядра, 4 потока.

  • Память: 16 ГБ

  • Дисковая подсистема: NVME SSD (база данных), SATA SSD (папка для хранения CSV-файлов).

  • Firebird 4.0.4 x64

Результаты:

bash
CSVExport.exe -H --table-filter="COLOR|BREED|HORSE|COVER|MEASURE|LAB_LINE|SEX" --parallel=1 \
  -d inet://localhost:3054/horses -u SYSDBA -p masterkey --charset=WIN1251 -o ./single

Elapsed time in milliseconds parallel_part: 35894 ms
Elapsed time in milliseconds: 36317 ms

CSVExport.exe -H --table-filter="COLOR|BREED|HORSE|COVER|MEASURE|LAB_LINE|SEX" --parallel=2 \
  -d inet://localhost:3054/horses -u SYSDBA -p masterkey --charset=WIN1251 -o ./multi

Elapsed time in milliseconds parallel_part: 19259 ms
Elapsed time in milliseconds: 20760 ms

CSVExport.exe -H --table-filter="COLOR|BREED|HORSE|COVER|MEASURE|LAB_LINE|SEX" --parallel=4 \
  -d inet://localhost:3054/horses -u SYSDBA -p masterkey --charset=WIN1251 -o ./multi

Elapsed time in milliseconds parallel_part: 19600 ms
Elapsed time in milliseconds: 21137 ms

Из результатов тестирования видно, что при использовании двух потоков ускорение составило 1,8 раза, что является хорошим результатом. Но параллельное выполнение экспорта в 4 потока также показало ускорение в 1,8 раза. Почему не в 3-4 раза? Дело в том, что сервер Firebird и утилита экспорта работают на одном компьютере, у которого всего 4 ядра. Таким образом, сам сервер Firebird использует 4 потока для чтения таблицы, а утилита FBCSVExport также использует 4 потока. Очевидно, что в этом случае довольно сложно добиться ускорения более чем в 2 раза. Поэтому попробуем на другом оборудовании, где количество ядер значительно больше.

Linux

  • Операционная система: CentOS 8.

  • Процессор: 2 процессора Intel Xeon E5-2603 v4, всего 12 ядер, 12 потоков.

  • Память: 32 ГБ

  • Дисковая подсистема: SAS HDD (RAID 10)

  • Firebird 4.0.4 x64

Результаты:

bash
[denis@copyserver build]$ ./CSVExport -H --table-filter="COLOR|BREED|HORSE|COVER|MEASURE|LAB_LINE|SEX" --parallel=1 \
  -d inet://localhost/horses -u SYSDBA -p masterkey --charset=UTF8 -o ./single

Elapsed time in milliseconds parallel_part: 57547 ms
Elapsed time in milliseconds: 57595 ms

[denis@copyserver build]$ ./CSVExport -H --table-filter="COLOR|BREED|HORSE|COVER|MEASURE|LAB_LINE|SEX" --parallel=4 \
  -d inet://localhost/horses -u SYSDBA -p masterkey --charset=UTF8 -o ./multi

Elapsed time in milliseconds parallel_part: 17755 ms
Elapsed time in milliseconds: 18148 ms

[denis@copyserver build]$ ./CSVExport -H --table-filter="COLOR|BREED|HORSE|COVER|MEASURE|LAB_LINE|SEX" --parallel=6 \
  -d inet://localhost/horses -u SYSDBA -p masterkey --charset=UTF8 -o ./multi

Elapsed time in milliseconds parallel_part: 13243 ms
Elapsed time in milliseconds: 13624 ms

[denis@copyserver build]$ ./CSVExport -H --table-filter="COLOR|BREED|HORSE|COVER|MEASURE|LAB_LINE|SEX" --parallel=12 \
  -d inet://localhost/horses -u SYSDBA -p masterkey --charset=UTF8 -o ./multi

Elapsed time in milliseconds parallel_part: 12712 ms
Elapsed time in milliseconds: 13140 ms

В этом случае оптимальное количество потоков для экспорта - 6 (6 потоков для Firebird и 6 потоков для утилиты FBCSVExport). При этом нам удалось добиться ускорения в 5 раз, что свидетельствует о довольно хорошей масштабируемости. На Linux-сервере и Windows-компьютере мы использовали идентичные базы данных, и вы, вероятно, заметили, что однопоточный экспорт на Windows был почти в 2 раза быстрее: это связано с более быстрой дисковой подсистемой (NVME-накопитель намного быстрее, чем SAS-диски, объединённые в RAID).

Резюме

В этой статье мы рассмотрели, как эффективно читать данные из таблиц СУБД Firebird с использованием параллелизма. Также был показан пример того, как можно использовать некоторые возможности СУБД Firebird для организации такого чтения в вашем программном обеспечении.

Огромное спасибо Владиславу Хорсуну, разработчику ядра Firebird, за помощь в подготовке этого материала.

По любым вопросам или замечаниям, пожалуйста, пишите на [email protected].