Ta strona została przetłumaczona maszynowo. Przeczytaj oryginał angielski. English

Biblioteka IBSurgeon

Równoległe odczytywanie danych w Firebird

D.Simonov, V.Horsun

wersja 1.0.5 z 05.12.2023

Ten materiał jest sponsorowany i stworzony przy wsparciu i patronacie IBSurgeon www.ib-aid.com, dostawcy HQbird (zaawansowanej dystrybucji Firebird) oraz usług optymalizacji wydajności, migracji i wsparcia technicznego dla Firebird.

Materiał jest licencjonowany na podstawie Public Documentation License https://www.firebirdsql.org/file/documentation/html/en/licenses/pdl/public-documentation-license.html

Powiązane materiały:

Przedmowa

Firebird 5.0 wprowadził możliwość wykorzystania równoległości podczas tworzenia kopii zapasowej za pomocą narzędzia gbak, między innymi funkcje równoległe. Początkowo funkcja ta pojawiła się w HQbird 2.5, następnie w HQbird 3.0 i 4.0, a później została przeniesiona do Firebird 5.0.

W tym artykule rozważymy funkcje używane podczas tworzenia równoległych kopii zapasowych wewnątrz narzędzia gbak. Pokażemy również, jak można je wykorzystać w aplikacjach do równoległego odczytu danych.

Ważne jest, aby zauważyć, że nie mówimy tutaj o równoległym skanowaniu tabel wewnątrz silnika Firebird podczas wykonywania zapytań SQL, ale o odczycie danych wewnątrz aplikacji w równoległych strumieniach.

Przykładowe narzędzie FBCSVExport

Aby zademonstrować równoległy odczyt danych z systemu DBMS Firebird, napisano przykładowe narzędzie, które eksportuje dane z jednej lub więcej tabel do formatu CSV.

Jego opis oraz otwarty kod źródłowy znajdują się tutaj: https://github.com/IBSurgeon/FBCSVExport.git

Jak widać w opisie narzędzia i z poniższego artykułu, dzięki przetwarzaniu równoległemu możliwe jest eksportowanie danych i wykonywanie innych operacji równoległych 2-10 razy szybciej niż w 1 wątku (w zależności od sprzętu).

W przypadku jakichkolwiek pytań prosimy o kontakt: [email protected].

Odczyt równoległy

Zastanówmy się, jak czytać dane z kilku tabel równolegle. Jak wiadomo, Firebird pozwala na wykonywanie zapytań równolegle tylko wtedy, gdy każde zapytanie jest wykonywane w osobnej sesji połączenia.

Utwórzmy pulę wątków roboczych. Główny wątek aplikacji jest również wątkiem roboczym, więc liczba dodatkowych wątków roboczych powinna wynosić N - 1, gdzie N to całkowita liczba równoległych pracowników. Każdy wątek roboczy będzie prowadził własne połączenie i transakcję.

Pierwszy problem: jak zapewnić spójność odczytywanych danych?

Spójny odczyt danych

Ponieważ każdy wątek roboczy używa własnego połączenia i własnej transakcji, pojawia się problem niespójnych odczytów - jeśli tabela jest jednocześnie zmieniana przez innych użytkowników, odczytane dane mogą być niespójne. W trybie jednowątkowym gbak używa transakcji z trybem izolacji SNAPSHOT, co umożliwia odczyt spójnych informacji na początku transakcji SNAPSHOT. Ale tutaj mamy wiele transakcji i konieczne jest, aby widziały one ten sam “migawkę” (snapshot), aby odczytywały te same niezmienne dane.

Mechanizm tworzenia wspólnej migawki dla różnych transakcji z trybem izolacji SNAPSHOT został wprowadzony w Firebird 4.0 (pierwotnie w HQBird 2.5, ale w Firebird 4.0/HQBird 4.0 jest prostszy i bardziej wydajny). Istnieją dwa sposoby utworzenia wspólnej migawki:

  1. Za pomocą SQL
  • pobierz numer migawki z głównej transakcji (która jest uruchamiana w głównym wątku roboczym).
sql
SELECT RDB$GET_CONTEXT('SYSTEM', 'SNAPSHOT_NUMBER') FROM RDB$DATABASE
  • uruchom inne transakcje za pomocą następującego SQL:
