Cette page a été traduite automatiquement. Lisez l'original en anglais. English

Bibliothèque IBSurgeon

Lecture parallèle des données dans Firebird

D.Simonov, V.Horsun

version 1.0.5 du 05.12.2023

Ce matériel est sponsorisé et créé avec le parrainage et le soutien d’IBSurgeon www.ib-aid.com, fournisseur de HQbird (distribution avancée de Firebird) et prestataire de services d’optimisation des performances, de migration et de support technique pour Firebird.

Ce matériel est sous licence Public Documentation License https://www.firebirdsql.org/file/documentation/html/en/licenses/pdl/public-documentation-license.html

Documents connexes :

Préface

Firebird 5.0 a introduit la possibilité d’utiliser le parallélisme lors de la création d’une sauvegarde avec l’utilitaire gbak, parmi d’autres fonctions parallèles. Initialement, cette fonction est apparue dans HQbird 2.5, puis dans HQbird 3.0 et 4.0, avant d’être portée dans Firebird 5.0.

Dans cet article, nous examinerons les fonctions utilisées lors de la création de sauvegardes parallèles dans l’utilitaire gbak. Nous montrerons également comment les utiliser dans vos applications pour la lecture parallèle des données.

Il est important de noter qu’il ne s’agit pas ici de l’analyse parallèle des tables dans le moteur Firebird lors de l’exécution de requêtes SQL, mais de la lecture des données dans des flux parallèles au sein de votre application.

Outil d’exemple FBCSVExport

Pour démontrer la lecture parallèle des données depuis le SGBD Firebird, un utilitaire d’exemple a été écrit pour exporter les données d’une ou plusieurs tables au format CSV.

Sa description et son code source ouvert sont disponibles ici : https://github.com/IBSurgeon/FBCSVExport.git

Comme vous pouvez le voir dans la description de l’utilitaire et dans l’article ci-dessous, grâce au traitement parallèle, il est possible d’exporter des données et d’effectuer d’autres opérations parallèles 2 à 10 fois plus rapidement qu’avec un seul thread (selon le matériel).

Pour toute question, veuillez contacter [email protected].

Lecture parallèle

Réfléchissons à la manière de lire des données de plusieurs tables en parallèle. Comme vous le savez, Firebird ne permet d’exécuter des requêtes en parallèle que si chaque requête est exécutée dans une connexion séparée.

Créons un pool de threads de travail. Le thread principal de l’application est également un thread de travail, donc le nombre de threads de travail supplémentaires doit être N - 1, où N est le nombre total de workers parallèles. Chaque thread de travail exécutera sa propre connexion et sa propre transaction.

Le premier problème : comment garantir la cohérence des données lues ?

Lecture cohérente des données

Comme chaque thread de travail utilise sa propre connexion et sa propre transaction, le problème des lectures incohérentes se pose - si la table est modifiée simultanément par d’autres utilisateurs, les données lues peuvent être incohérentes. En mode mono-thread, gbak utilise une transaction avec le mode d’isolation SNAPSHOT, ce qui permet de lire des informations cohérentes au début de la transaction SNAPSHOT. Mais ici, nous avons plusieurs transactions et il est nécessaire qu’elles voient le même « snapshot » afin de lire les mêmes données immuables.

Le mécanisme de création d’un snapshot partagé pour différentes transactions avec le mode d’isolation SNAPSHOT a été introduit dans Firebird 4.0 (à l’origine dans HQBird 2.5, mais dans Firebird 4.0/HQBird 4.0, il est plus simple et plus efficace). Il existe deux façons de créer un snapshot partagé :

  1. Avec SQL
  • obtenir le numéro de snapshot de la transaction principale (démarrée dans le thread de travail principal).
sql
SELECT RDB$GET_CONTEXT('SYSTEM', 'SNAPSHOT_NUMBER') FROM RDB$DATABASE
  • démarrer les autres transactions avec le SQL suivant :
sql
SET TRANSACTION SNAPSHOT AT NUMBER snapshot_number

snapshot_number est le numéro récupéré par la requête précédente.

  1. Avec l’API
  • obtenir le numéro de snapshot de la transaction principale (démarrée dans le thread de travail principal) avec la fonction isc_transaction_info ou ITransaction.getInfo avec la balise fb_info_tra_snapshot_number ;

  • démarrer les autres transactions avec la balise isc_tpb_at_snapshot_number avec le numéro de snapshot obtenu.

