Esta página fue traducida automáticamente. Lee el original en inglés. English

Biblioteca de IBSurgeon

Lectura paralela de datos en Firebird

D.Simonov, V.Horsun

versión 1.0.5 del 05.12.2023

Este material está patrocinado y creado con el patrocinio y apoyo de IBSurgeon www.ib-aid.com, proveedor de HQbird (distribución avanzada de Firebird) y proveedor de servicios de optimización de rendimiento, migración y soporte técnico para Firebird.

El material está licenciado bajo la Licencia de Documentación Pública https://www.firebirdsql.org/file/documentation/html/en/licenses/pdl/public-documentation-license.html

Materiales relacionados:

Prefacio

Firebird 5.0 introdujo la capacidad de usar paralelismo al crear una copia de seguridad utilizando la utilidad gbak, entre otras funciones paralelas. Inicialmente, esta función apareció en HQbird 2.5, luego en HQbird 3.0 y 4.0, y posteriormente fue portada a Firebird 5.0.

En este artículo consideraremos las funciones que se utilizan al crear copias de seguridad paralelas dentro de la utilidad gbak. También mostraremos cómo se pueden usar en sus aplicaciones para la lectura paralela de datos.

Es importante señalar que aquí no estamos hablando de escaneo paralelo de tablas dentro del motor Firebird al ejecutar consultas SQL, sino de leer datos dentro de su aplicación en flujos paralelos.

Herramienta de ejemplo FBCSVExport

Para demostrar la lectura paralela de datos del DBMS Firebird, se escribió una utilidad de ejemplo que exporta datos de una o más tablas a formato CSV.

Su descripción y su código fuente abierto están aquí: https://github.com/IBSurgeon/FBCSVExport.git

Como puede ver en la descripción de la utilidad y en el artículo a continuación, mediante el procesamiento paralelo es posible exportar datos y realizar otras operaciones paralelas 2-10 veces más rápido que en 1 hilo (dependiendo del hardware).

Para cualquier pregunta, contacte con [email protected].

Lectura paralela

Pensemos en cómo leer datos de varias tablas en paralelo. Como sabe, Firebird permite ejecutar consultas en paralelo solo si cada consulta se ejecuta en una conexión separada.

Creemos un grupo de hilos de trabajo. El hilo principal de la aplicación también es un hilo de trabajo, por lo que el número de hilos de trabajo adicionales debe ser N - 1, donde N es el número total de trabajadores paralelos. Cada hilo de trabajo ejecutará su propia conexión y transacción.

El primer problema: ¿cómo garantizar la consistencia de los datos leídos?

Lectura de datos consistente

Dado que cada hilo de trabajo usa su propia conexión y su propia transacción, surge el problema de lecturas inconsistentes: si la tabla es modificada simultáneamente por otros usuarios, los datos leídos pueden ser inconsistentes. En modo de un solo hilo, gbak usa una transacción con modo de aislamiento SNAPSHOT, lo que permite leer información consistente al inicio de la transacción SNAPSHOT. Pero aquí tenemos múltiples transacciones y es necesario que vean el mismo “snapshot” para que lean los mismos datos inmutables.

El mecanismo para crear un snapshot compartido para diferentes transacciones con el modo de aislamiento SNAPSHOT se introdujo en Firebird 4.0 (originalmente en HQBird 2.5, pero en Firebird 4.0/HQBird 4.0 es más simple y eficiente). Hay dos formas de crear un snapshot compartido:

  1. Con SQL
  • obtener el número de snapshot de la transacción principal (que se inicia en el hilo de trabajo principal).
sql
SELECT RDB$GET_CONTEXT('SYSTEM', 'SNAPSHOT_NUMBER') FROM RDB$DATABASE
  • iniciar otras transacciones con el siguiente SQL:
sql
SET TRANSACTION SNAPSHOT AT NUMBER snapshot_number

donde snapshot_number es el número recuperado por la consulta anterior.

  1. Con API
  • obtener el número de snapshot de la transacción principal (que se inicia en el hilo de trabajo principal) con la función isc_transaction_info o ITransaction.getInfo con la etiqueta fb_info_tra_snapshot_number;
  • iniciar otras transacciones con la etiqueta isc_tpb_at_snapshot_number con un número de snapshot obtenido.

