Proyecto 2: un millón de registros, tres mecanismos de sincronización y pruebas que hacen visible la condición de carrera

Por: Artiko
sistemas-operativosadaelixirconcurrenciacondiciones-de-carrerasemaforosmonitorespaso-de-mensajestestingproyecto

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:

  1. Generar input.txt con 1.000.000 de registros sintéticos, en Elixir, con procesamiento concurrente.
  2. Transformar input.txt en output.csv, también en Elixir y también concurrente.
  3. Producir output.py.csv con un programa Python estrictamente secuencial, para tener una línea base contra la cual medir.
  4. 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:

CampoFormatoEjemplo
Teléfonoprefijo 569 más 8 dígitos56912345678
Primer nombretexto alfabéticoCamila
Segundo nombretexto alfabéticoAndrea
Primer apellidotexto alfabéticoRojas
Segundo apellidotexto alfabéticoVergara
Emailusuario más dominio[email protected]
RUT7 u 8 dígitos más dígito verificador18345672-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:

FasePresupuestoNaturaleza del trabajo
Generación de input.txt~6 minutos, con desviación de ~60 sCómputo puro más escritura secuencial
Procesamiento a output.csvhasta ~60 minutosLectura, parseo, transformación y escritura
Línea base secuencial en Pythonsin límite, es la referenciaIgual 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

CriterioSemáforoMonitor (GenServer)Paso de mensajes
Dónde vive el estadoFuera, en una celda aparteDentro del procesoNo hay estado compartido
Quién garantiza la exclusiónEl programador, al parear adquirir/liberarLa estructura del procesoNo aplica
Error típicoOlvidar liberar tras una excepciónComponer dos llamadas creyendo que son unaPerder un mensaje y esperar por siempre
Detección del errorBloqueo permanente o permisos que se agotanValores menores a lo esperadoTimeout en receive
Costo por operaciónDos mensajes al semáforo más la operaciónUn mensajeUn mensaje por resultado, no por operación
Escala con N trabajadoresContención en el semáforoContención en el buzónEscala mientras el trabajo parcial sea grande
Ideal paraLimitar cupos y recursos escasosInvariantes sobre un dato compartidoAgregació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 de Task.async_stream/3 es 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. Con ordered: 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 es System.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écnicaQué haceCuándo usarlaRiesgo
VolumenEjecuta miles de operaciones por corridaSiempre; es la baseTest lento si se exagera
RepeticiónCorre la misma prueba N veces en un bucleCarreras de baja probabilidadEnmascara si N es bajo
InterferenciaInserta sleep(0) o yield en la ventanaCuando se conoce dónde está la ventanaPuede ocultar otras ventanas
BarreraSincroniza el arranque de todos los trabajadoresPara volver determinista un test inestableRequiere código extra en la prueba
InvariantesVerifica una propiedad global, no un valor puntualCuando el resultado exacto varíaUn 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:

  1. La barrera. en_paralelo/2 no 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.
  2. 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.
  3. async: false. Los tests de concurrencia compiten por los planificadores. Ejecutarlos en paralelo con otros archivos introduce ruido que vuelve inestables los resultados.
  4. Mensajes de fallo que dicen qué hacer. El texto del assert explica 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 un join explícito.
  • delay 0.0 fuerza 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. El delay no 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:

ConcurrenciaComportamiento típicoFactor limitante
1Referencia; equivale a la versión secuencialUn solo núcleo activo
2 a 4Mejora cercana al factor de paralelismoCómputo de generación
Igual a la cantidad de núcleosPunto donde suele estar el mejor rendimientoNúcleos disponibles
Muy por encima de los núcleosLa mejora se aplana o revierteCambios 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