L’outil d’exemple FBCSVExport, ainsi que gbak, utilise la deuxième approche. Ces approches peuvent être mélangées - par exemple, obtenir le numéro de snapshot avec SQL, puis utiliser le numéro obtenu pour démarrer d’autres transactions avec l’API, ou vice versa.

Dans FBCSVExport, nous obtenons un numéro de snapshot avec le code suivant :

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

Pour démarrer une transaction avec le numéro de snapshot, nous utilisons le code suivant :

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

Désormais, les données lues depuis différentes connexions seront cohérentes, nous pouvons donc répartir la charge entre les threads de travail.

Comment répartir exactement la charge entre les threads de travail ? Dans le cas d’une exportation complète de toutes les tables ou d’une copie de sauvegarde, l’option la plus simple sera un thread de travail par table. Mais avec cette approche, nous avons le problème suivant : s’il y a beaucoup de petites tables dans une base de données et une grande table, ou même une seule table énorme, nous ne verrons pas d’amélioration. Dans ce cas, un thread obtiendra une grande table et les autres threads resteront inactifs. Pour éviter cela, il est nécessaire de traiter une grande table par parties.

Note Le matériel ci-dessous est consacré à la lecture complète des tables. Si vous souhaitez organiser une lecture parallèle à partir d’une requête (ou d’une vue), cela nécessitera une approche légèrement différente, qui dépend des données réelles.

Diviser une grande table en parties

Supposons que nous ayons une seule grande table que nous voulons lire entièrement et le plus rapidement possible. Il est proposé de la diviser en plusieurs parties et de lire chaque partie depuis son propre flux de manière indépendante. Chaque thread doit avoir sa propre connexion à la base de données.

Dans ce cas, les questions suivantes se posent :

  • En combien de parties de traitement la table doit-elle être divisée ?

  • Quelle est la meilleure façon de diviser la table en termes d’accès aux données ?

Répondons à ces questions dans l’ordre.

En combien de parties de traitement la table doit-elle être divisée ?

Supposons le scénario idéal - le serveur et le client sont dédiés à Firebird, c’est-à-dire que tous les CPU sont entièrement à notre disposition. Alors il est recommandé :

a) Utiliser comme nombre maximal de parties parallèles le double du nombre de cœurs CPU du serveur. Pourquoi 2x cœurs ? Nous savons avec certitude qu’il y aura des délais liés aux E/S, nous pouvons donc permettre une utilisation supplémentaire du CPU. Cependant, ce nombre doit être considéré comme un réglage initial ; en pratique, il dépend des données.

b) Prendre en compte le nombre de cœurs du client : s’il y en a beaucoup plus sur le serveur (situation habituelle), il peut être judicieux de limiter davantage le nombre de parties de la partition, afin de ne pas surcharger le client (il ne pourra de toute façon pas en traiter davantage, et les coûts de commutation des flux ne disparaissent pas). Il sera possible de décider plus précisément en surveillant la charge CPU du client et du serveur - si elle est à 100 % sur le client mais nettement inférieure sur le serveur, il est judicieux de réduire le nombre de parties.

c) si le client et le serveur sont le même hôte, voir (a).

Si le client et/ou le serveur sont occupés par autre chose, vous devrez peut-être réduire le nombre de parties. Cela peut également être affecté par la capacité des disques du serveur à traiter de nombreuses requêtes d’E/S simultanément (surveillez la taille de la file d’attente et le temps de réponse).

Quelle est la meilleure façon de diviser la table en termes d’accès aux données ?

Pour mettre en œuvre un traitement parallèle efficace, il est important d’assurer une répartition uniforme des tâches entre les gestionnaires et de minimiser leur synchronisation mutuelle. De plus, il faut se rappeler que la synchronisation des gestionnaires peut se produire à la fois côté serveur et côté client. Par exemple, plusieurs gestionnaires ne doivent pas utiliser la même connexion à la base de données. Un exemple moins évident : il est mauvais que différents gestionnaires lisent des enregistrements depuis les mêmes pages de base de données. Par exemple, lorsque deux gestionnaires lisent des enregistrements pairs et impairs - ce n’est pas efficace. La synchronisation côté client peut se produire lors de la distribution des tâches, lors du traitement des données reçues (allocation de mémoire pour les résultats), etc.

L’un des problèmes du partitionnement « équitable » est que le client ne sait pas comment les enregistrements sont répartis entre les pages (et entre les clés d’index), combien d’enregistrements ou de pages de données existent (pour les grandes tables, il serait trop long de compter le nombre d’enregistrements à l’avance).

Voyons comment gbak résout ce problème.