La herramienta de ejemplo FBCSVExport, así como gbak, utiliza el segundo enfoque. Estos enfoques se pueden mezclar: por ejemplo, obtener el número de snapshot con SQL y usar el número de snapshot obtenido para iniciar otras transacciones con API, o viceversa.

En FBCSVExport obtenemos un número de snapshot con el siguiente código:

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

Para iniciar la transacción con el número de snapshot usamos el siguiente código:

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

Ahora los datos leídos desde diferentes conexiones serán consistentes, por lo que podemos distribuir la carga entre los hilos de trabajo.

¿Cómo distribuir exactamente la carga entre los hilos de trabajo? En caso de exportación completa de todas las tablas o una copia de seguridad, la opción más simple será un hilo de trabajo por tabla. Pero con este enfoque tenemos el siguiente problema: si hay muchas tablas pequeñas en una base de datos y una tabla grande, o incluso solo una tabla y es enorme, no veremos la mejora. En este caso, algún hilo obtendrá una tabla grande y los hilos restantes estarán inactivos. Para evitar que esto suceda, es necesario procesar una tabla grande en partes.

Nota El material a continuación está dedicado a la lectura completa de tablas; si desea organizar la lectura paralela desde alguna consulta (o vista), se requerirá un enfoque ligeramente diferente, que depende de los datos reales.

Dividir una tabla grande en partes

Digamos que tenemos solo una tabla grande que queremos leer en su totalidad y lo más rápido posible. Se propone dividirla en varias partes y leer cada parte desde su propio flujo de forma independiente. Cada hilo debe tener su propia conexión a la base de datos.

En este caso, surgen las siguientes preguntas:

  • ¿En cuántas partes de procesamiento se debe dividir la tabla?

  • ¿Cuál es la mejor manera de dividir la tabla en términos de acceso a datos?

Respondamos estas preguntas en orden.

¿En cuántas partes de procesamiento se debe dividir la tabla?

Supongamos el escenario ideal: el servidor y el cliente están dedicados a Firebird, es decir, todas las CPU están completamente a nuestra disposición. Entonces se recomienda:

a) Usar como número máximo de partes paralelas el doble del número de núcleos de CPU en el servidor. ¿Por qué 2x núcleos? Sabemos con certeza que habrá retrasos asociados con IO, por lo que podemos permitir un uso adicional de CPU. Sin embargo, este número debe considerarse como una configuración inicial; prácticamente depende de los datos.

b) Tener en cuenta el número de núcleos en el cliente: si hay muchos más en el servidor (situación habitual), entonces podría tener sentido limitar aún más el número de partes de la partición, para no sobrecargar al cliente (de todos modos no podrá procesar más, y los costos de cambio de flujos no desaparecen). Se podrá decidir con mayor precisión monitoreando la carga de CPU del cliente y del servidor: si está al 100% en el cliente, pero notablemente menos en el servidor, entonces tiene sentido reducir el número de partes.

c) si el cliente y el servidor son el mismo host, entonces vea (a).

Si el cliente y/o el servidor están ocupados con otra cosa, es posible que deba reducir el número de partes. Esto también puede verse afectado por la capacidad de los discos en el servidor para procesar muchas solicitudes IO simultáneamente (monitoree el tamaño de la cola y el tiempo de respuesta).

¿Cuál es la mejor manera de dividir la tabla en términos de acceso a datos?

Para implementar un procesamiento paralelo efectivo, es importante garantizar una distribución uniforme de trabajos entre los manejadores y minimizar su sincronización mutua. Además, debe recordar que la sincronización de manejadores puede ocurrir tanto en el lado del servidor como en el lado del cliente. Por ejemplo, varios manejadores no deben usar la misma conexión a la base de datos. Un ejemplo menos obvio: es malo si diferentes manejadores leen registros de las mismas páginas de base de datos. Por ejemplo, cuando dos manejadores leen registros pares e impares, no es efectivo. La sincronización en el cliente puede ocurrir durante la distribución de tareas, durante el procesamiento de los datos recibidos (asignación de memoria para resultados), y así sucesivamente.

Uno de los problemas con la partición “justa” es que el cliente no sabe cómo están distribuidos los registros entre las páginas (y entre las claves de índice), cuántos registros o páginas de datos hay (para tablas grandes, contar el número de registros de antemano tomaría demasiado tiempo).

Veamos cómo gbak resuelve este problema.