sql
SET TRANSACTION SNAPSHOT AT NUMBER snapshot_number

gdzie snapshot_number to numer pobrany przez poprzednie zapytanie.

  1. Za pomocą API
  • pobierz numer migawki z głównej transakcji (która jest uruchamiana w głównym wątku roboczym) za pomocą funkcji isc_transaction_info lub ITransaction.getInfo ze znacznikiem fb_info_tra_snapshot_number;

  • uruchom inne transakcje ze znacznikiem isc_tpb_at_snapshot_number z uzyskanym numerem migawki.

Przykładowe narzędzie FBCSVExport, podobnie jak gbak, używa drugiego podejścia. Te podejścia można mieszać - na przykład pobrać numer migawki za pomocą SQL, a następnie użyć uzyskanego numeru migawki do uruchomienia innych transakcji za pomocą API, lub odwrotnie.

W FBCSVExport pobieramy numer migawki za pomocą następującego kodu:

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;
}

Aby uruchomić transakcję z numerem migawki, używamy następującego kodu:

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)
    )
);

Teraz dane odczytywane z różnych połączeń będą spójne, więc możemy rozłożyć obciążenie na wątki robocze.

Jak dokładnie rozłożyć obciążenie między wątkami roboczymi? W przypadku pełnego eksportu wszystkich tabel lub kopii zapasowej najprostszą opcją będzie jeden wątek roboczy na tabelę. Ale przy tym podejściu mamy następujący problem: jeśli w bazie danych jest wiele małych tabel i jedna duża, albo nawet tylko jedna tabela i jest ogromna, nie zobaczymy poprawy. W takim przypadku jeden wątek otrzyma dużą tabelę, a pozostałe wątki będą bezczynne. Aby temu zapobiec, konieczne jest przetwarzanie dużej tabeli w częściach.

Uwaga Poniższy materiał poświęcony jest pełnemu odczytowi tabel. Jeśli chcesz zorganizować równoległy odczyt z jakiegoś zapytania (lub widoku), będzie wymagać to nieco innego podejścia, które zależy od rzeczywistych danych.

Podział dużej tabeli na części

Załóżmy, że mamy tylko jedną dużą tabelę, którą chcemy odczytać w całości i jak najszybciej. Proponuje się podzielić ją na kilka części i odczytywać każdą część z własnego strumienia niezależnie. Każdy wątek musi mieć własne połączenie z bazą danych.

W tym przypadku pojawiają się następujące pytania:

  • Na ile części przetwarzania należy podzielić tabelę?

  • Jaki jest najlepszy sposób podziału tabeli pod względem dostępu do danych?

Odpowiedzmy na te pytania po kolei.

Na ile części przetwarzania należy podzielić tabelę?

Załóżmy idealny scenariusz - serwer i klient są dedykowane dla Firebird, to znaczy wszystkie procesory są całkowicie do naszej dyspozycji. Wtedy zaleca się:

a) Jako maksymalną liczbę równoległych części użyć podwojonej liczby rdzeni procesora na serwerze. Dlaczego 2x rdzeni? Wiemy na pewno, że wystąpią opóźnienia związane z operacjami wejścia/wyjścia, więc możemy pozwolić na pewne dodatkowe wykorzystanie procesora. Jednak liczbę tę należy traktować jako ustawienie początkowe; w praktyce zależy ona od danych.

b) Wziąć pod uwagę liczbę rdzeni na kliencie: jeśli na serwerze jest ich znacznie więcej (typowa sytuacja), to może mieć sens dalsze ograniczenie liczby części podziału, aby nie przeciążać klienta (i tak nie będzie w stanie przetworzyć więcej, a koszty przełączania strumieni nie znikną). Będzie można zdecydować dokładniej, monitorując obciążenie procesora klienta i serwera - jeśli na kliencie wynosi 100%, a na serwerze zauważalnie mniej, to ma sens zmniejszenie liczby części.

c) jeśli klient i serwer to ten sam host, patrz (a).

Jeśli klient i/lub serwer są zajęte czymś innym, może być konieczne zmniejszenie liczby części. Na to może również wpływać zdolność dysków na serwerze do przetwarzania wielu żądań wejścia/wyjścia jednocześnie (monitoruj rozmiar kolejki i czas odpowiedzi).

