Proyecto 2: un millón de registros, tres mecanismos de sincronización y pruebas que hacen visible la condición de carrera
Proyecto 2: un millón de registros, tres mecanismos de sincronización y pruebas que hacen visible la condición de carrera
En el capítulo 15 construimos un intérprete de comandos: creamos procesos, los esperamos, redirigimos sus descriptores y los conectamos con tuberías. Ese proyecto trabajaba con procesos aislados: cada hijo tenía su propio espacio de direcciones y toda la comunicación pasaba por el kernel, así que las condiciones de carrera aparecían solo en los bordes —el orden de cierre de descriptores, la recolección de zombis, la señal que llega entre dos llamadas al sistema—.
Este segundo proyecto invierte la situación. Aquí los flujos de ejecución comparten estado a propósito, porque necesitan cooperar sobre un mismo archivo de salida y un mismo contador de progreso. Esa es la clase de concurrencia que estudiamos desde el capítulo 7 hasta el capítulo 11: sección crítica, exclusión mutua, semáforos, monitores y paso de mensajes. El objetivo del proyecto no es solo escribir un programa rápido, sino escribir un programa cuya corrección se pueda demostrar con pruebas que fallen de manera repetible cuando la sincronización se rompe.
El enunciado tiene tres piezas y una comparación:
- Generar
input.txtcon 1.000.000 de registros sintéticos, en Elixir, con procesamiento concurrente. - Transformar
input.txtenoutput.csv, también en Elixir y también concurrente. - Producir
output.py.csvcon un programa Python estrictamente secuencial, para tener una línea base contra la cual medir. - Escribir un informe que compare tiempos, dificultades y diferencias entre archivos.
Sobre ese esqueleto vamos a montar lo que realmente enseña el proyecto: identificar el punto de contención, resolverlo con las tres familias de sincronización, romperlo a propósito y escribir las pruebas que lo cazan.
El enunciado en detalle
El registro tiene seis campos separados por el carácter |, sin espacios alrededor del separador:
| Campo | Formato | Ejemplo |
|---|---|---|
| Teléfono | prefijo 569 más 8 dígitos | 56912345678 |
| Primer nombre | texto alfabético | Camila |
| Segundo nombre | texto alfabético | Andrea |
| Primer apellido | texto alfabético | Rojas |
| Segundo apellido | texto alfabético | Vergara |
| usuario más dominio | [email protected] | |
| RUT | 7 u 8 dígitos más dígito verificador | 18345672-K |
La salida es un CSV con cuatro columnas: nombre completo, RUT, email y teléfono. Es decir, el procesamiento fusiona los cuatro campos de nombre en uno solo y reordena.
Los presupuestos de tiempo del enunciado original marcan la diferencia entre las dos fases:
| Fase | Presupuesto | Naturaleza del trabajo |
|---|---|---|
Generación de input.txt | ~6 minutos, con desviación de ~60 s | Cómputo puro más escritura secuencial |
Procesamiento a output.csv | hasta ~60 minutos | Lectura, parseo, transformación y escritura |
| Línea base secuencial en Python | sin límite, es la referencia | Igual que la anterior, un flujo |
Esos números no son la meta. Son el orden de magnitud que hace que el problema sea interesante: con mil registros cualquier implementación funciona y ninguna carrera aparece; con un millón, el planificador tiene un millón de oportunidades de entrelazar operaciones de la forma exacta que rompe un contador mal escrito.
Arquitectura de la solución
Antes de escribir código conviene separar tres responsabilidades que suelen mezclarse: producir los datos, coordinar el acceso al recurso compartido y persistir el resultado. La coordinación es el único lugar donde vive la sincronización; si se filtra a las otras dos, el programa se vuelve imposible de probar.
flowchart LR
subgraph Productores["Productores concurrentes"]
W1["Worker 1<br/>lote 1..125000"]
W2["Worker 2<br/>lote 125001..250000"]
WN["Worker N<br/>lote ..1000000"]
end
subgraph Coordinacion["Punto de contención"]
SEM["Mecanismo de<br/>sincronización"]
end
subgraph Persistencia["Recursos compartidos"]
F["Descriptor de<br/>input.txt"]
C["Contador de<br/>registros escritos"]
end
W1 --> SEM
W2 --> SEM
WN --> SEM
SEM --> F
SEM --> C
F --> OUT["input.txt<br/>1.000.000 líneas"]
C --> LOG["Reporte final<br/>de progreso"]
Hay exactamente dos recursos compartidos: el descriptor del archivo y el contador. Todo lo demás —generar un RUT, armar un correo, formatear un teléfono— es cómputo local sin estado compartido, y por lo tanto no necesita sincronización de ningún tipo. Esta distinción es la primera decisión de diseño del proyecto: cuanto más pequeña sea la sección crítica, menos tiempo pasan los trabajadores esperando y menos superficie tiene el error.
Fase 1: el generador de registros
Empecemos por el cómputo local, que no tiene concurrencia y se puede probar con tests ordinarios.
Dígito verificador del RUT
El RUT chileno usa el módulo 11 con una serie de multiplicadores 2, 3, 4, 5, 6, 7 que se repite de derecha a izquierda. El resultado 11 se representa como 0 y el 10 como K.
defmodule Proyecto2.Rut do
@moduledoc """
Generación y validación de RUT con dígito verificador módulo 11.
Este módulo es puro: no toca estado compartido ni realiza E/S.
"""
@doc "Calcula el dígito verificador de un número de RUT."
@spec digito_verificador(pos_integer()) :: String.t()
def digito_verificador(numero) when is_integer(numero) and numero > 0 do
{suma, _factor} =
numero
|> Integer.digits()
|> Enum.reverse()
|> Enum.reduce({0, 2}, fn digito, {acumulado, factor} ->
siguiente = if factor == 7, do: 2, else: factor + 1
{acumulado + digito * factor, siguiente}
end)
case 11 - rem(suma, 11) do
11 -> "0"
10 -> "K"
resto -> Integer.to_string(resto)
end
end
@doc "Genera un RUT aleatorio en el rango de personas naturales."
@spec aleatorio() :: String.t()
def aleatorio do
numero = :rand.uniform(24_000_000) + 1_000_000
"#{numero}-#{digito_verificador(numero)}"
end
@doc "Verifica que un RUT con formato 'numero-dv' sea consistente."
@spec valido?(String.t()) :: boolean()
def valido?(rut) do
case String.split(rut, "-") do
[numero, dv] ->
case Integer.parse(numero) do
{n, ""} -> digito_verificador(n) == String.upcase(dv)
_ -> false
end
_ ->
false
end
end
end
El registro completo
defmodule Proyecto2.Registro do
@moduledoc "Construcción de un registro sintético del formato del proyecto."
alias Proyecto2.Rut
@nombres ~w(Camila Matias Valentina Benjamin Sofia Vicente Isidora Agustin
Antonia Lucas Josefa Martin Emilia Tomas Florencia Joaquin)
@apellidos ~w(Gonzalez Munoz Rojas Diaz Perez Soto Contreras Silva Martinez
Sepulveda Morales Fuentes Torres Araya Flores Valenzuela)
@dominios ~w(correo.cl mail.com ejemplo.org red.cl portal.net)
@doc """
Devuelve un registro como iodata. Se usa iodata y no un binario porque
la escritura a disco acepta listas anidadas sin concatenar en memoria.
"""
@spec generar() :: iodata()
def generar do
nombre1 = Enum.random(@nombres)
apellido1 = Enum.random(@apellidos)
telefono = "569" <> String.pad_leading(Integer.to_string(:rand.uniform(99_999_999)), 8, "0")
email =
String.downcase(nombre1) <> "." <> String.downcase(apellido1) <>
Integer.to_string(:rand.uniform(999)) <> "@" <> Enum.random(@dominios)
campos = [
telefono, nombre1, Enum.random(@nombres),
apellido1, Enum.random(@apellidos), email, Rut.aleatorio()
]
[Enum.intersperse(campos, "|"), "\n"]
end
end
Nótese que generar/0 devuelve iodata: una lista anidada de binarios. Esto importa para la concurrencia porque permite entregar el registro completo a la capa de escritura en una sola operación, en lugar de hacer siete escrituras separadas. Más adelante veremos que esa diferencia es exactamente la que separa un archivo correcto de un archivo con líneas entrelazadas.
El punto de contención
Con el cómputo local resuelto, queda la parte difícil. Un millón de registros repartidos entre N trabajadores, todos queriendo escribir en el mismo archivo y actualizar el mismo contador.
La condición de carrera clásica del contador es la de lectura-modificación-escritura perdida. Vale la pena verla con precisión temporal antes de resolverla:
sequenceDiagram
autonumber
participant A as Worker A
participant M as Memoria compartida<br/>(contador = 41)
participant B as Worker B
A->>M: leer contador
M-->>A: 41
Note over A: A calcula 41 + 1 = 42<br/>pero aún no escribe
B->>M: leer contador
M-->>B: 41
Note over B: B calcula 41 + 1 = 42
A->>M: escribir 42
B->>M: escribir 42
Note over M: contador = 42<br/>se perdió un incremento:<br/>debería ser 43
El incremento perdido es silencioso: no hay excepción, no hay corrupción de memoria, el programa termina normalmente y reporta un número que es un poco menor que el real. En una ejecución con mil operaciones puede no aparecer nunca; en una con un millón aparece decenas de veces. Esa es precisamente la razón por la que el proyecto pide un millón de registros y no mil.
La versión rota a propósito
Escribamos primero el contador incorrecto. Tenerlo en el repositorio, con un test que lo caza, es parte del entregable: demuestra que se entiende por qué la solución correcta lo es.
defmodule Proyecto2.Contador.Roto do
@moduledoc """
Contador con condición de carrera deliberada.
Existe únicamente para que las pruebas demuestren el fallo.
"""
use Agent
def start_link(_opts \\ []) do
Agent.start_link(fn -> 0 end)
end
@doc """
Lee, calcula fuera del Agent y escribe. Entre la lectura y la escritura
hay una ventana en la que otro proceso puede leer el mismo valor.
"""
def incrementar(pid) do
valor = Agent.get(pid, & &1)
# Punto de replanificación explícito: hace la ventana observable
# sin cambiar la naturaleza del error.
Process.sleep(0)
Agent.update(pid, fn _ignorado -> valor + 1 end)
end
def valor(pid), do: Agent.get(pid, & &1)
end
El detalle importante: Agent.get/2 y Agent.update/3 son individualmente atómicas —el Agent es un proceso y atiende un mensaje a la vez—, pero la secuencia de ambas no lo es. La atomicidad de las partes no da atomicidad al todo. Este es el error más frecuente de quien llega a Elixir desde otro lenguaje: se supone que “como todo es un proceso, no hay carreras”. Las hay, y viven exactamente en los huecos entre dos llamadas.
Las tres soluciones
El proyecto pide resolverlo con semáforos, monitores o mensajes. Vamos a implementar los tres para el mismo problema, porque la comparación es la parte formativa.
Solución 1: semáforo contador
El semáforo no protege un dato; limita la cantidad de trabajadores que pueden estar dentro de una región al mismo tiempo. Con un permiso funciona como exclusión mutua; con N permisos, como limitador de concurrencia. En este proyecto sirve para dos cosas distintas: proteger la sección crítica del contador (un permiso) y acotar cuántos trabajadores tienen un descriptor de archivo abierto a la vez (N permisos).
defmodule Proyecto2.Semaforo do
@moduledoc """
Semáforo contador implementado como GenServer.
Las llamadas que no consiguen permiso quedan encoladas en `from`
y se responden cuando alguien libera.
"""
use GenServer
def start_link(permisos, opts \\ []) when is_integer(permisos) and permisos > 0 do
GenServer.start_link(__MODULE__, permisos, opts)
end
@doc "Adquiere un permiso. Bloquea al llamador hasta obtenerlo."
def adquirir(sem, timeout \\ :infinity), do: GenServer.call(sem, :adquirir, timeout)
@doc "Devuelve un permiso y despierta a un esperante si lo hay."
def liberar(sem), do: GenServer.cast(sem, :liberar)
def disponibles(sem), do: GenServer.call(sem, :disponibles)
@doc "Ejecuta `fun` con un permiso tomado y lo devuelve siempre."
def con_permiso(sem, fun) do
:ok = adquirir(sem)
try do
fun.()
after
liberar(sem)
end
end
@impl true
def init(permisos), do: {:ok, %{permisos: permisos, cola: :queue.new()}}
@impl true
def handle_call(:adquirir, _from, %{permisos: p} = estado) when p > 0 do
{:reply, :ok, %{estado | permisos: p - 1}}
end
@impl true
def handle_call(:adquirir, from, %{cola: cola} = estado) do
# Sin permisos: no se responde ahora, se guarda `from` para responder después.
{:noreply, %{estado | cola: :queue.in(from, cola)}}
end
@impl true
def handle_call(:disponibles, _from, estado) do
{:reply, estado.permisos, estado}
end
@impl true
def handle_cast(:liberar, %{cola: cola, permisos: p} = estado) do
case :queue.out(cola) do
{{:value, esperante}, resto} ->
# El permiso pasa directamente al esperante: no vuelve al contador.
GenServer.reply(esperante, :ok)
{:noreply, %{estado | cola: resto}}
{:empty, _} ->
{:noreply, %{estado | permisos: p + 1}}
end
end
end
El uso queda así:
defmodule Proyecto2.Contador.ConSemaforo do
@moduledoc "Contador protegido por un semáforo binario."
alias Proyecto2.Semaforo
def start_link do
{:ok, sem} = Semaforo.start_link(1)
{:ok, celda} = Agent.start_link(fn -> 0 end)
{:ok, {sem, celda}}
end
def incrementar({sem, celda}) do
Semaforo.con_permiso(sem, fn ->
valor = Agent.get(celda, & &1)
Process.sleep(0)
Agent.update(celda, fn _ -> valor + 1 end)
end)
end
def valor({_sem, celda}), do: Agent.get(celda, & &1)
end
El cuerpo del incremento es literalmente el mismo código que la versión rota, incluido el Process.sleep(0). Lo único que cambió es que ahora está encerrado entre la adquisición y la liberación de un permiso. Que la versión con carrera y la versión correcta compartan el cuerpo hace evidente qué aporta el semáforo: no cambia el cálculo, cambia quién puede estar ejecutándolo.
Solución 2: monitor
Un monitor encapsula el dato y las operaciones que lo tocan, y garantiza que solo un flujo esté dentro en cada momento. En Elixir, un GenServer es un monitor: el dato vive en su estado, las operaciones son sus callbacks, y el buzón serializa las entradas. No hace falta ningún candado explícito porque la exclusión mutua es una propiedad estructural del proceso.
defmodule Proyecto2.Contador.Monitor do
@moduledoc """
Contador como monitor: el estado y las operaciones viven juntos
y el buzón del proceso impone la exclusión mutua.
"""
use GenServer
def start_link(opts \\ []), do: GenServer.start_link(__MODULE__, 0, opts)
@doc "Incremento de una sola operación: no hay ventana entre leer y escribir."
def incrementar(pid), do: GenServer.cast(pid, :incrementar)
@doc "Suma un lote completo con un solo mensaje."
def sumar(pid, n) when is_integer(n), do: GenServer.cast(pid, {:sumar, n})
def valor(pid), do: GenServer.call(pid, :valor)
@impl true
def init(inicial), do: {:ok, inicial}
@impl true
def handle_cast(:incrementar, cuenta), do: {:noreply, cuenta + 1}
@impl true
def handle_cast({:sumar, n}, cuenta), do: {:noreply, cuenta + n}
@impl true
def handle_call(:valor, _from, cuenta), do: {:reply, cuenta, cuenta}
end
Hay una trampa didáctica escondida en incrementar/1: usa cast, que es asíncrono. Si un test hace un millón de cast y consulta el valor inmediatamente, puede leer un número menor, no por una carrera sino porque los mensajes siguen en el buzón. La solución no es un candado: es una lectura síncrona posterior, que por ser un call se encola detrás de todos los cast pendientes y por lo tanto los “descarga”. El orden de mensajes entre dos procesos dados está garantizado en la máquina virtual de Erlang, y esa garantía es la que hace correcto el patrón.
Solución 3: paso de mensajes puro
La tercera familia elimina el estado compartido: no hay dato que proteger porque cada trabajador informa su resultado y un recolector suma. Es el modelo que el proyecto favorece de forma natural, porque Task.async_stream/3 ya lo implementa por debajo.
defmodule Proyecto2.Recolector do
@moduledoc """
Agregación por paso de mensajes: no existe estado compartido.
Cada trabajador envía su total parcial y el recolector los suma.
"""
@doc """
Lanza `n_trabajadores` procesos, cada uno con `por_trabajador` incrementos,
y devuelve la suma total recibida por mensajes.
"""
def contar(n_trabajadores, por_trabajador) do
padre = self()
referencia = make_ref()
for id <- 1..n_trabajadores do
spawn_link(fn ->
parcial = Enum.reduce(1..por_trabajador, 0, fn _, acc -> acc + 1 end)
send(padre, {referencia, id, parcial})
end)
end
recibir(referencia, n_trabajadores, 0)
end
defp recibir(_ref, 0, total), do: total
defp recibir(ref, pendientes, total) do
receive do
{^ref, _id, parcial} -> recibir(ref, pendientes - 1, total + parcial)
after
30_000 -> raise "timeout esperando #{pendientes} trabajadores"
end
end
end
El after 30_000 no es cosmético: convierte un interbloqueo en una excepción con mensaje. Una prueba que se cuelga para siempre no informa nada; una que falla por timeout en treinta segundos señala exactamente qué trabajador nunca respondió.
Comparación de los tres mecanismos
| Criterio | Semáforo | Monitor (GenServer) | Paso de mensajes |
|---|---|---|---|
| Dónde vive el estado | Fuera, en una celda aparte | Dentro del proceso | No hay estado compartido |
| Quién garantiza la exclusión | El programador, al parear adquirir/liberar | La estructura del proceso | No aplica |
| Error típico | Olvidar liberar tras una excepción | Componer dos llamadas creyendo que son una | Perder un mensaje y esperar por siempre |
| Detección del error | Bloqueo permanente o permisos que se agotan | Valores menores a lo esperado | Timeout en receive |
| Costo por operación | Dos mensajes al semáforo más la operación | Un mensaje | Un mensaje por resultado, no por operación |
| Escala con N trabajadores | Contención en el semáforo | Contención en el buzón | Escala mientras el trabajo parcial sea grande |
| Ideal para | Limitar cupos y recursos escasos | Invariantes sobre un dato compartido | Agregación de trabajo particionable |
La lectura práctica de la tabla: el paso de mensajes gana cuando el trabajo se puede particionar y cada parte produce un resultado independiente, que es exactamente el caso de generar un millón de registros. El semáforo gana cuando el recurso escaso es físico —descriptores de archivo, conexiones, memoria— y hay que poner un techo. El monitor gana cuando existe un invariante que debe cumplirse después de cada operación individual.
Fase 1 completa: generación concurrente de input.txt
Con las piezas listas, el generador queda corto. La estrategia es partir el rango en lotes, procesar cada lote en un proceso distinto, y hacer que cada uno escriba su propio archivo parcial. La concatenación final es secuencial pero costa casi nada porque es E/S continua.
defmodule Proyecto2.Generador do
@moduledoc "Generación concurrente de input.txt con archivos parciales."
alias Proyecto2.Registro
@total 1_000_000
@tam_lote 25_000
@doc """
Genera `total` registros usando `concurrencia` procesos simultáneos.
Devuelve {microsegundos, cantidad_de_registros}.
"""
def generar(ruta \\ "input.txt", total \\ @total, opts \\ []) do
concurrencia = Keyword.get(opts, :concurrencia, System.schedulers_online())
tam_lote = Keyword.get(opts, :tam_lote, @tam_lote)
dir_parciales = Path.join(Path.dirname(ruta), "parciales")
File.rm_rf!(dir_parciales)
File.mkdir_p!(dir_parciales)
{micros, escritos} =
:timer.tc(fn ->
1..total
|> Stream.chunk_every(tam_lote)
|> Stream.with_index()
|> Task.async_stream(
fn {lote, indice} -> escribir_parcial(dir_parciales, indice, length(lote)) end,
max_concurrency: concurrencia,
timeout: :infinity,
ordered: false
)
|> Enum.reduce(0, fn {:ok, n}, acc -> acc + n end)
end)
unir_parciales(dir_parciales, ruta)
File.rm_rf!(dir_parciales)
{micros, escritos}
end
defp escribir_parcial(dir, indice, cantidad) do
ruta = Path.join(dir, "parte-#{String.pad_leading(Integer.to_string(indice), 6, "0")}.txt")
File.write!(ruta, Enum.map(1..cantidad, fn _ -> Registro.generar() end))
cantidad
end
defp unir_parciales(dir, destino) do
salida = File.open!(destino, [:write, :raw, {:delayed_write, 512 * 1024, 2_000}])
dir
|> File.ls!()
|> Enum.sort()
|> Enum.each(fn n -> IO.binwrite(salida, File.read!(Path.join(dir, n))) end)
File.close(salida)
end
end
Tres decisiones merecen explicación:
timeout: :infinity: el valor por omisión deTask.async_stream/3es 5000 milisegundos por elemento. Con lotes de 25.000 registros, cualquier máquina con carga lo supera y la tarea muere con:timeout. Es el error número uno del proyecto.ordered: false: no interesa recibir los resultados en orden porque solo se suman cantidades. Conordered: true(el valor por omisión), un lote lento retiene los resultados de los que ya terminaron y baja el aprovechamiento.max_concurrency: por omisión esSystem.schedulers_online(), que suele igualar la cantidad de núcleos. Para trabajo de CPU es el punto de partida razonable; para trabajo dominado por E/S se puede subir.
Y una decisión de arquitectura: ningún trabajador comparte el descriptor del archivo final. Cada uno tiene el suyo. Esto elimina por construcción la carrera de escritura, y es la forma preferida cuando el problema lo permite. La sección siguiente muestra qué pasa cuando no lo permite.
Cuando el descriptor sí se comparte
Supongamos que el requerimiento cambia y hay que escribir en un único archivo abierto, sin parciales —por ejemplo, porque el destino es una tubería o un socket—. Ahora sí existe una carrera de E/S, y es más sutil que la del contador.
En Elixir, un archivo abierto en modo normal está respaldado por un proceso: cada IO.write/2 es un mensaje y el proceso los atiende de a uno. Eso significa que una escritura completa nunca se parte. Pero si un registro se escribe con siete llamadas —una por campo—, entre la tercera y la cuarta puede colarse la escritura de otro proceso, y el archivo queda con líneas mezcladas.
flowchart TD
A["Worker A tiene un registro<br/>de 7 campos"] --> B{"¿Cómo lo escribe?"}
B -->|"7 llamadas IO.write<br/>una por campo"| C["Cada llamada es un mensaje<br/>independiente al proceso archivo"]
C --> D["Worker B intercala<br/>sus propios mensajes"]
D --> E["Línea resultante:<br/>campos de A y B mezclados"]
E --> F["Parseo posterior falla<br/>o produce datos falsos"]
B -->|"1 llamada IO.write<br/>con iodata completa"| G["Un solo mensaje<br/>al proceso archivo"]
G --> H["El proceso lo escribe entero<br/>antes de atender el siguiente"]
H --> I["Líneas íntegras<br/>sin sincronización adicional"]
style E fill:#7a2d2d,color:#fff
style I fill:#1f5c3a,color:#fff
La conclusión operativa es doble. Primero: construir el registro completo en memoria y escribirlo con una sola operación elimina la carrera sin ningún mecanismo de sincronización. Segundo: cuando eso no es posible —porque el registro no cabe, o porque el destino es un archivo abierto en modo :raw donde la escritura va directo al sistema operativo sin proceso intermediario—, hay que introducir un escritor único.
El escritor único es un monitor:
defmodule Proyecto2.Escritor do
@moduledoc """
Monitor que posee el descriptor del archivo. Ningún otro proceso
lo toca: la exclusión mutua es consecuencia de la propiedad exclusiva.
"""
use GenServer
def start_link(ruta, opts \\ []) do
GenServer.start_link(__MODULE__, ruta, opts)
end
@doc "Escribe una línea completa. Asíncrono: no espera confirmación."
def escribir(pid, iodata), do: GenServer.cast(pid, {:escribir, iodata})
@doc "Fuerza el vaciado y devuelve cuántas líneas se escribieron."
def cerrar(pid), do: GenServer.call(pid, :cerrar, :infinity)
@impl true
def init(ruta) do
dispositivo = File.open!(ruta, [:write, :raw, {:delayed_write, 512 * 1024, 2_000}])
{:ok, %{dispositivo: dispositivo, lineas: 0}}
end
@impl true
def handle_cast({:escribir, iodata}, estado) do
:ok = IO.binwrite(estado.dispositivo, iodata)
{:noreply, %{estado | lineas: estado.lineas + 1}}
end
@impl true
def handle_call(:cerrar, _from, estado) do
:ok = File.close(estado.dispositivo)
{:reply, {:ok, estado.lineas}, estado}
end
end
La regla que hace correcto este diseño: el descriptor nunca sale del proceso. Si Escritor expusiera una función que devuelve estado.dispositivo, la garantía se evaporaría, porque cualquiera podría escribir en paralelo. Un monitor es tan fuerte como la disciplina con que se oculta su estado.
Fase 2: procesamiento a output.csv
La segunda fase lee input.txt, transforma cada línea y escribe el CSV. Estructuralmente es el mismo patrón, con la diferencia de que ahora la entrada es un flujo y no un rango.
defmodule Proyecto2.Procesador do
@moduledoc "Transforma input.txt en output.csv de forma concurrente."
@tam_lote 10_000
@doc """
Convierte cada registro de 7 campos en una fila CSV de 4 columnas:
nombre completo, RUT, email, teléfono.
"""
@spec transformar(String.t()) :: iodata()
def transformar(linea) do
[telefono, nom1, nom2, ape1, ape2, email, rut] =
linea |> String.trim_trailing("\n") |> String.split("|")
[Enum.join([nom1, nom2, ape1, ape2], " "), ",", rut, ",", email, ",", telefono, "\n"]
end
def procesar(entrada \\ "input.txt", salida \\ "output.csv", opts \\ []) do
concurrencia = Keyword.get(opts, :concurrencia, System.schedulers_online())
{:ok, escritor} = Proyecto2.Escritor.start_link(salida)
{micros, _} =
:timer.tc(fn ->
entrada
|> File.stream!([], :line)
|> Stream.chunk_every(@tam_lote)
|> Task.async_stream(
fn lote -> Enum.map(lote, &transformar/1) end,
max_concurrency: concurrencia,
timeout: :infinity,
ordered: false
)
|> Enum.each(fn {:ok, filas} -> Proyecto2.Escritor.escribir(escritor, filas) end)
end)
{:ok, lineas} = Proyecto2.Escritor.cerrar(escritor)
{micros, lineas}
end
end
Un detalle de rendimiento con consecuencias de corrección: cada lote transformado se envía al escritor como una sola llamada con toda la iodata del lote, no diez mil llamadas. Eso reduce el tráfico de mensajes en cuatro órdenes de magnitud y, de paso, mantiene íntegras las líneas.
Un detalle de semántica: con ordered: false, las filas del CSV no salen en el mismo orden que las de input.txt. Si el informe compara output.csv con output.py.csv línea a línea, van a diferir aunque ambos sean correctos. La comparación válida es por conjunto: mismo número de filas y mismo multiconjunto de contenidos. Ese punto —que dos salidas distintas puedan ser ambas correctas— es una de las diferencias que el informe del proyecto debe explicar.
Fase 3: la línea base secuencial en Python
El tercer entregable existe para tener contra qué medir. Debe ser deliberadamente simple y estrictamente de un flujo.
#!/usr/bin/env python3
"""Línea base secuencial: procesa input.txt y produce output.py.csv."""
import csv
import sys
import time
def transformar(linea: str) -> list[str]:
"""Convierte un registro de 7 campos en una fila de 4 columnas."""
tel, nom1, nom2, ape1, ape2, email, rut = linea.rstrip("\n").split("|")
return [f"{nom1} {nom2} {ape1} {ape2}", rut, email, tel]
def procesar(entrada: str, salida: str) -> tuple[float, int]:
inicio = time.perf_counter()
filas = 0
with open(entrada, "r", encoding="utf-8", buffering=1024 * 1024) as origen, \
open(salida, "w", encoding="utf-8", newline="", buffering=1024 * 1024) as destino:
escritor = csv.writer(destino)
for linea in origen:
if not linea.strip():
continue
escritor.writerow(transformar(linea))
filas += 1
return time.perf_counter() - inicio, filas
if __name__ == "__main__":
entrada = sys.argv[1] if len(sys.argv) > 1 else "input.txt"
salida = sys.argv[2] if len(sys.argv) > 2 else "output.py.csv"
segundos, filas = procesar(entrada, salida)
print(f"{filas} filas en {segundos:.2f} s ({filas / segundos:,.0f} filas/s)")
Este programa no tiene sincronización porque no tiene concurrencia. Su valor pedagógico es doble: da el denominador de la comparación y demuestra que la versión secuencial nunca falla las pruebas de carrera, lo que confirma que los fallos de la versión concurrente vienen de la concurrencia y no de la lógica de transformación.
Pruebas que exponen la condición de carrera
Llegamos al núcleo del capítulo. Una prueba ordinaria no detecta una condición de carrera porque la carrera depende de un entrelazado que normalmente no ocurre. Escribir pruebas de concurrencia consiste en subir la probabilidad del entrelazado dañino hasta que sea casi seguro.
flowchart TD
START["Código concurrente<br/>bajo prueba"] --> P{"¿Qué técnica<br/>aplico?"}
P --> T1["Volumen<br/>miles de operaciones<br/>por corrida"]
P --> T2["Repetición<br/>la misma prueba<br/>N veces seguidas"]
P --> T3["Interferencia<br/>sleep/yield dentro<br/>de la sección crítica"]
P --> T4["Barrera<br/>todos parten en el<br/>mismo instante"]
P --> T5["Invariantes<br/>propiedad que debe<br/>valer siempre"]
T1 --> V["Verificación"]
T2 --> V
T3 --> V
T4 --> V
T5 --> V
V --> R{"¿Falla?"}
R -->|"Sí, de forma repetible"| OK["La prueba sirve:<br/>documenta el defecto"]
R -->|"No falla nunca"| W["Subir volumen,<br/>agregar interferencia<br/>o revisar el invariante"]
R -->|"Falla a veces"| Z["Prueba inestable:<br/>agregar barrera para<br/>hacerla determinista"]
W --> P
Z --> T4
Las cinco técnicas del diagrama, en detalle:
| Técnica | Qué hace | Cuándo usarla | Riesgo |
|---|---|---|---|
| Volumen | Ejecuta miles de operaciones por corrida | Siempre; es la base | Test lento si se exagera |
| Repetición | Corre la misma prueba N veces en un bucle | Carreras de baja probabilidad | Enmascara si N es bajo |
| Interferencia | Inserta sleep(0) o yield en la ventana | Cuando se conoce dónde está la ventana | Puede ocultar otras ventanas |
| Barrera | Sincroniza el arranque de todos los trabajadores | Para volver determinista un test inestable | Requiere código extra en la prueba |
| Invariantes | Verifica una propiedad global, no un valor puntual | Cuando el resultado exacto varía | Un invariante débil no detecta nada |
El test que caza el contador roto
defmodule Proyecto2.ContadorTest do
use ExUnit.Case, async: false
alias Proyecto2.Contador
@trabajadores 50
@por_trabajador 200
@esperado @trabajadores * @por_trabajador
# Lanza `n` procesos que arrancan todos tras la misma señal (barrera)
# y espera a que todos terminen.
defp en_paralelo(n, fun) do
padre = self()
ref = make_ref()
pids =
for id <- 1..n do
spawn_link(fn ->
send(padre, {ref, :listo, self()})
receive do
{^ref, :arranca} -> :ok
end
fun.(id)
send(padre, {ref, :terminado, self()})
end)
end
# Fase 1: esperar que todos estén listos en la línea de partida.
for _ <- pids, do: assert_receive({^ref, :listo, _}, 5_000)
# Fase 2: soltar a todos en el mismo instante.
for pid <- pids, do: send(pid, {ref, :arranca})
# Fase 3: esperar el término.
for _ <- pids, do: assert_receive({^ref, :terminado, _}, 30_000)
:ok
end
describe "versión con condición de carrera" do
@tag :carrera
test "pierde incrementos bajo concurrencia" do
{:ok, pid} = Contador.Roto.start_link()
en_paralelo(@trabajadores, fn _id ->
for _ <- 1..@por_trabajador, do: Contador.Roto.incrementar(pid)
end)
obtenido = Contador.Roto.valor(pid)
assert obtenido < @esperado,
"""
Se esperaba observar incrementos perdidos.
Esperado teórico: #{@esperado}
Obtenido: #{obtenido}
Si este test falla, la carrera no se manifestó: subir
@trabajadores o @por_trabajador.
"""
end
end
describe "versión con semáforo" do
test "conserva todos los incrementos" do
{:ok, estado} = Contador.ConSemaforo.start_link()
en_paralelo(@trabajadores, fn _id ->
for _ <- 1..@por_trabajador, do: Contador.ConSemaforo.incrementar(estado)
end)
assert Contador.ConSemaforo.valor(estado) == @esperado
end
end
describe "versión monitor" do
test "conserva todos los incrementos" do
{:ok, pid} = Contador.Monitor.start_link()
en_paralelo(@trabajadores, fn _id ->
for _ <- 1..@por_trabajador, do: Contador.Monitor.incrementar(pid)
end)
# El call vacía los casts pendientes por el orden garantizado
# de mensajes entre dos procesos.
assert Contador.Monitor.valor(pid) == @esperado
end
test "resiste 30 repeticiones seguidas" do
for corrida <- 1..30 do
{:ok, pid} = Contador.Monitor.start_link()
en_paralelo(10, fn _ -> for _ <- 1..500, do: Contador.Monitor.incrementar(pid) end)
assert Contador.Monitor.valor(pid) == 5_000, "falló en la corrida #{corrida}"
end
end
end
describe "versión por mensajes" do
test "la suma de parciales iguala el total" do
assert Proyecto2.Recolector.contar(@trabajadores, @por_trabajador) == @esperado
end
end
end
Cuatro cosas que hacen que este archivo sea una prueba de concurrencia y no una prueba ordinaria disfrazada:
- La barrera.
en_paralelo/2no lanza y espera: lanza, espera a que todos lleguen a la línea de partida, y recién entonces los suelta. Sin la barrera, el primer proceso puede terminar antes de que nazca el último y la carrera no ocurre nunca. - El test de la versión rota afirma el fallo. Asegura
obtenido < @esperado. Si algún día ese test empieza a pasar de forma inversa —es decir, el contador roto da el número correcto—, significa que las condiciones dejaron de reproducir la carrera y el resto de la suite perdió poder de detección. async: false. Los tests de concurrencia compiten por los planificadores. Ejecutarlos en paralelo con otros archivos introduce ruido que vuelve inestables los resultados.- Mensajes de fallo que dicen qué hacer. El texto del
assertexplica cómo reaccionar cuando falla, no solo que falló.
Pruebas de integridad sobre los archivos generados
El contador es el caso didáctico; el archivo es el caso real. Estas pruebas verifican invariantes sobre la salida, que es donde una carrera de E/S se manifiesta.
defmodule Proyecto2.ArchivoTest do
use ExUnit.Case, async: false
alias Proyecto2.{Generador, Procesador, Rut}
@registros 20_000
setup do
dir = Path.join(System.tmp_dir!(), "proyecto2-#{System.unique_integer([:positive])}")
File.mkdir_p!(dir)
on_exit(fn -> File.rm_rf!(dir) end)
ruta = Path.join(dir, "input.txt")
Generador.generar(ruta, @registros, concurrencia: 8, tam_lote: 1_000)
{:ok, dir: dir, ruta: ruta}
end
defp campos_de(ruta) do
ruta
|> File.stream!([], :line)
|> Stream.map(fn l -> l |> String.trim_trailing("\n") |> String.split("|") end)
end
test "genera exactamente la cantidad pedida de líneas", %{ruta: ruta} do
assert Enum.count(File.stream!(ruta, [], :line)) == @registros
end
test "toda línea tiene exactamente 7 campos", %{ruta: ruta} do
malformadas =
ruta
|> campos_de()
|> Stream.reject(fn campos -> length(campos) == 7 end)
|> Enum.take(5)
assert malformadas == [],
"líneas con cantidad de campos distinta de 7 (escrituras entrelazadas): #{inspect(malformadas)}"
end
test "todos los RUT generados son consistentes", %{ruta: ruta} do
invalidos =
ruta
|> campos_de()
|> Stream.map(&List.last/1)
|> Stream.reject(&Rut.valido?/1)
|> Enum.take(5)
assert invalidos == []
end
test "la salida concurrente y la secuencial contienen el mismo conjunto", %{dir: dir, ruta: ruta} do
salida = Path.join(dir, "output.csv")
Procesador.procesar(ruta, salida, concurrencia: 8)
esperado =
ruta
|> File.stream!([], :line)
|> Enum.map(fn l -> l |> Procesador.transformar() |> IO.iodata_to_binary() end)
|> MapSet.new()
obtenido = salida |> File.stream!([], :line) |> MapSet.new()
assert MapSet.equal?(esperado, obtenido),
"diferencias: #{inspect(MapSet.difference(esperado, obtenido) |> Enum.take(3))}"
end
end
La última prueba resuelve el problema del orden: compara conjuntos, no secuencias. Es la formulación correcta del invariante cuando se usa ordered: false. Si el proyecto necesitara preservar el orden original, el invariante sería otro y habría que cambiar la implementación —por ejemplo, adjuntando el índice de lote a cada resultado y ordenando antes de escribir—.
La prueba de escritura entrelazada
Para ver la carrera de E/S en acción hace falta el escritor ingenuo que escribe campo por campo:
defmodule Proyecto2.EscrituraTest do
use ExUnit.Case, async: false
@procesos 16
@lineas_por_proceso 500
setup do
ruta = Path.join(System.tmp_dir!(), "escritura-#{System.unique_integer([:positive])}.txt")
on_exit(fn -> File.rm(ruta) end)
{:ok, ruta: ruta}
end
defp campos(id, n) do
["56910000000", "Nom#{id}", "Seg#{id}", "Ape#{id}", "Ap2#{id}",
"u#{id}@correo.cl", "1234567#{rem(n, 10)}-#{rem(id, 10)}"]
end
# Una llamada de E/S por campo: siete mensajes al proceso archivo.
defp escribir_por_campos(dispositivo, id, n) do
for campo <- campos(id, n), do: IO.write(dispositivo, [campo, "|"])
IO.write(dispositivo, "\n")
end
# Una sola llamada de E/S con el registro completo.
defp escribir_de_una(dispositivo, id, n) do
IO.write(dispositivo, [Enum.intersperse(campos(id, n), "|"), "|\n"])
end
defp correr(ruta, escritor_fun) do
dispositivo = File.open!(ruta, [:write, :utf8])
padre = self()
ref = make_ref()
for id <- 1..@procesos do
spawn_link(fn ->
for n <- 1..@lineas_por_proceso, do: escritor_fun.(dispositivo, id, n)
send(padre, {ref, :fin})
end)
end
for _ <- 1..@procesos, do: assert_receive({^ref, :fin}, 30_000)
File.close(dispositivo)
ruta
|> File.stream!([], :line)
|> Stream.map(fn l -> l |> String.trim_trailing("\n") |> String.split("|") end)
|> Enum.count(fn campos -> length(campos) != 8 end)
end
@tag :carrera
test "escribir campo por campo produce líneas mezcladas", %{ruta: ruta} do
rotas = correr(ruta, &escribir_por_campos/3)
assert rotas > 0,
"no se observó entrelazado; aumentar @procesos o @lineas_por_proceso"
end
test "escribir el registro en una sola operación mantiene la integridad", %{ruta: ruta} do
rotas = correr(ruta, &escribir_de_una/3)
assert rotas == 0, "#{rotas} líneas quedaron mal formadas"
end
end
Los dos tests corren exactamente el mismo esquema de concurrencia; lo único distinto es la granularidad de la escritura. Esa es la demostración más limpia del principio general: la atomicidad se define en la unidad de operación que el sistema garantiza, no en la unidad de significado que tiene el programa. Siete campos son una línea para el programador y siete operaciones para el sistema.
La misma solución en Ada
El proyecto original se centra en Elixir, pero implementar el mismo núcleo en Ada obliga a razonar con memoria realmente compartida, que es el modelo que corresponde a hilos en un sistema operativo tradicional. Aquí la carrera es una carrera de verdad sobre una variable en memoria, no una carrera entre mensajes.
with Ada.Text_IO; use Ada.Text_IO;
procedure Carrera_Contador is
N_Tareas : constant := 8;
Por_Tarea : constant := 50_000;
Esperado : constant := N_Tareas * Por_Tarea;
-- Versión sin protección: dos tareas pueden leer el mismo valor.
Global_Sin_Proteger : Integer := 0;
-- Monitor: objeto protegido. El compilador garantiza que solo una
-- tarea ejecute un procedure a la vez sobre la misma instancia.
protected Contador_Seguro is
procedure Incrementar;
function Valor return Integer;
private
Cuenta : Integer := 0;
end Contador_Seguro;
protected body Contador_Seguro is
procedure Incrementar is
begin
Cuenta := Cuenta + 1;
end Incrementar;
function Valor return Integer is
begin
return Cuenta;
end Valor;
end Contador_Seguro;
-- Semáforo contador con barrera de entrada.
protected type Semaforo (Inicial : Natural) is
entry Adquirir;
procedure Liberar;
function Disponibles return Natural;
private
Permisos : Natural := Inicial;
end Semaforo;
protected body Semaforo is
entry Adquirir when Permisos > 0 is
begin
Permisos := Permisos - 1;
end Adquirir;
procedure Liberar is
begin
Permisos := Permisos + 1;
end Liberar;
function Disponibles return Natural is
begin
return Permisos;
end Disponibles;
end Semaforo;
Mutex : Semaforo (1);
-- Tercera variante: la misma variable global, pero protegida por el
-- semáforo binario. El cálculo es idéntico al de la versión insegura.
Global_Con_Semaforo : Integer := 0;
task type Trabajador;
task body Trabajador is
Local : Integer;
begin
for I in 1 .. Por_Tarea loop
-- (1) Versión con carrera.
Local := Global_Sin_Proteger;
delay 0.0; -- punto de replanificación
Global_Sin_Proteger := Local + 1;
-- (2) Versión monitor.
Contador_Seguro.Incrementar;
-- (3) Versión semáforo.
Mutex.Adquirir;
Local := Global_Con_Semaforo;
Global_Con_Semaforo := Local + 1;
Mutex.Liberar;
end loop;
end Trabajador;
begin
-- El bloque interno no termina hasta que todas las tareas del arreglo
-- hayan finalizado: ese es el punto de sincronización.
declare
Equipo : array (1 .. N_Tareas) of Trabajador;
pragma Unreferenced (Equipo);
begin
null;
end;
Put_Line ("Esperado :" & Integer'Image (Esperado));
Put_Line ("Sin proteger :" & Integer'Image (Global_Sin_Proteger));
Put_Line ("Monitor (protected) :" & Integer'Image (Contador_Seguro.Valor));
Put_Line ("Semáforo binario :" & Integer'Image (Global_Con_Semaforo));
Put_Line ("Permisos al final :" & Integer'Image (Mutex.Disponibles));
if Global_Sin_Proteger = Esperado then
Put_Line ("Aviso: la carrera no se manifestó en esta corrida.");
else
Put_Line ("Incrementos perdidos :" &
Integer'Image (Esperado - Global_Sin_Proteger));
end if;
end Carrera_Contador;
Compilación y ejecución con GNAT:
gnatmake carrera_contador.adb -o carrera_contador
./carrera_contador
Dos observaciones sobre el código Ada:
- El bloque
declare ... begin null; end;que contiene el arreglo de tareas es el mecanismo de espera. En Ada, un bloque no puede completarse mientras existan tareas dependientes activas. No hace falta unjoinexplícito. delay 0.0fuerza un punto de replanificación. Sin él, en algunas implementaciones y con optimización activa, la secuencia leer-sumar-escribir puede ejecutarse tan rápido que la carrera casi no aparezca. Eldelayno crea el defecto: hace observable uno que ya existía.
La última línea del reporte, Permisos al final, es una prueba de invariante en sí misma: si el semáforo se inicializó con 1 y todas las adquisiciones tienen su liberación, debe terminar en 1. Cualquier otro valor indica un Liberar de más o de menos.
Medición y comparación
El informe del proyecto exige comparar tiempos. Conviene medir con instrumentación interna y confirmar con una herramienta externa.
defmodule Proyecto2.Bench do
@moduledoc "Barrido de concurrencia para el informe comparativo."
def barrido(total \\ 200_000) do
IO.puts("concurrencia | segundos | registros/s")
for c <- [1, 2, 4, 8, 16] do
{micros, _} = Proyecto2.Generador.generar("bench-#{c}.txt", total, concurrencia: c)
File.rm("bench-#{c}.txt")
s = micros / 1_000_000
IO.puts("#{c} | #{Float.round(s, 2)} | #{round(total / s)}")
end
end
end
Desde la línea de comandos, el tiempo total de cada entregable:
# Generación concurrente
/usr/bin/time -v elixir -e 'Proyecto2.Generador.generar("input.txt", 1_000_000)'
# Procesamiento concurrente
/usr/bin/time -v elixir -e 'Proyecto2.Procesador.procesar("input.txt", "output.csv")'
# Línea base secuencial
/usr/bin/time -v python3 procesar.py input.txt output.py.csv
# Verificación de que ambas salidas tienen las mismas filas
wc -l output.csv output.py.csv
sort output.csv > a.sorted && sort output.py.csv > b.sorted && diff -q a.sorted b.sorted
La comparación sort más diff es la traducción a shell del mismo invariante de conjunto que verificaba el test de ExUnit: dos salidas con distinto orden pero idéntico contenido pasan, dos con contenido distinto no.
Qué esperar del barrido de concurrencia, en términos cualitativos:
| Concurrencia | Comportamiento típico | Factor limitante |
|---|---|---|
| 1 | Referencia; equivale a la versión secuencial | Un solo núcleo activo |
| 2 a 4 | Mejora cercana al factor de paralelismo | Cómputo de generación |
| Igual a la cantidad de núcleos | Punto donde suele estar el mejor rendimiento | Núcleos disponibles |
| Muy por encima de los núcleos | La mejora se aplana o revierte | Cambios de contexto y contención de E/S |
La curva se aplana porque en algún punto el cuello deja de ser la CPU y pasa a ser el disco o el proceso escritor. Identificar dónde ocurre ese cambio, con datos de la propia máquina, es una de las conclusiones que el informe debe contener.
Errores comunes
| Error | Síntoma observable | Causa | Solución |
|---|---|---|---|
Task.async_stream con timeout por omisión | Excepción :timeout al procesar lotes grandes | El valor por omisión es 5000 ms por elemento | Pasar timeout: :infinity o reducir el tamaño del lote |
Componer Agent.get y Agent.update | El contador final es menor que el esperado | Cada llamada es atómica, la secuencia no | Un solo Agent.update(pid, &(&1 + 1)) o un GenServer |
| Escribir un registro con varias llamadas de E/S | Líneas con más o menos campos de los debidos | Otro proceso intercala sus escrituras entre las llamadas | Construir la iodata completa y escribir una sola vez |
Leer el contador justo después de miles de cast | Valor menor al esperado, sin carrera real | Los mensajes siguen en el buzón | Leer con call, que se encola detrás de los cast |
Semáforo sin after en la liberación | El programa se cuelga tras una excepción | El permiso nunca se devolvió | Envolver el cuerpo en try ... after liberar |
Comparar output.csv línea a línea con la salida secuencial | diff reporta miles de diferencias | ordered: false no preserva el orden | Comparar conjuntos, u ordenar ambos archivos antes |
| Test de concurrencia sin barrera de arranque | El test pasa siempre, incluso con el código roto | El primer trabajador termina antes de que arranque el último | Sincronizar el arranque con una señal común |
Test con async: true en la suite de carreras | Resultados que varían entre corridas sin causa clara | Otros archivos de prueba compiten por los planificadores | Marcar async: false en los archivos de concurrencia |
Un solo proceso escritor con call por registro | Rendimiento peor que la versión secuencial | Un viaje de ida y vuelta por cada registro | Enviar lotes completos, no registros sueltos |
max_concurrency muy por encima de los núcleos | El tiempo empeora al subir la concurrencia | Cambios de contexto y contención de disco | Ajustar a System.schedulers_online() y medir |
| Archivos parciales sin orden estable en la unión | El archivo final varía entre corridas | File.ls!/1 no garantiza orden | Ordenar los nombres con relleno de ceros antes de unir |
Olvidar File.close/1 con delayed_write | Faltan las últimas líneas del archivo | El búfer diferido no se vació | Cerrar siempre el dispositivo, idealmente en un after |
| RUT generado sin verificar el dígito | La prueba de validez falla en un pequeño porcentaje | El caso 11 se representa como 0 y el 10 como K | Cubrir ambos casos en el cálculo del dígito |
delay 0.0 ausente en la versión Ada con carrera | La carrera no se manifiesta | La secuencia leer-sumar-escribir es demasiado breve | Insertar un punto de replanificación explícito |
Entregables y verificación
Una lista de comprobación para cerrar el proyecto:
| Entregable | Verificación objetiva |
|---|---|
Generador concurrente de input.txt | wc -l input.txt devuelve 1000000 |
| Todas las líneas bien formadas | La prueba de siete campos pasa sobre el archivo completo |
Procesador concurrente a output.csv | Misma cantidad de filas que la entrada |
| Procesador secuencial en Python | output.py.csv con la misma cantidad de filas |
| Equivalencia entre ambas salidas | sort de ambos archivos más diff sin diferencias |
| Implementación con semáforo | Módulo propio con adquirir/liberar y su suite de pruebas |
| Implementación con monitor | GenServer u objeto protegido con el estado encapsulado |
| Implementación con mensajes | Recolector sin estado compartido y con timeout en receive |
| Pruebas que exponen la carrera | Al menos un test que falla con la versión insegura |
| Medición comparativa | Tabla de tiempos por nivel de concurrencia |
| Informe | Tiempos, dificultades, diferencias entre archivos y por qué existen |
El punto más exigente es el noveno. Un proyecto que solo incluye la versión correcta demuestra que el programa funciona; un proyecto que además incluye la versión rota y la prueba que la caza demuestra que se entiende por qué funciona.
Ejercicios propuestos
-
Semáforo de descriptores. Modifica
Generadorpara que cada trabajador abra su archivo parcial solo después de adquirir un permiso de un semáforo inicializado en 4. Mide si el tiempo total cambia y explica el resultado con la tabla de barrido de concurrencia. -
Invariante del semáforo. Escribe una prueba que, tras completar una carga de trabajo con 200 trabajadores, verifique que
Semaforo.disponibles/1devuelve el valor inicial. Luego elimina elafterdecon_permiso/2, provoca una excepción dentro de la función y observa qué reporta la prueba. -
Orden preservado. Reimplementa
Procesador.procesar/3de modo queoutput.csvrespete el orden deinput.txt, sin renunciar a la concurrencia. Sugerencia: adjunta el índice del lote al resultado y ordena antes de escribir. Mide el costo en tiempo y en memoria de esa garantía adicional. -
Detección de duplicados. Agrega al generador la restricción de que ningún RUT se repita. Resuélvelo de dos formas: con un monitor que guarde el conjunto de RUT ya emitidos, y particionando el espacio de números entre trabajadores sin ninguna sincronización. Compara ambos tiempos.
-
Carrera en Ada sobre un arreglo. Adapta
Carrera_Contadorpara que las tareas escriban en posiciones de un arreglo compartido en lugar de un entero. Determina experimentalmente si escribir en posiciones distintas del mismo arreglo requiere sincronización, y contrasta el resultado con lo que ocurre cuando dos tareas escriben en la misma posición. -
Prueba basada en propiedades. Escribe una prueba que genere lotes de tamaño aleatorio entre 1 y 5000, con concurrencia aleatoria entre 1 y 32, y verifique en cada caso que la cantidad de líneas producidas coincide con la pedida. Ejecútala 100 veces y registra si alguna combinación falla.
-
Interbloqueo provocado. Construye dos semáforos y dos trabajadores que los adquieran en orden opuesto. Escribe una prueba que detecte el interbloqueo mediante un timeout y que reporte cuál trabajador quedó esperando. Después corrige el defecto imponiendo un orden global de adquisición y verifica que la misma prueba pasa.
-
Presión de memoria. Sube
@tam_lotea 500.000 y observa el consumo de memoria con/usr/bin/time -v. Explica por qué construir la iodata de un lote completo antes de escribirlo intercambia memoria por cantidad de operaciones de E/S, y determina cuál es el tamaño de lote con mejor comportamiento en tu máquina.
Lo que queda demostrado
Este proyecto reúne todo lo estudiado sobre concurrencia en una sola pieza de software con salida verificable. El recorrido tuvo una estructura fija: separar el cómputo local del recurso compartido, identificar el punto de contención exacto, resolverlo con las tres familias de sincronización, romperlo a propósito y escribir pruebas cuyo fallo sea reproducible. Ese último paso es el que convierte la concurrencia en algo que se puede afirmar y no solo suponer: una condición de carrera que ninguna prueba puede provocar es indistinguible de un programa correcto, hasta el día en que la carga cambia.
También quedó claro que el mecanismo de sincronización más eficaz suele ser el que se puede evitar. Los archivos parciales eliminaron la carrera de escritura sin ningún candado; la iodata construida de una vez eliminó el entrelazado de campos sin ningún semáforo; la partición del rango eliminó la contención del contador sin ningún monitor. Los mecanismos explícitos quedan para lo que realmente no se puede particionar.
En el capítulo 17 llevamos estas mismas ideas a un escenario donde la contención no se puede evitar por diseño: un servidor de videojuegos con múltiples clientes conectados que comparten un estado de mundo mutable y en tiempo real. Ahí no hay archivos parciales que valgan —el estado es uno solo, todos lo leen y todos lo modifican— y aparecen restricciones nuevas: latencia acotada, tolerancia a fallos de un cliente sin arrastrar al resto, y supervisión de procesos. Ese capítulo cierra el curso y recoge además la bibliografía completa. El índice general está en /tecnologias/sistemas-operativos/00-indice/.