Para gbak, una unidad de trabajo es un conjunto de registros de páginas de datos (DP) que pertenecen a la misma página de puntero (PP). Por un lado, es un número bastante grande de registros para mantener al manejador ocupado sin tener que solicitar con frecuencia un nuevo fragmento de datos (sincronización). Por otro lado, incluso si estos conjuntos de registros no tienen exactamente el mismo tamaño, permitirá cargar a los trabajadores de manera relativamente uniforme. Es decir, es bastante posible que un trabajador lea N registros de una PP, y el otro M registros, y M sea bastante diferente de N. Este enfoque no es ideal, pero es bastante simple de implementar y generalmente es bastante efectivo, al menos a gran escala (con decenas o cientos (o más) PP).

¿Cómo obtener el número de PP (Páginas de Puntero) para una tabla determinada? Es bastante fácil y, lo más importante, rápido calcularlo desde la tabla 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

A continuación, podríamos simplemente dividir el número de PP por el número de trabajadores y dar a cada trabajador su propia parte. Esto está bien para el escenario en el que el procesamiento paralelo lo realiza el desarrollador que conoce la distribución de datos. Pero, para el escenario más común, no hay garantía de que tales “grandes” partes signifiquen la misma cantidad de trabajo. No nos interesa ver la situación en la que 15 trabajadores terminaron su trabajo y están inactivos, y el 16º lee sus 10M de registros durante mucho tiempo.

Por eso gbak lo hace de manera diferente. Hay un coordinador de trabajo que asigna a cada procesador 1 PP a la vez. El coordinador sabe cuántas PP hay en total y cuántas ya se han emitido para trabajar. Cuando el trabajador completa la lectura de sus registros, contacta al coordinador para obtener un nuevo número de PP. Continúa hasta que las PP se agotan (o hay trabajadores activos). Por supuesto, tal interacción de los trabajadores con el coordinador requiere sincronización. La experiencia muestra que la cantidad de trabajo dada por una PP permite no sincronizar con demasiada frecuencia. Este enfoque permite cargar prácticamente de manera uniforme a todos los trabajadores (y por lo tanto los núcleos de CPU) con trabajo, independientemente del número real de registros que pertenecen a cada PP.

¿Cómo lee el handler los registros de su PP? Para ello, a partir de Firebird 4.0 (apareció por primera vez en HQBird 2.5) existe una función incorporada MAKE_DBKEY(). Con su ayuda, se puede obtener el RDB$DB_KEY (número de registro físico) para el primer registro en el PP especificado.

Y con la ayuda de estos RDB$DB_KEY se seleccionan los registros necesarios:

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)

Por ejemplo, si establece loPP = 0 y hiPP = 1, entonces se leerán todos los registros con PP = 0, y solo de ese.

Ahora que tiene una idea de cómo funciona gbak, puede pasar a una descripción de la implementación de la utilidad FBCSVExport.

Implementación de la utilidad FBCSVExport

La utilidad FBCSVExport está diseñada para exportar datos de tablas de bases de datos Firebird a formato CSV.

Cada tabla se exporta a un archivo llamado .csv. En modo normal (de un solo hilo), los datos de las tablas se exportan secuencialmente en orden alfabético de los nombres de las tablas.

En modo paralelo, las tablas se exportan en paralelo, cada tabla en un hilo separado. Si la tabla es muy grande, se divide en partes, y cada parte se exporta en un flujo separado. Para cada parte de una tabla grande, se crea un archivo separado con el nombre .csv.partN, donde N es el número de parte.

Cuando todas las partes de una tabla grande se exportan, los archivos de partes se fusionan en un archivo llamado .csv.

Se utiliza una expresión regular para especificar qué tablas se exportarán. Solo se pueden exportar tablas regulares (las tablas del sistema, GTT, vistas y tablas externas no son compatibles). Las expresiones regulares deben estar en sintaxis SQL, es decir, las que se usan en el predicado SIMILAR TO.

Para seleccionar una lista de tablas exportadas, así como una lista de sus PP en modo de subprocesos múltiples, usamos la siguiente consulta:

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 modo de un solo hilo, esta consulta se puede simplificar a

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 modo de un solo hilo, los valores de los campos PAGE_SEQUENCE y PP_CNT no se utilizan; se agregan a la solicitud para unificar los mensajes de salida.