Jaki jest najlepszy sposób podziału tabeli pod względem dostępu do danych?

Aby zaimplementować efektywne przetwarzanie równoległe, ważne jest zapewnienie równomiernego rozłożenia zadań między obsługującymi i zminimalizowanie ich wzajemnej synchronizacji. Ponadto należy pamiętać, że synchronizacja obsługujących może występować zarówno po stronie serwera, jak i po stronie klienta. Na przykład kilka obsługujących nie powinno używać tego samego połączenia z bazą danych. Mniej oczywisty przykład: źle jest, jeśli różne obsługujące odczytują rekordy z tych samych stron bazy danych. Na przykład, gdy dwóch obsługujących odczytuje rekordy parzyste i nieparzyste - to nie jest efektywne. Synchronizacja po stronie klienta może wystąpić podczas dystrybucji zadań, podczas przetwarzania otrzymanych danych (przydzielania pamięci na wyniki) i tak dalej.

Jednym z problemów ze “sprawiedliwym” podziałem jest to, że klient nie wie, jak rekordy są rozmieszczone na stronach (i kluczach indeksu), ile jest rekordów lub stron danych (dla dużych tabel zbyt długo trwałoby wcześniejsze policzenie liczby rekordów).

Zobaczmy, jak gbak rozwiązuje ten problem.

Dla gbak jednostką pracy jest zestaw rekordów ze stron danych (DP) należących do tej samej strony wskaźników (PP). Z jednej strony jest to dość duża liczba rekordów, aby utrzymać obsługującego zajętego bez konieczności częstego proszenia o nową porcję danych (synchronizacja). Z drugiej strony, nawet jeśli takie zestawy rekordów nie mają dokładnie tego samego rozmiaru, pozwoli to na stosunkowo równomierne obciążenie pracowników. To znaczy, całkiem możliwe, że jeden pracownik odczyta N rekordów z jednej PP, a drugi M rekordów, gdzie M może się znacznie różnić od N. To podejście nie jest idealne, ale jest dość proste do wdrożenia i zwykle dość skuteczne, przynajmniej na dużą skalę (przy dziesiątkach lub setkach (lub więcej) PP).

Jak uzyskać liczbę PP (Pointer Pages) dla danej tabeli? Jest to dość łatwe i, co najważniejsze, szybkie do obliczenia z tabeli 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

Następnie moglibyśmy po prostu podzielić liczbę PP przez liczbę pracowników i dać każdemu pracownikowi ich własną część. Jest to w porządku w scenariuszu, gdy przetwarzanie równoległe jest wykonywane przez programistę, który zna rozkład danych. Ale w bardziej typowym scenariuszu nie ma gwarancji, że takie “duże” części będą oznaczać tę samą ilość pracy. Nie chcemy widzieć sytuacji, w której 15 pracowników skończyło swoją pracę i stoi bezczynnie, a 16. czyta swoje 10 milionów rekordów przez długi czas.

Dlatego gbak robi to inaczej. Istnieje koordynator pracy, który wydaje każdemu procesorowi 1 PP na raz. Koordynator wie, ile PP jest łącznie i ile już wydano do pracy. Gdy pracownik zakończy odczyt swoich rekordów, kontaktuje się z koordynatorem w celu uzyskania nowego numeru PP. Kontynuuje to, dopóki PP się nie wyczerpią (lub dopóki są aktywni pracownicy). Oczywiście taka interakcja pracowników z koordynatorem wymaga synchronizacji. Doświadczenie pokazuje, że ilość pracy przypadająca na jedną PP pozwala nie synchronizować się zbyt często. To podejście pozwala praktycznie równomiernie obciążyć wszystkich pracowników (a więc i rdzenie procesora) pracą, niezależnie od rzeczywistej liczby rekordów należących do każdej PP.

Jak handler odczytuje rekordy ze swojego PP? W tym celu, począwszy od Firebird 4.0 (po raz pierwszy pojawił się w HQBird 2.5) istnieje wbudowana funkcja MAKE_DBKEY(). Za jej pomocą można uzyskać RDB$DB_KEY (fizyczny numer rekordu) dla pierwszego rekordu na określonym PP.

