Cerraduras distribuidas en ZooKeeper con la API nativa y con Apache Curator

  1. Implementación de una cerradura dsitribuida con la API nativa de ZooKeeper

1.1. Enfoque basado en la unicidad de un nodo

ZooKeeper no permite crear dos veces el mismo nodo. Esta restricción se puede explotar directamente como mecanismo de exclusión mutua: quien consigue crear el nodo, posee la cerradura.

Flujo

  1. Al solicitar la cerradura se intenta crear el nodo que la representa.
  2. Si la creación falla, el nodo ya existe y la cerradura está ocupada: se registra un watcher sobre ese nodo y se espera.
  3. Si la creación tiene éxito, la cerradura ha sido adquirida.
  4. La liberación voluntaria consiste en borrar dicho nodo.
  5. Si la sesión que posee la cerradura expira o se desconecta, el nodo se elimina automáticamente porque es efímero.
  6. Cuando el nodo desaparece, los watchers registrados despiertan a sus hilos y estos vuelven a competir por la cerradura.

Código

public class CerraduraNodoExclusivo implements Cerradura {

    private static final String RAIZ = "/cerraduras_distribuidas";

    private final ZooKeeper zk;
    private final String nombreRecurso;

    public CerraduraNodoExclusivo(ZooKeeper zk, String nombreRecurso) throws KeeperException, InterruptedException {
        this.zk = zk;
        this.nombreRecurso = nombreRecurso;
        crearNodoPadre(RAIZ);
    }

    private void crearNodoPadre(String ruta) throws KeeperException, InterruptedException {
        try {
            zk.create(ruta, new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
        } catch (KeeperException e) {
            if (e.code() != KeeperException.Code.NODEEXISTS) {
                throw e;
            }
        }
    }

    @Override
    public void adquirir() throws KeeperException, InterruptedException {
        String ruta = RAIZ + "/" + nombreRecurso;

        while (true) {
            try {
                zk.create(ruta, new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL);
                return;
            } catch (KeeperException e) {
                if (e.code() != KeeperException.Code.NODEEXISTS) {
                    throw e;
                }
            }

            CountDownLatch avisoLiberacion = new CountDownLatch(1);
            Stat estado = zk.exists(ruta, evento -> {
                if (evento.getType() == Watcher.Event.EventType.NodeDeleted) {
                    avisoLiberacion.countDown();
                }
            });

            if (estado == null) {
                // El nodo desapareció antes de registrar el watcher: reintentar la creación.
                continue;
            }

            avisoLiberacion.await();
        }
    }

    @Override
    public void liberar() throws KeeperException, InterruptedException {
        zk.delete(RAIZ + "/" + nombreRecurso, -1);
    }
}

Prueba

public class PruebaCerraduraNodoExclusivo {

    private static final String CADENA_CONEXION = "10.5.31.155:2181";
    private static final String RECURSO = "recurso_exclusivo";
    private static final int NUM_HILOS = 9;

    @Test
    public void pruebaExclusionMutua() throws Exception {
        CountDownLatch conectado = new CountDownLatch(1);
        ZooKeeper zk = new ZooKeeper(CADENA_CONEXION, 60000, evento -> {
            if (evento.getState() == Watcher.Event.KeeperState.SyncConnected) {
                conectado.countDown();
            }
        });
        conectado.await();

        CountDownLatch fin = new CountDownLatch(NUM_HILOS);
        ExecutorService pool = Executors.newFixedThreadPool(NUM_HILOS);

        for (int i = 1; i <= NUM_HILOS; i++) {
            final int id = i;
            pool.submit(() -> {
                try {
                    Cerradura cerradura = new CerraduraNodoExclusivo(zk, RECURSO);
                    cerradura.adquirir();
                    System.out.println("hilo-" + id + " obtuvo la cerradura");
                    Thread.sleep(2000);
                    cerradura.liberar();
                    System.out.println("hilo-" + id + " liberó la cerradura");
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    fin.countDown();
                }
            });
        }

        fin.await();
        pool.shutdown();
        zk.close();
    }
}

Ventajas

  1. Implementación sencilla y reutilizable tal cual.
  2. El mecanismo de notificación ofrece tiempos de respuesta bajos.
  3. Los nodos efímeros garantizan que la cerradura se libere aunque el proceso caiga.

Inconvenientes

Provoca el efecto estampida (thundering herd): al eliminarse un nodo, todos los hilos suscritos a la eliminación de ese nodo reciben la notificación y se desipertan a la vez, aunque solo uno podrá obtener la cerradura.

1.2. Enfoque basado en nodos secuenciales

Para evitar el efecto estampida del enfoque anterior se aprovecha la creación de nodos secuenciales: cada solicitante ocupa una posición en una cola ordenada y solo vigila al inmediatamente anterior.

  1. Al solicitar la cerradura se crea un nodo temporal y secuencial.
  2. Se comprueba si ese nodo es el menor de todos. Si lo es, la cerradura se ha adquirido; en caso contrario, se registra un watcher únicamente sobre el nodo inmediatamente anterior y se espera.
  3. La liberación voluntaria consiste en borrar el nodo propio.
  4. Si la sesión que posee la cerradura expira o se desconecta, el nodo se elimina automáticamente porque es efímero.
  5. Cuando el nodo vigilado desaparece, el siguiente hilo vuelve a comprobar si ya es el menor y repite el paso 2.

Código

public class CerraduraNodoSecuencial implements Cerradura {

