Packages

Librería base para ejecución de comandos del sistema y gestión de tareas asíncronas con resultados estructurados.

Retired package: Release invalid - Versión publicada por error

Current section

Files

Jump to
argos lib parallel example.ex
Raw

lib/parallel/example.ex

defmodule Argos.Parallel.Example do
@moduledoc """
Ejemplo completo de uso de `Argos.Parallel` para ejecución paralela de tareas.
Este módulo sirve como referencia y guía para implementar procesamiento paralelo
usando `Argos.Parallel`. Demuestra cómo:
- Iniciar el sistema paralelo
- Crear workers con múltiples tareas
- Suscribirse a eventos del Leader (eventos individuales)
- Suscribirse a eventos del Monitor (estado consolidado)
- Procesar eventos en tiempo real
- Obtener estadísticas y resultados finales
- Detener el sistema correctamente
## Tipos de Suscripción
`Argos.Parallel` ofrece dos tipos de suscripción:
### Leader (Eventos Individuales)
- Recibe eventos individuales de cada worker en tiempo real
- Útil para visualización detallada o procesamiento específico por evento
- Mensaje: `{:parallel_event, leader, event}`
- Eventos: `:result`, `:progress`, `:finished`, `:error`
### Monitor (Estado Consolidado)
- Recibe actualizaciones del estado consolidado de todos los workers
- Útil para visualización general o dashboards
- Mensaje: `{:parallel_monitor_update, monitor_state}`
- Contiene el estado completo de todos los workers
## Ejecución
Para ejecutar este ejemplo:
mix run -e "Argos.Parallel.Example.start()"
## Estructura del Ejemplo
El ejemplo crea 4 workers con diferentes características:
1. **data_processor**: Procesamiento de datos con diferentes duraciones (con logging habilitado)
2. **calculator**: Cálculos matemáticos (suma, factorial, potencia) (con logging habilitado)
3. **io_simulator**: Simulación de operaciones de E/S con posibilidad de error (sin logging)
4. **fast_worker**: Operaciones rápidas múltiples (sin logging)
Los workers con logging habilitado generan archivos de log individuales en
`~/.argos/logs/parallel_workers/worker_<worker_id>.jsonl` que se eliminan
automáticamente cuando se detiene el sistema.
## Flujo de Trabajo
1. Iniciar el sistema con `Argos.Parallel.start_system()`
2. Suscribirse a los canales deseados (`subscribe_leader()` y/o `subscribe_monitor()`)
3. Crear especificaciones de workers con `create_worker_spec/3` (puede incluir `log: true`)
4. Iniciar workers con `start_workers/1`
5. Escuchar eventos en un loop con `receive`
6. Procesar eventos según el tipo
7. Detener el sistema con `stop_system()` cuando termine (los logs se eliminan automáticamente)
## Ejemplo de Worker
### Worker sin logging
worker_spec = Argos.Parallel.create_worker_spec(:my_worker, [
fn ->
Process.sleep(1000)
{:ok, :result_1}
end,
fn ->
Process.sleep(2000)
{:ok, :result_2}
end
])
### Worker con logging habilitado
worker_spec_with_log = Argos.Parallel.create_worker_spec(:my_worker, [
fn ->
Process.sleep(1000)
{:ok, :result_1}
end,
fn ->
Process.sleep(2000)
{:ok, :result_2}
end
], log: true)
Los workers con `log: true` generan un archivo de log individual en formato JSON Lines
en `~/.argos/logs/parallel_workers/worker_<worker_id>.jsonl`. Todos los logs se eliminan
automáticamente cuando se llama a `Argos.Parallel.stop_system()`.
## Procesamiento de Eventos
### Eventos del Leader
receive do
{:parallel_event, _leader, event} ->
case event do
%{type: :result, worker_id: worker_id, data: data} ->
IO.puts("Worker " <> to_string(worker_id) <> " completó tarea " <> Integer.to_string(data.task_index))
%{type: :progress, worker_id: worker_id, data: data} ->
IO.puts("Worker " <> to_string(worker_id) <> " progreso: " <> Float.to_string(data.percent) <> "%")
%{type: :finished, worker_id: worker_id} ->
IO.puts("Worker " <> to_string(worker_id) <> " terminó")
%{type: :error, worker_id: worker_id, data: data} ->
IO.puts("Worker " <> to_string(worker_id) <> " error: " <> to_string(data.reason))
end
end
### Eventos del Monitor
receive do
{:parallel_monitor_update, monitor_state} ->
stats = Argos.Parallel.get_stats()
IO.puts("Progreso general: " <> Float.to_string(stats.progress_percentage) <> "%")
monitor_state.workers
|> Enum.each(fn {id, worker} ->
IO.puts(to_string(id) <> ": " <> to_string(worker.status) <> " - " <> Integer.to_string(worker.progress) <> "%")
end)
end
## Manejo de Errores
Cuando un worker encuentra un error:
- El worker se detiene inmediatamente
- Se envía un evento `:error` con los detalles
- Las tareas restantes del worker no se ejecutan
- Otros workers continúan normalmente
## Estadísticas
Obtener estadísticas generales:
stats = Argos.Parallel.get_stats()
## Sistema de Logging
`Argos.Parallel` incluye un sistema de logging automático que permite generar
archivos de log individuales por worker. Para habilitar el logging en un worker,
usa la opción `log: true` al crear la especificación:
Argos.Parallel.create_worker_spec(:worker_id, tasks, log: true)
Los logs se guardan en formato JSON Lines (una línea JSON por evento) en:
`~/.argos/logs/parallel_workers/worker_<worker_id>.jsonl`
Todos los logs se eliminan automáticamente cuando se detiene el sistema con
`Argos.Parallel.stop_system()`.
## Ver También
- `Argos.Parallel`: API principal del sistema paralelo
- `Argos.Parallel.Leader`: Gestión de eventos individuales
- `Argos.Parallel.Monitor`: Gestión de estado consolidado
- `Argos.Parallel.WorkerState`: Estado de un worker individual
- `Argos.Parallel.Logger`: Sistema de logging automático
"""
@doc """
Inicia el ejemplo completo de procesamiento paralelo.
Este es el punto de entrada principal que demuestra todo el flujo de trabajo.
"""
alias Argos.Parallel.WorkerState
@spec start() :: :ok
def start do
case Argos.Parallel.start_system() do
{:ok, _pid} ->
IO.puts("✅ Sistema paralelo iniciado")
{:error, {:already_started, _pid}} ->
IO.puts("✅ Sistema paralelo ya estaba ejecutándose")
{:error, reason} ->
IO.puts("❌ Error iniciando sistema: #{inspect(reason)}")
end
:ok = Argos.Parallel.subscribe_leader()
:ok = Argos.Parallel.subscribe_monitor()
IO.puts("✅ Suscrito a Leader y Monitor")
worker_specs = create_sample_workers()
IO.puts("🚀 Iniciando workers...")
results = Argos.Parallel.start_workers(worker_specs)
IO.puts("✅ Workers iniciados: #{length(results)}")
listen()
end
@spec create_sample_workers() :: [%{id: any, tasks: list, log: boolean}]
defp create_sample_workers do
[
Argos.Parallel.create_worker_spec(
:data_processor,
[
fn ->
Process.sleep(1000)
{:ok, :processed_data_1, %{size: 100, timestamp: DateTime.utc_now()}}
end,
fn ->
Process.sleep(2000)
{:ok, :processed_data_2, %{size: 200, timestamp: DateTime.utc_now()}}
end,
fn ->
Process.sleep(1500)
{:ok, :processed_data_3, %{size: 150, timestamp: DateTime.utc_now()}}
end
],
log: true
),
Argos.Parallel.create_worker_spec(
:calculator,
[
fn ->
Process.sleep(800)
result = Enum.sum(1..1000)
{:sum, result}
end,
fn ->
Process.sleep(1200)
result = Enum.reduce(1..100, 1, &(&1 * &2))
{:factorial, result}
end,
fn ->
Process.sleep(500)
result = :math.pow(2, 16) |> round
{:power, result}
end
],
log: true
),
Argos.Parallel.create_worker_spec(:io_simulator, [
fn ->
Process.sleep(3000)
{:read, "file1.txt", "Content of file 1"}
end,
fn ->
Process.sleep(2500)
{:write, "file2.txt", 2048}
end,
fn ->
Process.sleep(1800)
if :rand.uniform(10) == 1 do
raise "Simulated IO error: Device not ready"
else
{:read, "file3.txt", "Content of file 3"}
end
end
]),
Argos.Parallel.create_worker_spec(:fast_worker, [
fn ->
Process.sleep(100)
:quick_task_1
end,
fn ->
Process.sleep(150)
:quick_task_2
end,
fn ->
Process.sleep(200)
:quick_task_3
end,
fn ->
Process.sleep(100)
:quick_task_4
end,
fn ->
Process.sleep(50)
:quick_task_5
end
])
]
end
@doc """
Loop principal de escucha de eventos.
Procesa eventos del Leader (individuales) y del Monitor (consolidados).
Termina después de 30 segundos o cuando todos los workers terminan.
"""
@spec listen() :: :ok
def listen do
receive do
{:parallel_event, _leader, event} ->
handle_parallel_event(event)
listen()
{:parallel_monitor_update, monitor_state} ->
handle_monitor_update(monitor_state)
listen()
after
30_000 ->
IO.puts("\n🎉 Proceso terminado (timeout)")
show_final_results()
Argos.Parallel.stop_system()
IO.puts("\n🛑 Sistema paralelo detenido")
IO.puts("📝 Logs de workers eliminados automáticamente")
end
end
defp handle_parallel_event(event) do
System.cmd("clear", []) |> elem(0) |> IO.puts()
IO.puts("")
IO.puts("\n--- EVENTO INDIVIDUAL ---")
IO.puts("Evento: #{format_event_type(event)}")
case event do
%{type: :result, worker_id: worker_id, data: data} ->
IO.puts("✅ Worker #{inspect(worker_id)} - Tarea #{data.task_index} completada: #{inspect(data.result)}")
%{type: :progress, worker_id: worker_id, data: data} ->
IO.puts("📊 Worker #{inspect(worker_id)} - Progreso: #{data.percent}% (tarea #{data.task_index}/#{data.total})")
%{type: :finished, worker_id: worker_id} ->
IO.puts("🎉 Worker #{inspect(worker_id)} - COMPLETADO")
%{type: :error, worker_id: worker_id, data: data} ->
IO.puts("❌ Worker #{inspect(worker_id)} - ERROR en tarea #{data.task_index}: #{data.reason}")
_ ->
:ok
end
end
defp handle_monitor_update(monitor_state) do
System.cmd("clear", []) |> elem(0) |> IO.puts()
IO.puts("\n--- ESTADO CONSOLIDADO ---")
stats = Argos.Parallel.get_stats()
IO.puts("📈 Progreso general: #{stats.progress_percentage}%")
IO.puts("👥 Workers totales: #{stats.total_workers}")
IO.puts("Estados workers: #{format_status_counts(stats.status_counts)}")
monitor_state.workers
|> Enum.each(&print_worker_status/1)
end
defp print_worker_status({worker_id, worker}) do
status_icon = get_status_icon(worker.status)
elapsed = WorkerState.elapsed_time(worker)
IO.puts("#{status_icon} #{inspect(worker_id)}: #{worker.status} - #{worker.progress || 0}% (#{elapsed}ms)")
end
defp get_status_icon(:running), do: "🔄"
defp get_status_icon(:finished), do: "✅"
defp get_status_icon(:error), do: "❌"
defp get_status_icon(:started), do: "🚀"
defp get_status_icon(_), do: "⏸️"
defp format_event_type(%{type: type}), do: "#{type}"
defp format_event_type(_), do: "unknown"
defp format_status_counts(counts) do
Enum.map_join(counts, ", ", fn {status, count} -> "#{status}: #{count}" end)
end
defp show_final_results do
IO.puts("\n📋 RESULTADOS FINALES:")
final_state = Argos.Parallel.get_monitor_state()
final_state.workers
|> Enum.each(&print_worker_result/1)
end
defp print_worker_result({worker_id, worker}) do
IO.puts("\n=== Worker: #{inspect(worker_id)} ===")
IO.puts("Estado: #{worker.status}")
IO.puts("Tareas completadas: #{length(worker.results)}/#{worker.total}")
IO.puts("Tiempo total: #{WorkerState.elapsed_time(worker)}ms")
print_worker_error(worker)
print_worker_results(worker)
end
defp print_worker_error(%{error: nil}), do: :ok
defp print_worker_error(%{error: error}) do
IO.puts("❌ Error: #{inspect(error)}")
end
defp print_worker_results(%{results: []}), do: :ok
defp print_worker_results(%{results: results}) do
IO.puts("Resultados:")
results
|> Enum.reverse()
|> Enum.each(fn result ->
IO.puts(" 📝 Tarea #{result.task_index}: #{inspect(result.result)}")
end)
end
end