El resultado de esta consulta se forma en un vector de estructuras:

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

Este vector se llena usando una función declarada como:

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

El último parámetro singleWorker cambia el modo de llenado de std::vector; si singleWorker = true, se usa la solicitud para modo de un solo hilo; si singleWorker = false, se usa una consulta más costosa y compleja para modo de subprocesos múltiples. No daré la implementación en sí, es bastante simple y puede verla en el código fuente del proyecto.

Para exportar una tabla a formato CSV, se ha desarrollado la clase CSVExportTable, que contiene los siguientes métodos:

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

El método prepare está diseñado para construir y preparar una consulta que se usa para exportar una tabla a formato CSV. La consulta interna se construye de manera diferente según el parámetro withDbkeyFilter. Si withDbkeyFilter = true, entonces la consulta se construye con filtrado por el rango 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, ?)

de lo contrario, se usa una consulta simplificada:

sql
SELECT *
FROM tableName

El valor del parámetro withDbkeyFilter se establece en true si se usa el modo de subprocesos múltiples y la tabla es grande. Consideramos que la tabla es grande si pp_cnt > 1.

El método printHeader está diseñado para imprimir el encabezado de un archivo CSV (nombres de columnas de la tabla).

El método printData imprime datos de la tabla en un archivo CSV desde el número de página PP ppNum, si la solicitud se preparó usando un filtro por rango RDB$DB_KEY, y todos los datos de la tabla en caso contrario.

Ahora veamos el código para el modo de un solo hilo

cpp
...

// Abriendo la conexión principal
Firebird::AutoRelease att(
    provider->attachDatabase(
        &status,
        m_database.c_str(),
        dbpLength,
        dpb
    )
);

// Iniciando la transacción principal en el modo de aislamiento 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)
    )
);
// Obtener una lista de tablas usando la expresión regular en m_filter.
// m_parallel establece el número de hilos paralelos; cuando es igual a 1,
// se usa una consulta simplificada para obtener la lista de tablas,
// de lo contrario, se genera una lista de PP y su número para cada tabla.
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) {
        // no tiene sentido usar un filtro de rango RDB$DB_KEY aquí
        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);
    }
}

Aquí todo es bastante simple y no requiere explicación adicional, así que pasemos a la parte de subprocesos múltiples.

Para que la exportación ocurra en modo de subprocesos múltiples, es necesario crear m_parallel - 1 hilos de trabajo adicionales. ¿Por qué el número de hilos adicionales es 1 menos? Sí, porque el hilo principal también exportará datos y es igual a los hilos adicionales. Movamos la parte común del flujo principal y adicional a una función separada:

cpp
void ExportApp::exportByTableDesc(Firebird::ThrowStatusWrapper* status, FBExport::CSVExportTable& csvExport, const TableDesc& tableDesc)
{
    // Si tableDesc tiene pp_cnt > 1, entonces describe solo una parte de la tabla, y es necesario construir
    // la consulta usando un filtro por rango 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 esta no es la primera parte de la tabla, escriba esta parte en el archivo .csv.part, donde
    // N - número de PP. Más tarde, las partes de la tabla se combinarán en un solo archivo .csv
    if (tableDesc.page_sequence > 0) {
        fileName += ".part_" + std::to_string(tableDesc.page_sequence);
    }
    csv::CSVFile csv(m_outputDir / fileName);
    // El encabezado del archivo CSV debe imprimirse solo en la primera parte de la tabla.
    if (tableDesc.page_sequence == 0 && m_printHeader) {
        csvExport.printHeader(status, csv);
    }
    csvExport.printData(status, csv, tableDesc.page_sequence);
}

Las descripciones de tablas o sus partes se encuentran en un vector común con estructuras TableDesc. De este vector, cada hilo de trabajo toma una tabla o la siguiente parte. Para evitar condiciones de carrera, es necesario sincronizar el acceso al recurso compartido. Pero std::vector en sí no cambia, por lo que solo se puede sincronizar la variable compartida, que es el índice en este vector. Esto se puede hacer fácilmente usando std::atomic como tal variable.