A za pomocą tych RDB$DB_KEY wybierane są potrzebne rekordy:

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)

Na przykład, jeśli ustawisz loPP = 0 i hiPP = 1, wtedy wszystkie rekordy z PP = 0 zostaną odczytane, i tylko z niego.

Teraz, gdy masz pojęcie o tym, jak działa gbak, możesz przejść do opisu implementacji narzędzia FBCSVExport.

Implementacja narzędzia FBCSVExport

Narzędzie FBCSVExport jest przeznaczone do eksportowania danych z tabel bazy danych Firebird do formatu CSV.

Każda tabela jest eksportowana do pliku o nazwie .csv. W normalnym (jednowątkowym) trybie dane z tabel są eksportowane sekwencyjnie w kolejności alfabetycznej nazw tabel.

W trybie równoległym tabele są eksportowane równolegle, każda tabela w osobnym wątku. Jeśli tabela jest bardzo duża, jest dzielona na części, a każda część jest eksportowana w osobnym strumieniu. Dla każdej części dużej tabeli tworzony jest osobny plik o nazwie .csv.partN, gdzie N to numer części.

Gdy wszystkie części dużej tabeli zostaną wyeksportowane, pliki części są scalane w plik o nazwie .csv.

Do określenia, które tabele zostaną wyeksportowane, używane jest wyrażenie regularne. Mogą być eksportowane tylko zwykłe tabele (tabele systemowe, GTT, widoki, tabele zewnętrzne nie są obsługiwane). Wyrażenia regularne muszą być w składni SQL, to znaczy takie, które są używane w predykacie SIMILAR TO.

Do wybrania listy eksportowanych tabel, a także listy ich PP w trybie wielowątkowym, używamy następującego zapytania:

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

W trybie jednowątkowym to zapytanie można uprościć do

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

W trybie jednowątkowym wartości pól PAGE_SEQUENCE i PP_CNT nie są używane; są dodawane do zapytania w celu ujednolicenia komunikatów wyjściowych.

Wynik tego zapytania jest formowany w wektor struktur:

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;
};

Ten wektor jest wypełniany za pomocą funkcji zadeklarowanej jako:

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

Ostatni parametr singleWorker przełącza tryb wypełniania std::vector - jeśli singleWorker = true, używane jest zapytanie dla trybu jednowątkowego, jeśli singleWorker = false, używane jest bardziej kosztowne i złożone zapytanie dla trybu wielowątkowego. Nie podam samej implementacji, jest dość prosta i możesz ją zobaczyć w kodzie źródłowym projektu.

Do eksportowania tabeli do formatu CSV opracowano klasę CSVExportTable, która zawiera następujące metody:

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);

Metoda prepare jest przeznaczona do zbudowania i przygotowania zapytania, które jest używane do eksportowania tabeli do formatu CSV. Wewnętrzne zapytanie jest konstruowane inaczej w zależności od parametru withDbkeyFilter. Jeśli withDbkeyFilter = true, zapytanie jest budowane z filtrowaniem po zakresie RDB$DB_KEY:

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

w przeciwnym razie używane jest uproszczone zapytanie:

sql
SELECT *
FROM tableName

Wartość parametru withDbkeyFilter jest ustawiana na true, jeśli używany jest tryb wielowątkowy i tabela jest duża. Uznajemy tabelę za dużą, jeśli pp_cnt > 1.

Metoda printHeader jest przeznaczona do wydrukowania nagłówka pliku CSV (nazwy kolumn tabeli).

Metoda printData drukuje dane tabeli do pliku CSV ze strony PP o numerze ppNum, jeśli zapytanie zostało przygotowane z użyciem filtra po zakresie RDB$DB_KEY, a wszystkie dane tabeli w przeciwnym razie.

Teraz przyjrzyjmy się kodowi dla trybu jednowątkowego

cpp
...

// Otwarcie głównego połączenia
Firebird::AutoRelease att(
    provider->attachDatabase(
        &status,
        m_database.c_str(),
        dbpLength,
        dpb
    )
);