Pour gbak, une unité de travail est un ensemble d’enregistrements provenant de pages de données (DP) appartenant à la même page de pointeurs (PP). D’une part, c’est un nombre assez important d’enregistrements pour maintenir le gestionnaire occupé sans avoir à demander fréquemment un nouveau morceau de données (synchronisation). D’autre part, même si ces ensembles d’enregistrements n’ont pas exactement la même taille, cela permettra de charger relativement uniformément les workers. C’est-à-dire qu’il est tout à fait possible qu’un worker lise N enregistrements d’une PP et qu’un autre lise M enregistrements, M étant assez différent de N. Cette approche n’est pas idéale, mais elle est assez simple à mettre en œuvre et généralement assez efficace, du moins à grande échelle (avec des dizaines ou des centaines (ou plus) de PP).

Comment obtenir le nombre de PP (Pointer Pages) pour une table donnée ? C’est assez facile et, surtout, rapide à calculer à partir de la table 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

Ensuite, nous pourrions simplement diviser le nombre de PP par le nombre de workers et donner à chaque worker sa propre partie. C’est acceptable pour le scénario où le traitement parallèle est effectué par le développeur qui connaît la distribution des données. Mais, pour un scénario plus courant, il n’y a aucune garantie que de telles « grandes » parties signifient la même quantité de travail. Nous ne souhaitons pas voir la situation où 15 workers ont terminé leur travail et restent inactifs, tandis que le 16e lit ses 10 millions d’enregistrements pendant longtemps.

C’est pourquoi gbak procède différemment. Il existe un coordinateur de travail qui attribue à chaque processeur 1 PP à la fois. Le coordinateur sait combien de PP existent au total et combien ont déjà été attribuées pour le travail. Lorsque le worker termine la lecture de ses enregistrements, il contacte le coordinateur pour obtenir un nouveau numéro de PP. Cela continue jusqu’à épuisement des PP (ou jusqu’à ce qu’il y ait des workers actifs). Bien sûr, cette interaction des workers avec le coordinateur nécessite une synchronisation. L’expérience montre que la quantité de travail donnée par une PP permet de ne pas synchroniser trop souvent. Cette approche permet de charger pratiquement uniformément tous les workers (et donc les cœurs CPU) avec du travail, quel que soit le nombre réel d’enregistrements appartenant à chaque PP.

Comment le gestionnaire lit-il les enregistrements depuis son PP ? Pour ce faire, à partir de Firebird 4.0 (apparu pour la première fois dans HQBird 2.5), il existe une fonction intégrée MAKE_DBKEY(). Grâce à elle, vous pouvez obtenir le RDB$DB_KEY (numéro physique d’enregistrement) pour le premier enregistrement sur le PP spécifié.

Et à l’aide de ces RDB$DB_KEY, les enregistrements nécessaires sont sélectionnés :

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)

Par exemple, si vous définissez loPP = 0 et hiPP = 1, alors tous les enregistrements avec PP = 0 seront lus, et uniquement à partir de celui-ci.

Maintenant que vous avez une idée de la façon dont gbak fonctionne, vous pouvez passer à une description de l’implémentation de l’utilitaire FBCSVExport.

Implémentation de l’utilitaire FBCSVExport

L’utilitaire FBCSVExport est conçu pour exporter des données depuis les tables de bases de données Firebird vers le format CSV.

Chaque table est exportée dans un fichier nommé .csv. En mode normal (mono-thread), les données des tables sont exportées séquentiellement dans l’ordre alphabétique des noms de tables.

En mode parallèle, les tables sont exportées en parallèle, chaque table dans un thread séparé. Si la table est très grande, elle est divisée en parties, et chaque partie est exportée dans un flux séparé. Pour chaque partie d’une grande table, un fichier séparé est créé avec le nom .csv.partN, où N est le numéro de la partie.

Lorsque toutes les parties d’une grande table sont exportées, les fichiers de parties sont fusionnés dans un fichier appelé .csv.

Une expression régulière est utilisée pour spécifier quelles tables seront exportées. Seules les tables régulières peuvent être exportées (les tables système, GTT, vues, tables externes ne sont pas prises en charge). Les expressions régulières doivent être en syntaxe SQL, c’est-à-dire celles utilisées dans le prédicat SIMILAR TO.

Pour sélectionner une liste de tables exportées, ainsi qu’une liste de leurs PP en mode multi-thread, nous utilisons la requête suivante :

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

En mode mono-thread, cette requête peut être simplifiée en

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

En mode mono-thread, les valeurs des champs PAGE_SEQUENCE et PP_CNT ne sont pas utilisées ; elles sont ajoutées à la requête pour unifier les messages de sortie.

