Programación concurrente (IX): Generación dinámica de tareas

1. Introducción

Imagen generada con inteligencia artificial. Uso con fines divulgativos
En este artículo revisitaremos el problema planteado en el post anterior de esta serie. Dada la ruta (path) de un directorio (directory), deseamos generar un resumen estadístico del contenido del árbol de directorios: para cada extensión de archivo encontrada (.txt, .zip, etc.), obtendremos el número de ficheros y el tamaño total acumulado. El análisis proporcionará además el número total de subdirectorios descubiertos durante el recorrido.

La solución presentada en el artículo original, aunque sencilla y funcional, adolecía de una limitación importante: requería recorrer secuencialmente todo el árbol mediante un iterador std::filesystem::recursive_directory_iterator [1] para construir un vector con las rutas de todos los archivos y subdirectorios encontrados. Sólo una vez completada esa fase preliminar era posible distribuir el trabajo de análisis entre varios hilos, que procesaban de forma paralela las rutas recopiladas según un esquema fork-join.

En esta nueva implementación, más avanzada, paralelizaremos el propio recorrido del árbol de directorios. Cada hilo trabajador irá procesando directorios extraídos de una cola concurrente compartida, descubrirá nuevos subdirectorios contenidos en ellos y los incorporará dinámicamente a la cola para que puedan ser procesados posteriormente por cualquier otro trabajador. De este modo, múltiples hilos podrán explorar simultáneamente distintas ramas del árbol y acumular estadísticas parciales que se fusionarán finalmente en un único resultado global.

Este análisis permitirá generar, dado un directorio raíz, desgloses como el del ejemplo siguiente:
 .docx:  16 files (147.938.289 bytes)
  .pdf:  28 files (72.191.677 bytes)
  .png:   1 files (90.575 bytes)
 .xlsx:   1 files (88.162 bytes)
_____________________________________
contains: 46 files, 5 folders
size: 210.1 MiB (220.308.703 bytes)
Para implementar nuestra solución, emplearemos el estándar C++26 (particularmente, contratos y std::optional<T&>) y el sistema de módulos introducido en C++20 (lo que nos permitirá organizar el código en componentes bien definidos). Esta elección responde a criterios de simplicidad y modernidad del código; de ser necesario, su adaptación al esquema clásico basado en archivos de cabecera (.hpp) y de implementación (.cpp) resultaría inmediata.

2. Módulo concurrency_tools.dynamic_task_queue

En primer lugar, implementaremos una cola de trabajo concurrente bloqueante dynamic_task_queue<T>, específicamente diseñada para situaciones en las que las tareas puedan generar dinámicamente nuevas tareas durante su ejecución. Su empleo simplificará la paralelización de recorridos de grafos, como las jerarquías de directorios que nos ocupan en este artículo.

La clase almacenará las tareas pendientes de tipo T en una cola privada estándar std::queue<T> de nombre tasks_. Asimismo, mantendrá un contador de tareas actualmente en proceso denominado active_, cuyo valor se incrementará al adquirir una tarea y se decrementará al completarla. Ello permitirá detectar automáticamente la finalización global del proceso cuando no queden elementos en la cola ni tareas en ejecución. En tal caso, se cumplirá la condición tasks_.empty() and active_ == 0. Tanto la cola como el contador estarán protegidos por un objeto std::mutex, mientras que una variable de condición std::condition_variable permitirá bloquear los trabajadores cuando no exista trabajo disponible y despertarlos cuando se incorporen nuevas tareas o se alcance la finalización global.
Es importante observar que una cola vacía no implica necesariamente que la computación haya concluido, ya que algún trabajador podría estar procesando una tarea y generar nuevas tareas posteriormente. El contador active_ permite, precisamente, distinguir entre ambas situaciones.
La interfaz de la cola se basará en un protocolo acquire/complete, donde cada adquisición de una tarea por parte de un trabajador a través de la función miembro acquire() deberá emparejarse necesariamente con una llamada posterior a la función complete(). Esta última podrá registrar nuevas tareas derivadas. Concretamente:
  • acquire(): si existe trabajo disponible, extrae una tarea e incrementa internamente el número de tareas activas active_ en ejecución. Si no hay tareas pendientes pero todavía existen trabajadores procesando otras tareas, la llamada queda bloqueada hasta que aparezca nuevo trabajo o la computación finalice. Cuando no quedan tareas pendientes ni activas, devuelve un resultado vacío (std::nullopt) para indicar que la exploración ha concluido y que los trabajadores pueden finalizar.
  • complete() realiza dos acciones de forma atómica: (i) registra la finalización de una tarea previamente adquirida mediante acquire() reduciendo el contador de tareas activas active_ y (ii) incorpora nuevas tareas descubiertas a la cola. Si tras completar la operación no quedan tareas pendientes ni activas, notifica la finalización global de la computación para que todos los trabajadores puedan terminar. Si hay trabajo pendiente por realizar (tasks_.empty() == false), despierta a un trabajador bloqueado.