// Rozpoczęcie głównej transakcji w trybie izolacji 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)
    )
);
// Pobranie listy tabel za pomocą wyrażenia regularnego w m_filter.
// m_parallel ustawia liczbę równoległych wątków, gdy jest równa 1,
// wtedy używane jest uproszczone zapytanie do uzyskania listy tabel,
// w przeciwnym razie dla każdej tabeli generowana jest lista PP i ich liczba.
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) {
        // nie ma sensu używać filtra zakresu RDB$DB_KEY tutaj
        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);
    }
}

Tutaj wszystko jest dość proste i nie wymaga dodatkowego wyjaśnienia, więc przejdźmy do części wielowątkowej.

Aby eksport odbywał się w trybie wielowątkowym, konieczne jest utworzenie dodatkowych m_parallel - 1 wątków roboczych. Dlaczego liczba dodatkowych wątków jest o 1 mniejsza? Tak, ponieważ główny wątek również będzie eksportował dane i jest równy z dodatkowymi wątkami. Przenieśmy wspólną część głównego i dodatkowego przepływu do osobnej funkcji:

cpp
void ExportApp::exportByTableDesc(Firebird::ThrowStatusWrapper* status, FBExport::CSVExportTable& csvExport, const TableDesc& tableDesc)
{
    // Jeśli tableDesc ma pp_cnt > 1, to opisuje tylko część tabeli i konieczne jest zbudowanie
    // zapytania z użyciem filtra po zakresie 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";

    // Jeśli to nie jest pierwsza część tabeli, zapisz tę część do pliku .csv.part, gdzie
    // N - numer PP. Później części tabeli zostaną połączone w jeden plik .csv
    if (tableDesc.page_sequence > 0) {
        fileName += ".part_" + std::to_string(tableDesc.page_sequence);
    }
    csv::CSVFile csv(m_outputDir / fileName);
    // Nagłówek pliku CSV powinien być wydrukowany tylko w pierwszej części tabeli.
    if (tableDesc.page_sequence == 0 && m_printHeader) {
        csvExport.printHeader(status, csv);
    }
    csvExport.printData(status, csv, tableDesc.page_sequence);
}

Opisy tabel lub ich części znajdują się we wspólnym wektorze ze strukturami TableDesc. Z tego wektora każdy wątek roboczy pobiera tabelę lub następną część. Aby zapobiec wyścigom danych, konieczne jest zsynchronizowanie dostępu do współdzielonego zasobu. Ale sam std::vector się nie zmienia, więc można zsynchronizować tylko współdzieloną zmienną, która jest indeksem w tym wektorze. Można to łatwo zrobić za pomocą std::atomic jako takiej zmiennej.

cpp
if (m_parallel == 1) {
    ...
}
else {
    // Określenie liczby dodatkowych wątków roboczych
    const auto workerCount = m_parallel - 1;

    // Pobranie numeru snapshotu z głównej transakcji
    auto snapshotNumber = getSnapshotNumber(&status, tra);
    // zmienna do przechowywania wyjątku w wątku
    std::exception_ptr exceptionPointer = nullptr;
    std::mutex m;
	// licznik atomowy
    // jest indeksem następnej tabeli lub jej części
    std::atomic<size_t> counter = 0;
    // pula wątków roboczych
    std::vector<std::thread> thread_pool;
    thread_pool.reserve(workerCount);
    for (int i = 0; i < workerCount; i++) {
		// dla każdego wątku tworzymy własne połączenie
        Firebird::AutoRelease workerAtt(
            provider->attachDatabase(
                &status,
                m_database.c_str(),
                dbpLength,
                dpb
            )
        );
		// i własną transakcję, do której przekazujemy numer snapshotu
        // aby utworzyć wspólny snapshot
        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)
            )
        );
		// utworzenie wątku
        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) {
                    // inkrementacja licznika atomowego
                    size_t localCounter = counter++;
					// jeśli tabele lub ich części się skończyły, wyjdź
                    // z nieskończonej pętli i zakończ wątek
                    if (localCounter >= tables.size())
                        break;
                    // pobierz opis tabeli lub jej części
                    const auto& tableDesc = tables[localCounter];
                    // i wykonaj eksport
                    exportByTableDesc(&status, csvExport, tableDesc);
                }
                if (tra) {
                    tra->commit(&status);
                    tra.release();
                }

                if (att) {
                    att->detach(&status);
                    att.release();
                }
            }
            catch (...) {
				// jeśli wystąpi wyjątek, zapisz go dla
                // późniejszego zwolnienia w głównym wątku
                std::unique_lock<std::mutex> lock(m);

exceptionPointer = std::current_exception();
            }
        });
        thread_pool.push_back(std::move(t));
    }

    // export in the main thread
    FBExport::CSVExportTable csvExport(att, tra, fb_master);
    while (true) {
        // increment the atomic counter
        size_t localCounter = counter++;
        if (localCounter >= tables.size())
            break;
        // if the tables or their parts are over, exit
        // from an endless loop
        const auto& tableDesc = tables[localCounter];
        exportByTableDesc(&status, csvExport, tableDesc);
    }
    // wait for the worker threads to complete
    for (auto& th : thread_pool) {
        th.join();
    }
    // if there was an exception in the worker threads, throw it again
    if (exceptionPointer) {
        std::rethrow_exception(exceptionPointer);
    }
    ...