Le résultat de cette requête est formé dans un vecteur de structures :

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

Ce vecteur est rempli à l’aide d’une fonction déclarée comme :

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

Le dernier paramètre singleWorker bascule le mode de remplissage de std::vector : si singleWorker = true, alors la requête pour le mode mono-thread est utilisée ; si singleWorker = false, alors une requête plus coûteuse et complexe est utilisée pour le mode multi-thread. Je ne donnerai pas l’implémentation elle-même, elle est assez simple et vous pouvez la voir dans le code source du projet.

Pour exporter une table au format CSV, la classe CSVExportTable a été développée, qui contient les méthodes suivantes :

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

La méthode prepare est destinée à construire et préparer une requête utilisée pour exporter une table au format CSV. La requête interne est construite différemment selon le paramètre withDbkeyFilter. Si withDbkeyFilter = true, alors la requête est construite avec un filtrage par plage 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, ?)

sinon, une requête simplifiée est utilisée :

sql
SELECT *
FROM tableName

La valeur du paramètre withDbkeyFilter est définie sur true si le mode multi-thread est utilisé et que la table est grande. Nous considérons que la table est grande si pp_cnt > 1.

La méthode printHeader est destinée à imprimer l’en-tête d’un fichier CSV (noms des colonnes de la table).

La méthode printData imprime les données de la table dans un fichier CSV à partir du numéro de page PP ppNum, si la requête a été préparée avec un filtre par plage RDB$DB_KEY, et toutes les données de la table sinon.

Examinons maintenant le code du mode mono-thread

cpp
...

// Ouverture de la connexion principale
Firebird::AutoRelease att(
    provider->attachDatabase(
        &status,
        m_database.c_str(),
        dbpLength,
        dpb
    )
);

// Démarrage de la transaction principale en mode d'isolation 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)
    )
);
// Obtention d'une liste de tables à l'aide de l'expression régulière dans m_filter.
// m_parallel définit le nombre de threads parallèles ; lorsqu'il est égal à 1,
// une requête simplifiée est utilisée pour obtenir la liste des tables,
// sinon, une liste de PP et leur nombre est générée pour chaque table.
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) {
        // il n'y a aucun intérêt à utiliser un filtre de plage RDB$DB_KEY ici
        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);
    }
}

Tout ici est assez simple et ne nécessite pas d’explication supplémentaire, passons donc à la partie multi-thread.

Pour que l’exportation se fasse en mode multi-thread, il est nécessaire de créer m_parallel - 1 threads de travail supplémentaires. Pourquoi le nombre de threads supplémentaires est-il inférieur de 1 ? Oui, parce que le thread principal exportera également des données et il est égal aux threads supplémentaires. Déplaçons la partie commune du flux principal et des flux supplémentaires dans une fonction séparée :

cpp
void ExportApp::exportByTableDesc(Firebird::ThrowStatusWrapper* status, FBExport::CSVExportTable& csvExport, const TableDesc& tableDesc)
{
    // Si tableDesc a pp_cnt > 1, alors il décrit seulement une partie de la table, et il est nécessaire de construire
    // une requête en utilisant un filtre par plage 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";

    // Si ce n'est pas la première partie de la table, alors écrire cette partie dans le fichier .csv.part, où
    // N - numéro PP. Plus tard, les parties de la table seront combinées dans un seul fichier .csv
    if (tableDesc.page_sequence > 0) {
        fileName += ".part_" + std::to_string(tableDesc.page_sequence);
    }
    csv::CSVFile csv(m_outputDir / fileName);
    // L'en-tête du fichier CSV doit être imprimé uniquement dans la première partie de la table.
    if (tableDesc.page_sequence == 0 && m_printHeader) {
        csvExport.printHeader(status, csv);
    }
    csvExport.printData(status, csv, tableDesc.page_sequence);
}

Les descriptions des tables ou de leurs parties sont situées dans un vecteur commun avec les structures TableDesc. À partir de ce vecteur, chaque thread de travail prend une table ou la partie suivante. Pour éviter les courses de données, il est nécessaire de synchroniser l’accès à la ressource partagée. Mais std::vector lui-même ne change pas, donc vous pouvez uniquement synchroniser la variable partagée, qui est l’index dans ce vecteur. Cela peut être facilement fait en utilisant std::atomic comme telle variable.