ErrorSíntoma observableCausaSolución
Task.async_stream con timeout por omisiónExcepción :timeout al procesar lotes grandesEl valor por omisión es 5000 ms por elementoPasar timeout: :infinity o reducir el tamaño del lote
Componer Agent.get y Agent.updateEl contador final es menor que el esperadoCada llamada es atómica, la secuencia noUn solo Agent.update(pid, &(&1 + 1)) o un GenServer
Escribir un registro con varias llamadas de E/SLíneas con más o menos campos de los debidosOtro proceso intercala sus escrituras entre las llamadasConstruir la iodata completa y escribir una sola vez
Leer el contador justo después de miles de castValor menor al esperado, sin carrera realLos mensajes siguen en el buzónLeer con call, que se encola detrás de los cast
Semáforo sin after en la liberaciónEl programa se cuelga tras una excepciónEl permiso nunca se devolvióEnvolver el cuerpo en try ... after liberar
Comparar output.csv línea a línea con la salida secuencialdiff reporta miles de diferenciasordered: false no preserva el ordenComparar conjuntos, u ordenar ambos archivos antes
Test de concurrencia sin barrera de arranqueEl test pasa siempre, incluso con el código rotoEl primer trabajador termina antes de que arranque el últimoSincronizar el arranque con una señal común
Test con async: true en la suite de carrerasResultados que varían entre corridas sin causa claraOtros archivos de prueba compiten por los planificadoresMarcar async: false en los archivos de concurrencia
Un solo proceso escritor con call por registroRendimiento peor que la versión secuencialUn viaje de ida y vuelta por cada registroEnviar lotes completos, no registros sueltos
max_concurrency muy por encima de los núcleosEl tiempo empeora al subir la concurrenciaCambios de contexto y contención de discoAjustar a System.schedulers_online() y medir
Archivos parciales sin orden estable en la uniónEl archivo final varía entre corridasFile.ls!/1 no garantiza ordenOrdenar los nombres con relleno de ceros antes de unir
Olvidar File.close/1 con delayed_writeFaltan las últimas líneas del archivoEl búfer diferido no se vacióCerrar siempre el dispositivo, idealmente en un after
RUT generado sin verificar el dígitoLa prueba de validez falla en un pequeño porcentajeEl caso 11 se representa como 0 y el 10 como KCubrir ambos casos en el cálculo del dígito
delay 0.0 ausente en la versión Ada con carreraLa carrera no se manifiestaLa secuencia leer-sumar-escribir es demasiado breveInsertar un punto de replanificación explícito

Entregables y verificación

Una lista de comprobación para cerrar el proyecto:

EntregableVerificación objetiva
Generador concurrente de input.txtwc -l input.txt devuelve 1000000
Todas las líneas bien formadasLa prueba de siete campos pasa sobre el archivo completo
Procesador concurrente a output.csvMisma cantidad de filas que la entrada
Procesador secuencial en Pythonoutput.py.csv con la misma cantidad de filas
Equivalencia entre ambas salidassort de ambos archivos más diff sin diferencias
Implementación con semáforoMódulo propio con adquirir/liberar y su suite de pruebas
Implementación con monitorGenServer u objeto protegido con el estado encapsulado
Implementación con mensajesRecolector sin estado compartido y con timeout en receive
Pruebas que exponen la carreraAl menos un test que falla con la versión insegura
Medición comparativaTabla de tiempos por nivel de concurrencia
InformeTiempos, 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

  1. Semáforo de descriptores. Modifica Generador para 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.

  2. Invariante del semáforo. Escribe una prueba que, tras completar una carga de trabajo con 200 trabajadores, verifique que Semaforo.disponibles/1 devuelve el valor inicial. Luego elimina el after de con_permiso/2, provoca una excepción dentro de la función y observa qué reporta la prueba.

  3. Orden preservado. Reimplementa Procesador.procesar/3 de modo que output.csv respete el orden de input.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.

  4. 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.

  5. Carrera en Ada sobre un arreglo. Adapta Carrera_Contador para 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.

  6. 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.

  7. 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.

  8. Presión de memoria. Sube @tam_lote a 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/.