Pozostaje tylko połączyć pliki utworzone dla części tabel w jeden plik dla każdej z tych tabel.

cpp
for (size_t i = 0; i < tables.size(); i++) {
    const auto& tableDesc = tables[i];
    // if the number of PP is greater than 1,
    // then the table is large and there were several parts for it
    if (tableDesc.pp_cnt > 1) {
        // main file for the table
        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++) {
            // files of table parts
            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();
    }
}

Zmierzmy wydajność narzędzia w trybie jednowątkowym i wielowątkowym.

Benchmark narzędzia FBCSVExport

Najpierw przyjrzyjmy się wynikom porównania trybu wielowątkowego i jednowątkowego eksportu na przeciętnym komputerze domowym. === Windows

  • System operacyjny: Windows 10 x64.

  • Procesor: Intel Core i3 8100, 4 rdzenie, 4 wątki.

  • Pamięć: 16 GB

  • Podsystem dyskowy: NVME SSD (baza danych), SATA SSD (folder do przechowywania plików CSV).

  • Firebird 4.0.4 x64

Wyniki:

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

Z wyników testów jasno wynika, że przy użyciu dwóch wątków przyspieszenie wyniosło 1,8 raza, co jest dobrym wynikiem. Jednak równoległe wykonanie eksportu w 4 wątkach również wykazało poprawę o 1,8 raza. Dlaczego nie 3-4 razy? Faktem jest, że serwer Firebird i narzędzie eksportujące działają na tym samym komputerze, który ma tylko 4 rdzenie. Zatem sam serwer Firebird używa 4 wątków do odczytu tabeli, a narzędzie FBCSVExport również używa 4 wątków. Oczywiście w tym przypadku dość trudno osiągnąć przyspieszenie większe niż 2 razy. Dlatego spróbujemy na innym sprzęcie, gdzie liczba rdzeni jest znacznie większa.

Linux

  • System operacyjny: CentOS 8.

  • Procesor: 2 procesory Intel Xeon E5-2603 v4, łącznie 12 rdzeni, 12 wątków.

  • Pamięć: 32 GB

  • Podsystem dyskowy: SAS HDD (RAID 10)

  • Firebird 4.0.4 x64

Wyniki:

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

W tym przypadku optymalna liczba wątków do eksportu to 6 (6 wątków dla Firebird i 6 wątków dla narzędzia FBCSVExport). Jednocześnie udało nam się osiągnąć 5-krotne przyspieszenie, co wskazuje na dość dobrą skalowalność. Na serwerze Linux i komputerze Windows użyliśmy identycznych baz danych i prawdopodobnie zauważyłeś, że eksport jednowątkowy na Windows był prawie 2 razy szybszy: wynika to z szybszego podsystemu dyskowego (dysk NVME jest znacznie szybszy niż dyski SAS połączone w RAID).

Podsumowanie

W tym artykule rozważyliśmy, jak efektywnie odczytywać dane z tabel systemu DBMS Firebird przy użyciu równoległości. Pokazano również przykład, jak można wykorzystać niektóre możliwości systemu DBMS Firebird do zorganizowania takiego odczytu w swoim oprogramowaniu.

Serdeczne podziękowania dla Vladislava Khorsuna, głównego programisty Firebird, za pomoc w przygotowaniu tego materiału.

W przypadku pytań lub uwag prosimy o kontakt e-mail: [email protected].