Para evitar la invocación manual de complete() por parte del programador y garantizar la correcta finalización de las tareas incluso en presencia de excepciones, la función acquire() no devuelve directamente un objeto de tipo T, sino un valor opcional std::optional de una clase handle auxiliar acquired_task. Esta clase actúa como un objeto RAII que representa una tarea actualmente en proceso. Además del valor adquirido, almacena un vector con las nuevas tareas descubiertas durante su procesamiento. Cuando el objeto abandona su ámbito de definición, su destructor invoca automáticamente complete(), transfiriendo a la cola las nuevas tareas acumuladas y notificando la finalización de la tarea original. De este modo, el protocolo acquire/complete queda garantizado por construcción mediante la técnica RAII:

export module concurrency_tools.dynamic_task_queue; import std; export namespace concurrency_tools { // cola de trabajo concurrente bloqueante para algoritmos en // los que las tareas pueden generar dinámicamente nuevas tareas: template<typename T> requires std::is_nothrow_move_constructible_v<T> class dynamic_task_queue { public: class acquired_task { public: acquired_task( dynamic_task_queue<T>& queue, T&& value ) noexcept : queue_opt_{queue}, value_{std::move(value)}, new_tasks_{} { } acquired_task(acquired_task const&) = delete; auto operator=(acquired_task const&) -> acquired_task& = delete; acquired_task(acquired_task&& other) noexcept : queue_opt_{other.queue_opt_}, value_{std::move(other.value_)}, new_tasks_{std::move(other.new_tasks_)} { other.queue_opt_ = std::nullopt; } auto operator=(acquired_task&&) -> acquired_task& = delete; auto get() const noexcept -> T const& { return value_; } auto get() noexcept -> T& { return value_; } template<typename S> requires std::constructible_from<T, S&&> void add_task(S&& task) { new_tasks_.emplace_back(std::forward<S>(task)); } ~acquired_task() { if (queue_opt_) { queue_opt_->complete(std::move(new_tasks_)); } } private: std::optional<dynamic_task_queue<T>&> queue_opt_; T value_; std::vector<T> new_tasks_; }; // -------------------------------------------------- template<typename S> requires std::constructible_from<T, S&&> explicit dynamic_task_queue(S&& initial) { tasks_.emplace(std::forward<S>(initial)); } dynamic_task_queue(dynamic_task_queue const&) = delete; auto operator=(dynamic_task_queue const&) -> dynamic_task_queue& = delete; [[nodiscard]] auto acquire() -> std::optional<acquired_task> { auto lock = std::unique_lock{mtx_}; cv_.wait(lock, [this]{ return active_ == 0 or not tasks_.empty(); }); if (tasks_.empty()) { // ⇔ finished() == true return std::nullopt; } auto res = std::optional<acquired_task>{ std::in_place, *this, std::move(tasks_.front()) }; tasks_.pop(); ++active_; return res; } [[nodiscard]] auto empty() const -> bool { auto lock = std::lock_guard{mtx_}; return tasks_.empty(); } [[nodiscard]] auto active() const -> std::size_t { auto lock = std::lock_guard{mtx_}; return active_; } [[nodiscard]] auto finished() const -> bool { auto lock = std::lock_guard{mtx_}; return tasks_.empty() and active_ == 0; } private: std::queue<T> tasks_; std::size_t active_ = 0; mutable std::mutex mtx_; std::condition_variable cv_; // nota sobre noexcept: un fallo durante la incorporación de nuevas // tareas se considera irrecuperable y provoca std::terminate() template<typename S> requires std::is_nothrow_constructible_v<T, S&&> void complete(std::vector<S>&& new_tasks) noexcept { auto globally_finished = false; auto has_pending_work = false; { auto lock = std::lock_guard{mtx_}; contract_assert(active_ != 0); --active_; for (auto& t : new_tasks) { tasks_.emplace(std::move(t)); } globally_finished = tasks_.empty() and active_ == 0; has_pending_work = not tasks_.empty(); } if (globally_finished) { cv_.notify_all(); } else if (has_pending_work) { cv_.notify_one(); } } }; } // namespace concurrency_tools

3. Módulo statistics.directory

El siguiente módulo implementa el análisis paralelo de un árbol de directorios. Dado un directorio raíz (root), produciremos un objeto de tipo directory_statistics que contendrá:
  • Un contenedor asociativo std::map que relacione cada extensión de fichero encontrada (.txt, .zip, etc.) con un objeto extension_statistics. Este último almacenará tanto el número de archivos de dicha extensión presentes en el árbol de directorios como el tamaño total acumulado por todos ellos.
  • El número total de subdirectorios descubiertos.
  • El número total de errores acontecidos en el análisis, si no se pudo (i) abrir un directorio, (ii) obtener información de una entrada concreta o (iii) continuar iterando el directorio.
La función run_directory_statistics() será la encargada de coordinar la ejecución paralela. Para ello, inicializará una cola concurrente de tareas dynamic_task_queue<std::filesystem::path> con el directorio raíz y hará participar a un número configurable num_workers de trabajadores, cada uno de los cuales devolverá un objeto directory_statistics con las estadísticas parciales obtenidas durante su trabajo. En este caso, num_workers - 1 trabajadores se lanzan de forma asíncrona mediante std::async ([2], véase este artículo previo para más detalles),  mientras que el hilo principal actúa como un trabajador adicional procesando la cola compartida.

Cada directorio dir pendiente de procesar se representa aquí como una tarea almacenada en la cola dynamic_task_queue<std::filesystem::path>, con std::filesystem::is_directory(dir) == true. Los trabajadores ejecutarán la función process_directories(), que adquiere directorios de la cola mediante acquire(), inspecciona su contenido y contabiliza los archivos encontrados por extensión y tamaño acumulado. Los subdirectorios descubiertos durante la exploración no se procesan inmediatamente, sino que se registran como nuevas tareas en el objeto acquired_task asociado al directorio actual. Al finalizar el procesamiento, dicho objeto transfiere automáticamente las tareas acumuladas a la cola y notifica la conclusión de la tarea original mediante RAII.
A diferencia de la solución presentada en el artículo anterior de esta serie, basada en el iterador recursivo  std::filesystem::recursive_directory_iterator, la implementación de process_directories() se apoya en std::filesystem::directory_iterator [3], que únicamente recorre las entradas contenidas en un directorio y no visita recursivamente sus subdirectorios. El orden de iteración no está especificado por el estándar, excepto que cada entrada de directorio se visita una única vez.
Remarquemos una vez más que, en este esquema, cada trabajador construye su propio objeto directory_statistics, en el que acumula las estadísticas correspondientes a los directorios que procesa.
Una vez finalizada la exploración, los resultados parciales se fusionan mediante merge_statistics(), acumulando para cada extensión el número de archivos y el tamaño total, así como el número total de subdirectorios visitados.

module; #include <cstdio> export module statistics.directory; import std; import concurrency_tools.dynamic_task_queue; namespace statistics { export struct extension_statistics {    std::uintmax_t num_files = 0;    std::uintmax_t total_size = 0; }; // resultado total, o parcial producido por un worker: export struct directory_statistics {    std::map<std::string, extension_statistics> files;    std::uintmax_t num_directories = 0;
   std::uintmax_t num_errors = 0;
}; using directory_queue = concurrency_tools::dynamic_task_queue<std::filesystem::path>; [[nodiscard]] auto process_directories(directory_queue& directories) -> directory_statistics {    auto res = directory_statistics{};    while (auto directory = directories.acquire()) {       // sólo procesamos el directorio 'dir'; los subdirectorios       // encontrados pasarán a la cola y serán procesados // posteriormente por algún worker:       auto const& dir = directory->get();       auto ec = std::error_code{};       auto it = std::filesystem::directory_iterator{dir, ec};       if (ec) {          // no podemos acceder al directorio; la tarea actual se          // completará automáticamente tras 'continue':          ++res.num_errors;          continue;       }     auto end = std::filesystem::directory_iterator{};       for (; it != end; it.increment(ec)) {          if (ec) { // error al avanzar por el directorio ++res.num_errors;          break;          }          auto const& entry = *it;          auto const status = entry.symlink_status(ec);        if (ec) { ++res.num_errors;           continue;          }        if (std::filesystem::is_symlink(status)) {             continue;          }          else if (std::filesystem::is_directory(status)) {             ++res.num_directories;           directory->add_task(entry.path());         }          else if (std::filesystem::is_regular_file(status)) {             auto const size = entry.file_size(ec);           if (ec) { ++res.num_errors;             continue;           } auto const extension = entry.path().extension().string();           auto& [num_files, total_size] = res.files[extension];           ++num_files;             total_size += size;         } }    } // destrucción del acquired_task 'directory' y llamada automática // a complete() antes de la próxima iteración del bucle, incorporando // los subdirectorios encontrados a la cola de tareas    return res; } // fusiona el resultado de un worker en el resultado global: void merge_statistics(    directory_statistics& destination,    directory_statistics const& source ){    destination.num_directories += source.num_directories; destination.num_errors += source.num_errors;    for (auto const& [extension, info] : source.files) {       auto& destination_info = destination.files[extension];       destination_info.num_files += info.num_files;       destination_info.total_size += info.total_size;    } } export [[nodiscard]] auto run_directory_statistics(    std::filesystem::path const& root,    std::size_t num_workers )    -> directory_statistics    pre(std::filesystem::is_directory(root))    pre(num_workers >= 1) {    auto directories = directory_queue{root};    auto futures = std::views::indices(num_workers - 1)       | std::views::transform([&]([[maybe_unused]] auto worker_id)            -> std::future<directory_statistics> {            return std::async(               std::launch::async,               process_directories,               std::ref(directories)            );         })       | std::ranges::to<std::vector>();    auto total = process_directories(directories);    for (auto& f : futures) {       merge_statistics(total, f.get());    }    return total; } } // namespace statistics

4. Módulos auxiliares

Antes de abordar el programa principal, definiremos tres módulos auxiliares sencillos que proporcionan compatibilidad con componentes de la biblioteca C y utilidades de formato para la salida del programa. Analizaremos en primer lugar los módulos de compatibilidad:
  • c_tools.exit_codes: expone los valores enteros de finalización estándar del lenguaje, EXIT_SUCCESS y EXIT_FAILURE, a través del espacio de nombres c_tools. Su objetivo es facilitar el uso de dichos valores desde código modular sin depender directamente de las macros definidas en <cstdlib>.
  • c_tools.standard_streams: proporciona funciones noexcept que devuelven punteros std::FILE* a los flujos estándar stdinstdout y stderr definidos en <cstdio>.
module; #include <cstdlib> export module c_tools.exit_codes; export namespace c_tools { constexpr int exit_success = EXIT_SUCCESS; constexpr int exit_failure = EXIT_FAILURE; } // namespace c_tools

module; #include <cstdio> export module c_tools.standard_streams; export namespace c_tools { [[nodiscard]] auto stdin_stream() noexcept -> std::FILE* { return stdin; } [[nodiscard]] auto stdout_stream() noexcept -> std::FILE* { return stdout; } [[nodiscard]] auto stderr_stream() noexcept -> std::FILE* { return stderr; } } // namespace c_tools

Finalmente:
  • format_tools: amplía std::format mediante especializaciones de std::formatter, permitiendo representar  tamaños binarios en unidades más legibles (KiB, MiB y GiB) y cantidades enteras con separadores de miles.
export module format_tools; import std; export namespace format_tools { struct binary_size { std::uint64_t value; }; struct thousands_separated { std::uint64_t value; }; } // namespace format_tools template<> struct std::formatter<format_tools::binary_size> : std::formatter<std::string_view> { template<typename format_context> auto format( format_tools::binary_size const& bsz, format_context& ctx ) const { constexpr auto KiB = 1024.0; constexpr auto MiB = 1024.0*KiB; constexpr auto GiB = 1024.0*MiB; auto format_size = [](std::uint64_t value) -> std::pair<double, std::string_view> { if (value >= GiB) { return {value/GiB, "GiB"}; } if (value >= MiB) { return {value/MiB, "MiB"}; } if (value >= KiB) { return {value/KiB, "KiB"}; } return {static_cast<double>(value), "bytes"}; }; auto const [size, unit] = format_size(bsz.value); auto const str = std::format("{:.1f} {}", size, unit); return std::formatter<std::string_view>::format(str, ctx); } }; template<> struct std::formatter<format_tools::thousands_separated> : std::formatter<std::string_view> { template<typename format_context> auto format( format_tools::thousands_separated const& n, format_context& ctx ) const { auto str = std::to_string(n.value); for (auto i = str.size(); i > 3uz; i -= 3) { str.insert(i - 3, 1, '.'); } return std::formatter<std::string_view>::format(str, ctx); } };

5. Benchmarks

La siguiente función principal main() actúa como banco de pruebas de los distintos módulos desarrollados anteriormente. Ejecuta el análisis estadístico de un árbol de directorios indicado como argumento en la línea de comandos, empleando un número creciente de hilos trabajadores y midiendo el tiempo de ejecución en cada caso. A partir de estas mediciones, se calculan e imprimen métricas de rendimiento como el speedup respecto a la ejecución secuencial (un único hilo trabajador) y la mejora porcentual obtenida al incrementar el número de trabajadores:

import std; import c_tools.exit_codes; import c_tools.standard_streams; import format_tools; import statistics.directory; namespace stdc = std::chrono; namespace stdf = std::filesystem; namespace stdv = std::views; struct benchmark_result { statistics::directory_statistics statistics; stdc::milliseconds duration; }; [[nodiscard]] auto run_benchmark( stdf::path const& root, std::size_t num_workers ) -> benchmark_result { using clock = stdc::steady_clock; auto const start = clock::now(); auto const total = statistics::run_directory_statistics(root, num_workers); using ms = stdc::milliseconds; return { .statistics = total, .duration = stdc::duration_cast<ms>(clock::now() - start) }; } auto main(int argc, char* argv[]) -> int { if (argc != 2) { std::println( c_tools::stderr_stream(), "usage: {} <directory>", argv[0] ); return c_tools::exit_failure; } auto const root = stdf::path{argv[1]}; auto ec = std::error_code{}; if (not stdf::is_directory(root, ec)) { std::println( c_tools::stderr_stream(), "error: '{}': {}", root.string(), ec? ec.message() : "not a directory" ); return c_tools::exit_failure; } auto const concurrency = std::thread::hardware_concurrency(); auto const max_workers = std::max(1u, concurrency); std::println( "{:>7} {:>10} {:>10} {:>15} {:>14}", "Workers", "Time (ms)", "Speedup", "vs. 1 worker", "vs. previous" ); auto baseline = stdc::milliseconds{}; auto previous = 0.0; for (auto const num_workers : stdv::iota(1u, max_workers + 1)) { auto const [stats, duration] = run_benchmark(root, num_workers); if (num_workers == 1) { baseline = duration; } auto const current = static_cast<double>(duration.count()); auto const speedup = static_cast<double>(baseline.count())/current; std::println( "{:>7} {:>10} {:>9.1f}× {:>14.1f}% {:>14}", num_workers, current, speedup, 100.0*(1.0 - 1.0/speedup), previous? std::format("{:.1f}%", 100*(1.0 - current/previous)) : "-" ); previous = current; if (num_workers == max_workers) { auto total_files = std::uintmax_t{}; auto total_size = std::uintmax_t{}; for (auto const& [extension, info] : stats.files) { total_files += info.num_files; total_size += info.total_size; } std::println( "{:_^60}\ncontains: {} files, {} folders\nsize: {} ({} bytes)\nerrors: {}", "", format_tools::thousands_separated{total_files}, format_tools::thousands_separated{stats.num_directories}, format_tools::binary_size{total_size}, format_tools::thousands_separated{total_size}, stats.num_errors ); } } return c_tools::exit_success; }
El número de hilos trabajadores num_workers se establece por defecto a partir de std::thread::hardware_concurrency() [4], que proporciona una estimación del grado de concurrencia disponible en el sistema, habitualmente coincidente con el número de hilos de hardware. Sin embargo, este valor debe interpretarse como un límite superior orientativo y no necesariamente como el número óptimo de trabajadores. En la práctica, el rendimiento puede saturarse con un número menor de hilos debido a factores como la contención por la cola compartida, el tipo de dispositivo de almacenamiento o las limitaciones del propio sistema de ficheros.
A modo de ejemplo, sobre un equipo con un procesador Intel Core i5-1135G7 de 11ª generación a 2.40 GHz (2.42 GHz efectivos), 8 hilos de ejecución hardware y una unidad SSD SK hynix HFM256GD3HX015N, se obtuvieron los siguientes resultados al analizar un árbol de directorios de prueba:
Workers  Time (ms)    Speedup    vs. 1 worker   vs. previous
      1      16331       1.0×            0.0%              -
      2      10194       1.6×           37.6%          37.6%
      3       7777       2.1×           52.4%          23.7%
      4       6242       2.6×           61.8%          19.7%
      5       4994       3.3×           69.4%          20.0%
      6       4545       3.6×           72.2%           9.0%
      7       4281       3.8×           73.8%           5.8%
      8       4125       4.0×           74.7%           3.6%
____________________________________________________________ contains: 136.613 files, 7.695 folders size: 7.2 GiB (7.679.411.327 bytes) errors: 0
Aquí, la columna "vs. 1 worker" expresa la reducción porcentual del tiempo de ejecución respecto a la versión secuencial.

El contenedor asociativo contenido en el resultado global directory_statistics permite, entre otras acciones, imprimir desgloses como el mostrado en la introducción de este artículo, o la búsqueda de información acerca de una extensión concreta. A modo de ejemplo:

auto const [files, num_directories, num_errors] = statistics::run_directory_statistics(root, 8/*workers*/); auto const extension = ".txt"; if ( auto const it = files.find(extension); it != files.end() ){ auto const& [num_files, total_size] = it->second; std::println( "{} files, {} ({} bytes)", format_tools::thousands_separated{num_files}, format_tools::binary_size{total_size}, format_tools::thousands_separated{total_size} ); } else { std::println("no {} files found", extension); }

Referencias bibliográficas

  1. cppreference – std::filesystem::recursive_directory_iterator – https://en.cppreference.com/cpp/filesystem/recursive_directory_iterator
  2. cppreference – std::async– https://en.cppreference.com/cpp/thread/async
  3. cppreference – std::filesystem::directory_iterator – https://en.cppreference.com/cpp/filesystem/directory_iterator
  4. cppreference – std::thread::hardware_concurrency – https://cppreference.com/cpp/thread/thread/hardware_concurrency

Comentarios