    private static final String RAIZ = "/cerraduras_distribuidas";
    private static final String PREFIJO = "cerradura-";

    private final ZooKeeper zk;
    private final String nombreRecurso;
    private String rutaPropia;

    public CerraduraNodoSecuencial(ZooKeeper zk, String nombreRecurso) throws KeeperException, InterruptedException {
        this.zk = zk;
        this.nombreRecurso = nombreRecurso;
        crearNodoPadre(RAIZ);
        crearNodoPadre(RAIZ + "/" + nombreRecurso);
    }

    private void crearNodoPadre(String ruta) throws KeeperException, InterruptedException {
        try {
            zk.create(ruta, new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
        } catch (KeeperException e) {
            if (e.code() != KeeperException.Code.NODEEXISTS) {
                throw e;
            }
        }
    }

    @Override
    public void adquirir() throws KeeperException, InterruptedException {
        String base = RAIZ + "/" + nombreRecurso;

        if (rutaPropia == null) {
            rutaPropia = zk.create(base + "/" + PREFIJO, new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL);
        }

        String propia = rutaPropia.substring(base.length() + 1);

        while (true) {
            List<String> cola = new ArrayList<>(zk.getChildren(base, false));
            Collections.sort(cola);

            int posicion = cola.indexOf(propia);
            if (posicion == 0) {
                // No hay nadie delante: la cerradura es nuestra.
                return;
            }

            String predecesora = cola.get(posicion - 1);
            CountDownLatch avisoLiberacion = new CountDownLatch(1);
            Stat estado = zk.exists(base + "/" + predecesora, evento -> {
                if (evento.getType() == Watcher.Event.EventType.NodeDeleted) {
                    avisoLiberacion.countDown();
                }
            });

            if (estado == null) {
                // La predecesora ya no existe: volver a evaluar la cola.
                continue;
            }

            avisoLiberacion.await();
        }
    }

    @Override
    public void liberar() throws KeeperException, InterruptedException {
        if (rutaPropia != null) {
            zk.delete(rutaPropia, -1);
            rutaPropia = null;
        }
    }
}

Prueba

public class PruebaCerraduraNodoSecuencial {

    private static final String CADENA_CONEXION = "10.5.31.155:2181";
    private static final String RECURSO = "recurso_secuencial";
    private static final int NUM_HILOS = 9;

    @Test
    public void pruebaExclusionMutua() throws Exception {
        CountDownLatch conectado = new CountDownLatch(1);
        ZooKeeper zk = new ZooKeeper(CADENA_CONEXION, 60000, evento -> {
            if (evento.getState() == Watcher.Event.KeeperState.SyncConnected) {
                conectado.countDown();
            }
        });
        conectado.await();

        CountDownLatch fin = new CountDownLatch(NUM_HILOS);
        ExecutorService pool = Executors.newFixedThreadPool(NUM_HILOS);

        for (int i = 1; i <= NUM_HILOS; i++) {
            final int id = i;
            pool.submit(() -> {
                try {
                    Cerradura cerradura = new CerraduraNodoSecuencial(zk, RECURSO);
                    cerradura.adquirir();
                    System.out.println("hilo-" + id + " obtuvo la cerradura");
                    Thread.sleep(2000);
                    cerradura.liberar();
                    System.out.println("hilo-" + id + " liberó la cerradura");
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    fin.countDown();
                }
            });
        }

        fin.await();
        pool.shutdown();
        zk.close();
    }
}
  1. Implementación de una cerradura distribuida con Apache Curator

Todo lo anterior supone reimplementar un mecanismo que ya existe. Curator es un cliente de ZooKeeper de código abierto que ofrece un nivel de abstracción superior al cliente nativo y reduce notablemente el código necesario. Con Curator, obtener una cerradura distribuida se reduce a unas pocas líneas, ya que implementa internamente el algoritmo de nodos secuanciales con vigilancia únicamente sobre el predecesor.

public class PruebaCerraduraCurator {

    private static final String CADENA_CONEXION = "10.5.31.155:2181";
    private static final String RUTA_CERRADURA = "/cerraduras_distribuidas/recurso_curator";
    private static final int NUM_HILOS = 9;

    @Test
    public void pruebaExclusionMutua() throws Exception {
        CuratorFramework cliente = CuratorFrameworkFactory.newClient(CADENA_CONEXION, new RetryNTimes(3, 5000));
        cliente.start();

        CountDownLatch fin = new CountDownLatch(NUM_HILOS);
        ExecutorService pool = Executors.newFixedThreadPool(NUM_HILOS);

        for (int i = 1; i <= NUM_HILOS; i++) {
            final int id = i;
            pool.submit(() -> {
                InterProcessMutex cerradura = new InterProcessMutex(cliente, RUTA_CERRADURA);
                try {
                    cerradura.acquire();
                    System.out.println("hilo-" + id + " obtuvo la cerradura");
                    Thread.sleep(2000);
                    cerradura.release();
                    System.out.println("hilo-" + id + " liberó la cerradura");
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    fin.countDown();
                }
            });
        }

        fin.await();
        pool.shutdown();
        cliente.close();
    }
}

Etiquetas: Zookeeper Apache Curator distributed lock ephemeral sequential nodes Watcher

Publicado el 9-29 20:02