cpp
if (m_parallel == 1) {
    ...
}
else {
    // Détermination du nombre de threads de travail supplémentaires
    const auto workerCount = m_parallel - 1;

    // Obtention du numéro de snapshot depuis la transaction principale
    auto snapshotNumber = getSnapshotNumber(&status, tra);
    // variable pour stocker l'exception dans le thread
    std::exception_ptr exceptionPointer = nullptr;
    std::mutex m;
	// compteur atomique
    // est l'index de la table suivante ou de sa partie
    std::atomic<size_t> counter = 0;
    // pool de threads de travail
    std::vector<std::thread> thread_pool;
    thread_pool.reserve(workerCount);
    for (int i = 0; i < workerCount; i++) {
		// pour chaque thread, nous créons notre propre connexion
        Firebird::AutoRelease workerAtt(
            provider->attachDatabase(
                &status,
                m_database.c_str(),
                dbpLength,
                dpb
            )
        );
		// et notre transaction à laquelle nous passons le numéro de snapshot
        // pour créer un snapshot partagé
        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)
            )
        );
		// création d'un thread
        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) {
                    // incrémenter le compteur atomique
                    size_t localCounter = counter++;
					// si les tables ou leurs parties sont épuisées, sortir
                    // de la boucle infinie et terminer le thread
                    if (localCounter >= tables.size())
                        break;
                    // obtenir une description de la table ou de sa partie
                    const auto& tableDesc = tables[localCounter];
                    // et faire l'exportation
                    exportByTableDesc(&status, csvExport, tableDesc);
                }
                if (tra) {
                    tra->commit(&status);
                    tra.release();
                }

                if (att) {
                    att->detach(&status);
                    att.release();
                }
            }
            catch (...) {
				// si une exception se produit, la sauvegarder pour
                // une libération ultérieure dans le thread principal
                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);
    }
    ...

Il ne reste plus qu’à combiner les fichiers créés pour les parties des tables en un seul fichier pour chacune de ces tables.

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

Mesurons les performances de l’outil en mode mono-thread et multi-thread.

Benchmark de l’outil FBCSVExport

Tout d’abord, examinons les résultats de la comparaison des modes d’exportation multi-thread et mono-thread sur un ordinateur domestique modéré. === Windows

  • Système d’exploitation : Windows 10 x64.

  • Processeur : Intel Core i3 8100, 4 cœurs, 4 threads.

  • Mémoire : 16 Go

  • Sous-système de disque : SSD NVME (base de données), SSD SATA (dossier de stockage des fichiers CSV).

  • Firebird 4.0.4 x64

Résultats :

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

D’après les résultats des tests, il est clair qu’avec deux threads, l’accélération était de 1,8 fois, ce qui est un bon résultat. Mais l’exécution parallèle de l’exportation avec 4 threads a également montré une amélioration de 1,8 fois. Pourquoi pas 3-4 ? Le fait est que le serveur Firebird et l’utilitaire d’exportation s’exécutent sur le même ordinateur, qui ne dispose que de 4 cœurs. Ainsi, le serveur Firebird lui-même utilise 4 threads pour lire la table, et l’utilitaire FBCSVExport utilise également 4 threads. Évidemment, dans ce cas, il est assez difficile d’obtenir une accélération de plus de 2 fois. Par conséquent, nous essaierons sur un autre matériel, où le nombre de cœurs est nettement plus important.

Linux

  • Système d’exploitation : CentOS 8.

  • Processeur : 2 processeurs Intel Xeon E5-2603 v4, total 12 cœurs, 12 threads.

  • Mémoire : 32 Go

  • Sous-système de disque : disque dur SAS (RAID 10)

  • Firebird 4.0.4 x64

Résultats :

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

Dans ce cas, le nombre optimal de threads pour l’exportation est de 6 (6 threads pour Firebird et 6 threads pour l’utilitaire FBCSVExport). En même temps, nous avons réussi à obtenir une accélération de 5 fois, ce qui indique une assez bonne évolutivité. Sur le serveur Linux et l’ordinateur Windows, nous avons utilisé des bases de données identiques, et vous avez probablement remarqué que l’exportation mono-thread sur Windows était presque 2 fois plus rapide : cela est dû à un sous-système de disque plus rapide (le lecteur NVME est beaucoup plus rapide que les disques SAS combinés en RAID).

Résumé

Dans cet article, nous avons examiné comment lire efficacement les données des tables du SGBD Firebird en utilisant le parallélisme. De plus, l’exemple a montré comment utiliser certaines des capacités du SGBD Firebird pour organiser une telle lecture dans vos logiciels.

Un grand merci à Vladislav Khorsun, développeur principal de Firebird, pour son aide sur ce sujet.

Pour toute question ou commentaire, veuillez envoyer un e-mail à [email protected].