cpp
if (m_parallel == 1) {
    ...
}
else {
    // Determinando el número de hilos de trabajo adicionales
    const auto workerCount = m_parallel - 1;

    // Obteniendo el número de instantánea de la transacción principal
    auto snapshotNumber = getSnapshotNumber(&status, tra);
    // variable para almacenar la excepción dentro del hilo
    std::exception_ptr exceptionPointer = nullptr;
    std::mutex m;
	// contador atómico
    // es el índice de la siguiente tabla o parte de ella
    std::atomic<size_t> counter = 0;
    // grupo de hilos de trabajo
    std::vector<std::thread> thread_pool;
    thread_pool.reserve(workerCount);
    for (int i = 0; i < workerCount; i++) {
		// para cada hilo creamos nuestra propia conexión
        Firebird::AutoRelease workerAtt(
            provider->attachDatabase(
                &status,
                m_database.c_str(),
                dbpLength,
                dpb
            )
        );
		// y nuestra transacción a la que pasamos el número de instantánea
        // para crear una instantánea compartida
        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)
            )
        );
		// crear un hilo
        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) {
                    // incrementar el contador atómico
                    size_t localCounter = counter++;
					// si las tablas o sus partes se han terminado, salir
                    // del bucle infinito y terminar el hilo
                    if (localCounter >= tables.size())
                        break;
                    // obtener una descripción de la tabla o parte de ella
                    const auto& tableDesc = tables[localCounter];
                    // y hacer la exportación
                    exportByTableDesc(&status, csvExport, tableDesc);
                }
                if (tra) {
                    tra->commit(&status);
                    tra.release();
                }

                if (att) {
                    att->detach(&status);
                    att.release();
                }
            }
            catch (...) {
				// si ocurre una excepción, guárdela para
                // liberación posterior en el hilo 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);
    }
    ...

Todo lo que queda es combinar los archivos que se crearon para las partes de las tablas en un solo archivo para cada una de estas tablas.

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

Midamos el rendimiento de la herramienta en modo de un solo hilo y en modo de múltiples hilos.

Prueba de rendimiento de la herramienta FBCSVExport

Primero, veamos los resultados de comparar los modos de exportación de múltiples hilos y de un solo hilo en una computadora doméstica moderada. === Windows

  • Sistema operativo: Windows 10 x64.

  • Procesador: Intel Core i3 8100, 4 núcleos, 4 hilos.

  • Memoria: 16 GB

  • Subsistema de disco: NVME SSD (base de datos), SATA SSD (carpeta para almacenar archivos CSV).

  • Firebird 4.0.4 x64

Resultados:

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

De los resultados de la prueba queda claro que al usar dos hilos, la aceleración fue de 1.8 veces, lo cual es un buen resultado. Pero la ejecución paralela de la exportación en 4 hilos también mostró una mejora de 1.8 veces. ¿Por qué no 3-4? El hecho es que el servidor Firebird y la utilidad de exportación se ejecutan en la misma computadora, que solo tiene 4 núcleos. Así, el servidor Firebird usa 4 hilos para leer la tabla, y la utilidad FBCSVExport también usa 4 hilos. Obviamente, en este caso es bastante difícil lograr una aceleración de más de 2 veces. Por lo tanto, probaremos en otro hardware, donde el número de núcleos es significativamente mayor.

Linux

  • Sistema operativo: CentOS 8.

  • Procesador: 2 procesadores Intel Xeon E5-2603 v4, total 12 núcleos, 12 hilos.

  • Memoria: 32 GB

  • Subsistema de disco: SAS HDD (RAID 10)

  • Firebird 4.0.4 x64

Resultados:

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

En este caso, el número óptimo de hilos para la exportación es 6 (6 hilos para Firebird y 6 hilos para la utilidad FBCSVExport). Al mismo tiempo, logramos alcanzar una aceleración de 5 veces, lo que indica una escalabilidad bastante buena. En el servidor Linux y la computadora Windows usamos bases de datos idénticas, y probablemente notaste que la exportación de un solo hilo en Windows fue casi 2 veces más rápida: se debe a un subsistema de disco más rápido (la unidad NVME es mucho más rápida que las unidades SAS combinadas en RAID).

Resumen

En este artículo, consideramos cómo leer eficazmente datos de las tablas del DBMS Firebird usando paralelismo. También se mostró un ejemplo de cómo puedes usar algunas de las capacidades del DBMS Firebird para organizar dicha lectura en tu software.

Muchas gracias a Vladislav Khorsun, desarrollador principal de Firebird, por su ayuda con este material.

Para cualquier pregunta o comentario, por favor escribe a [email protected].