Mostrando entradas con la etiqueta Programación. Mostrar todas las entradas
Mostrando entradas con la etiqueta Programación. Mostrar todas las entradas

viernes, 2 de junio de 2023

Desarrollo seguro: Protección de datos

 La seguridad de los datos debe estar amparada bajo los siguientes preceptos:

  • Autenticación. Se verifica la identidad del usuario para poder acceder al dato.
  • Autorización. El usuario solo puede realizar las operaciones que le permite los privilegios otorgados a sus roles sobre los datos a los que tiene permiso.
  • Confidencialidad. Cuando los datos se almacenan o están en tránsito deben estar protegidos contra la observación o divulgación no autorizada.
  • Integridad. Los datos deben estar protegidos contra manipulaciones por parte de posibles atacantes.
  • Disponibilidad. Los datos siempre han de estar disponibles para los usuarios autorizados. Esto incluye también las políticas de copias de respaldo.

Posibles riesgos

A continuación se describen algunos de los posibles riesgos relativos a la protección de datos:
  • Incumplimiento de normas y leyes relativas al tratamiento de datos personales, lo que supone sanciones legales.
  • Pérdida de información sensible.
  • Pérdida de información sensible de terceros. Esto puede originar acciones legales y sanciones.
  • Pérdida de reputación de la organización.
  • Pérdida de certificaciones.

Mejores prácticas

He aquí algunas recomendaciones para elevar la seguridad en la protección de datos:
  • Identificar toda la información relativa a datos personales y aplicar las políticas de tratamiento de datos personales (acceso, procesamiento, almacenamiento, etc), acorde a las leyes y legislaciones del país.
  • Los datos confidenciales no pueden viajar en los parámetros de las urls. Deben hacerlo en el cuerpo o en las cabeceras HTTP.
  • Los canales de comunicación en los que se transmiten los datos deben utilizar cifrado robusto.
  • Los datos confidenciales deben cifrarse de forma segura antes de su almacenamiento.
  • Evitar que los datos confidenciales se almacenen en la caché de los navegadores.
  • Sobreescribir la memoria con ceros cuando la información confidencial ya no se utilice.
  • Aplicar una política de borrado seguro para los datos que alcanzan el final de su ciclo de vida.



Desarrollo seguro: comunicaciones

La transmisión de información en las aplicaciones es un punto clave de seguridad, por lo que es esencial asegurar la confidencialidad, la integridad y la disponibilidad de las comunicaciones.

Imagen: public domain pictures


Posibles riesgos

A continuación se describen algunos de los posibles riesgos relativos a las comunicaciones:

  • Robo de información confidencial.
  • Interrupción del servicio.
  • Suplantación de identidad.
  • Alteración o destrucción de datos.
  • Propagación de software malicioso.


Mejores prácticas

He aquí algunas recomendaciones para elevar la seguridad en las comunicaciones:

  • Utilizar siempre canales de comunicación cifrados mediante TLS o WebSocket.
  • Usar criptografía robusta. 
  • Utilizar protocolos seguros, como HTTPS o FTPS

Referencias

Desarrollo seguro: transacciones

Las transacciones son operaciones críticas de negocio, como pagos o actividades financieras. Debido a su especial sensibilidad requieren fuertes medidas de seguridad, tales como:

  • Autenticación. Uso de contraseñas seguras y MFA para autorizar el acceso a las transacciones.
  • Criptografía. Asegurar la confidencialidad de las transacciones.
  • Validación. Aseguramiento de la legitimidad de las transacciones, evitando ataques y fraudes.
  • Monitorización. Permite detectar y prevenir de ataques y fraudes en tiempo real.
  • Copias de seguridad. Permite recuperar la disponibilidad de los datos en caso de ataques, caídas del sistema o pérdidas de datos.
Las transacciones deben cumplir las propiedades ACID: Atomicity, Consistency, Isolation y Durability.


Imagen: linnworks

Posibles riesgos

A continuación se describen algunos de los posibles riesgos relativos a las transacciones:
  • Payment Bypass. Vulnerabilidad de configuración que permite a un atacante manipular, en un sistema de pago, los parámetros y la respuesta entre cliente y servidor para eludir el sistema de pago.
  • Fraude. El atacante se beneficia de forma ilícita, manipulando las transacciones o usando información falsa.
  • Robo de información confidencial.  Para ser usada en futuros ataques.
  • Interrupción de las transacciones. Con ataques DoS/DDoS.
  • Pérdida de datos. Puede afectar a la confidencialidad y a la integridad de la información.

Mejores prácticas

He aquí algunas recomendaciones para elevar la seguridad en las transacciones:
  • Evitar un Payment Bypass:
    • Toda compra ha de autorizarse y confirmarse en el servidor.
    • Validar la autenticidad de las firmas durante la comunicación con la pasarela de pago.
    • Verificar el precio correcto en el servidor.
    • Asegurar que los pagos no se reutilizan.
    • Comprobar la fase actual de la transacción en el servidor de pago.
  • Para cada transacción, configurar un tiempo corto de expiración.
  • Llevar una trazabilidad fehaciente de las transacciones.
  • Utilizar comunicaciones con cifrado asimétrico.
  • Registrar el detalle de cada operación, anonimizando los datos sensibles.


Referencias

Desarrollo seguro: gestión de archivos

Los archivos son activos esenciales en la seguridad de un sistema de información, pues contiene y persiste gran parte de la información. Por ello, la gestión de archivos debe contemplarse desde la fase de diseño, identificando los archivos internos, los archivos que se cargan, su ubicación, el control de acceso, la copia de seguridad, el flujo de la información, los datos que contiene, etc.

Imagen: imageapi


Posibles riesgos

A continuación se describen algunos de los posibles riesgos relativos a la gestión de archivos:

  • Acceso no autorizado. Implicaría la revelación, manipulación, pérdida o borrado de datos.
  • Carga de archivos maliciosos. Podría ejecutar archivos de forma remota, provocar una infección de malware o realizar un ataque de denegación de servicio.
  • Carencia de copias de seguridad. Provocaría la pérdida de datos a la organización, la pérdida de datos de usuario o la falta de disponibilidad del servicio. 


Mejores prácticas

He aquí algunas recomendaciones para elevar la seguridad en la gestión de archivos:

  • Exigir la autenticación y la autorización de cada usuario al acceder a un archivo.
  • El nombrado de archivos y directorios no puede ser realizado con entradas de usuario.
  • Validar el tipo de contenido por encima de la extensión del archivo.
  • Verificar el tipo MIME del archivo.
  • Impedir la subida de archivos ejecutables.
  • Para evitar problemas de disponibilidad y de impacto en el servidor, limitar el tamaño del archivo.
  • Procesar los archivos mediante un antivirus y antimalwares.
  • En el directorio de subida de archivos de usuarios, desactivar los privilegios de ejecución.
  • Al proporcionar un enlace de descarga a un usuario, no facilita rutas absolutas (canonización de path).
  • Evitar el nombrado secuencial de archivos.
  • En el nombre de un archivo, no utilizar datos sensibles.
  • El acceso al archivo debe ser de sólo lectura.
  • Poner límite al número de archivos que puede subir un usuario.
  • Para asegurar su integridad, almacenar los hashes de cada archivo subido.

Referencias

Desarrollo seguro: el registro de seguridad

Es imprescindible tener un registro (o log) en un entorno protegido, el cual registre toda actividad o evento del sistema o de las aplicaciones, en lo concerniente a la seguridad. Este registro es vital para un posible análisis forense en un futuro.

El registro de seguridad ha de cumplir los siguientes requisitos:

  • Trazabilidad. Debe tener un formato temporal que facilite la trazabilidad de los eventos.
  • Auditabilidad. Ha de almacenarse de forma segura y durante un tiempo mínimo de retención, a efectos de auditoría.
  • Autenticación/autorización. El acceso a este registro sólo estará disponible para personas autenticadas y autorizadas.
  • Confidencialidad. El registro ha de asegurar que no puede ser accedido por medios distintos a la autenticación y autorización. Los datos sensibles deben ser cifrados.
  • Integridad. El registro debe asegurar que no hay manipulaciones a nivel de registro ni de entradas. Para ello, se recomienda el uso de firmas de integridad que se actualicen con cada nueva entrada.
  • Disponibilidad. Su almacenamiento debe ser redundante y contar con copias de respaldo.


Posibles riesgos

A continuación se describen algunos de los posibles riesgos relativos al registro de seguridad:

  • Fuga de información. Si hubiera vulnerabilidad en la autorización, un atacante podría acceder a la información sensible del registro.
  • Falsificación de registros. Una falta de integridad permitiría a un usuario no autorizado a manipular las entradas del registro, con lo que rompería la trazabilidad. Esto también podría llegar a provocar la ejecución de código malicioso.
  • Eliminación de registros. Un atacante podría eliminar sus propias entradas en el registro, para eliminar su actividad.

Mejores prácticas

He aquí algunas recomendaciones para elevar la seguridad en el registro de seguridad:
  • Evitar registrar datos sensibles (como datos personales, contraseñas, tarjetas de crédito, información financiera, etc). Si se registra, cifrar esta información mediante criptografía, tokens o anonimización.
  • Validar los datos antes de registrar la entrada, a fin de garantizar la información y evitar inyecciones.
  • Registrar el detalle de cada entrada:
    • Intentos de autenticación, sobre todo, los fallidos.
    • Accesos concedidos con roles, incluyendo el usuario.
    • Acceso a datos sensibles: qué usuario, qué roles, qué acciones se han realizado.
    • Errores de validación de las entradas.
    • Excepciones del sistema.
    • Amenazas e intentos de amenazas detectados.
  • Usar funciones hash para verificar la integridad de las entradas.
  • Mantener el registro de seguridad en un entorno protegido e independiente de otros registros.

Referencias



jueves, 11 de marzo de 2021

Idempotencia de producer y consumer en Kafka (exactly-once semantics)

 Kafka trabaja muy bien procesando colas de eventos de forma asíncrona, y es muy eficiente gracias al uso de clusters y brokers, por lo que la alta disponibilidad está (casi) asegurada.

Sin embargo, existen algunas casuísticas por la cuales el flujo de mensajes podría fallar, y eso es imperdonable en un sistema crítico de alto estrés, en el que cada mensaje cuenta. 

Nota: La primera parte del artículo expone los problemas y el planteamiento desde el punto de vista teórico. La segunda parte se abordará un ejemplo usando Java y la librería kafka-clients


Problemas

El orden de los mensajes

Cuando se envían multitud de mensajes, Kafka va a optimizar la eficiencia de tiempos y el rendimiento. Para ello, balancea la carga de las particiones, enviando cada mensaje a una partición que tenga menos carga de trabajo. Cada partición se ejecuta de forma asíncrona y en paralelo, por lo que el orden natural de los mensajes se diluye, y unos mensajes se ejecutan antes o después que otros, desde el punto de vista del orden.

En un escenario en que el orden en que se procesen los mensajes no sea relevante, esta es una solución óptima. En un flujo de logs se podría incluir el timestamp de origen para poder ordenar posteriormente la secuencia de estos logs.

Pero cuando trabajamos en una arquitectura EDA (Event Driven Architecture), el orden de la secuencia de mensajes es primordial, de cara a mantener correctamente el Event-Sourcing, y registrar correctamente la realidad. Este orden es clave para poder mantener y regenerar el estado de los agregados y de las entidades en caso de pérdida.

Para solucionar este problema, hemos de especificar, de forma explícita, una clave (key o ID) a los mensajes. Esto afecta a los mensajes referidos a una entidad en concreto. Por ejemplo, si enviamos mensajes sobre operaciones bancarias, la clave podría ser el número de cuenta bancaria o/y el número de tarjeta de crédito. Con esta clave, Kafka dirige todos los mensajes con la misma clave, siempre a la misma partición, con lo que la secuencia, en principio se mantendría en orden.

Pero puede ocurrir que, si hay una gran carga de trabajo en la red u ocurre algún problema transitorio, un mensaje no se pueda guardar en ese momento por un timeout. Para no perderlo, la propiedad "retries" permite reintentar especificar el número de reintentos para guardar un mensaje fallido más adelante. Por defecto, Kafka lo establece en un valor de 2147483647, por lo que no hay problema si no lo configuramos.

Pero, por otra parte, Kafka establece un número máximo de peticiones al mismo tiempo, por lo que, si uno de los mensajes falla, aunque luego lo reintente, las otras peticiones procesarán los siguientes mensajes. Cuando en el siguiente reintento se guarde el mensaje que falló, al hacerlo lo hará después de n siguientes mensajes, con lo que el orden se rompe.

Para solucionar esto podemos restringir el número de peticiones concurrentes mediante la propiedad "max.in.flight.requests.per.connection". Por defecto, Kafka establece este valor a 5 conexiones, por lo que, para asegurar el orden en caso de un fallo puntual en algún mensaje, estableceremos este valor a 1. 

Lo anterior asegurará el orden de los mensajes, pero la eficacia tiene el coste del rendimiento, el cual disminuye dependiendo del escenario y de la carga. No obstante, esta bajada de rendimiento no es crítica ni demasiado relevante.


La idempotencia

La idempotencia se refiere a la característica de que un mensaje siempre se entrega, y éste se entrega exactamente una vez. Es decir, la idempotencia garantiza la consistencia de datos sin duplicidad.

Suena obvio, pero en un sistema distribuido esta característica es la más compleja de cumplir. Vamos a tratar de explicarlo en detalle y con claridad.


¿Cómo asegura Kafka la entrega de un mensaje?

El flujo de entrega de un mensaje en Kafka, por parte de un producer (productor o emisor de mensajes), se resume en 3 pasos:

  • El producer envía el mensaje al broker de Kafka
  • El broker guarda el mensaje en el siguiente offset en la partición correspondiente
  • El broker retorna un ack (acknwoledgement o reconocimiento) con valor 1, confirmando al producer de que el mensaje se entregó correctamente.  

 Este flujo se representa fácilmente en el siguiente esquema:

Conociendo este flujo, surgen cuatro casuísticas o semánticas.


Semántica "no-guarantee" (sin garantía)

El mensaje se puede procesar una vez, varias veces o ninguna. Este escenario no viene establecido por defecto en Kafka, pues no es habitual. 


Semántica "at-least-once"  (al menos una vez)

En este escenario, si se produce algún problema en la entrega del mensaje, el productor realizará reintentos hasta conseguir que la entrega esté confirmada. 

Aparentemente, esta solución suena bien y puede ser correcta para la mayoría de aplicaciones. Pero existe una casuística que puede generar un problema para conseguir la idempotencia

Si el broker guarda correctamente el mensaje en el offset, pero en el momento de retornar el ack al producer se produce alguna incidencia transitoria, el producer no tendrá la confirmación de entrega del mensaje, aunque en realidad el broker lo haya guardado. El producer reintentará enviar el mensaje nuevamente tantas veces como sea necesario hasta que obtenga el ack de confirmación.

Este escenario asegura que, al menos una vez, el mensaje haya sido enviado con éxito, pero también puede generar mensajes duplicados en el offset de la partición.

El siguiente diagrama ilustra esta casuística:


Semántica "at-most-once" (como máximo una vez)

El offset es comprometido (committed) tan pronto como el lote de mensajes llegue al broker. Si bien este escenario evita las duplicidades, si algo fue mal en el envío, el mensaje se perderá y no podrá ser leído. 

Suele usarse para un enfoque fire-and-forget (dispara y olvídate), es decir, se envía el mensaje a Kafka sin reintentos, ignorando cualquier respuesta desde el broker.

Es recomendable guardar primero el progreso y los datos antes de enviar el mensaje a Kafka


Semántica "exactly-once" (exactamente una vez)

En este escenario, es necesario asegurar que el mensaje sea enviado a Kafka sin pérdida ni duplicidades. Es recomendable que el envío de mensajes se realice dentro de una transacción, a fin de asegurar la llegada de los mensajes, su secuencia y la unicidad de los mensajes, para que el consumidor de Kafka pueda operar correctamente con los offsets



La idempotencia en el producer

¿Cómo conseguir la idempotencia en el producer?

Para que se cumpla la idempotencia en el producer, será tan simple como configurar las siguientes propiedades:

  • enable.idempotence=true
  • acks=all
  • retries>0
  • max.in.flight.requests.per.connection<=5

Con la configuración anterior, ya tenemos nuestro producer listo para trabajar bajo la semántica de la idempotencia.

Pero si queremos que nuestro consumer también trabaje de forma idempotente, debemos ayudarle desde el producer

Tal y como lo hemos dejado ahora, el consumer debería añadir un esfuerzo extra en trabajar en el control y gestión de los offsets de forma manual para conseguir la idempotencia.


Para evitar este trabajo, desde el producer vamos realizar toda la operativa de envíos atómicos de mensajes dentro de una transacción. Para que esta transacción atómica sea más efectiva, definiremos la siguiente propiedad:

  • transactional.id=<id_transaccion>

Esta propiedad define una parte de la clave que utilizará el algoritmo de transacción para genera el id de la transacción junto a la secuencia.

La transacción asegurará que todos los mensajes de la transacción sean enviados correctamente, definiendo la secuencia y el offset preparado para el consumidor. 

De forma similar a una base de datos, una transacción prepara y ejecuta múltiples operaciones de escritura. En el caso de que ocurra algún problema, no se generan cambios en el offset y se informa del problema al producer, para que pueda decidir si volver a intentar la transacción o realizar alguna otra operación. Por tanto, si todo es correcto, la entrega los mensajes está asegurada sin incoherencias. Además, por parte del consumer podemos especificar que lea del offset solamente los mensajes una vez finalizada (y asegurada) la transacción, evitando leer mensajes aún en tránsito.



Planteamiento del ejemplo

Para el ejemplo, usaremos Java y la librería kafka-clients. Para configurar el pom.xml de Maven, añadiremos las siguientes dependencias:

<dependencies>

   <dependency>

      <groupId><org.apache.kafka/groupId>

      <artifactId>kafka-clients</artifactId>

      <version>2.3.0</version>

   </dependency>


   <dependency>

      <groupId>org.slf4j</groupId>

      <artifactId>slf4j-simple</artifactId>

      <version>1.6.1</version>

   </dependency>

</dependencies>


En este ejemplo vamos a simular el envío masivo de eventos (mensajes) de varios sensores de temperatura de un motor crítico (por ejemplo, el de un avión o el de un coche de Fórmula 1). Esta simulación enviará mil eventos, eligiendo aleatoriamente un determinado sensor (de entre 10 posibles), con una temperatura aleatoria entre 100 y 300 grados centígrados.


Código del producer

El siguiente código pone en práctica los conceptos vistos anteriormente. 

Nota: Antes de ejecutarlo, hemos de ejecutar primero el consumer, para que éste escuche los mensajes que genere el producer.


import java.text.SimpleDateFormat;

import java.util.Date;

import java.util.Properties;

import java.util.Random;

import java.util.TimeZone;


import org.apache.kafka.clients.producer.KafkaProducer;

import org.apache.kafka.clients.producer.Producer;

import org.apache.kafka.clients.producer.ProducerRecord;

import org.slf4j.Logger;

import org.slf4j.LoggerFactory;


/**

 * IdempotentProducer example

 * @author: Rafael Hernamperez

 *

 */

public class IdempotentProducer {

   private static final int numberOfSensors = 10;

   private static final int numberOfEvents = 1000;

   private static final String eventTopic = "temperature-changed";

   private static Logger log = LoggerFactory.getLogger(IdempotentProducer.class) {


   public static void main(String[] args) {

      Random rand = new Random();


      // Producer properties

      Properties props = new Properties();

      props.put("bootstrap.servers", "localhost:9092");

      props.put("acks", "all");

      props.put("linger.ms", "10");

      props.put("retries", "50");

      props.put("max.in.flight.requests.per.connection", "1");

      props.put("enable.idempotence", "true");

      props.put("transacional.id", "temp-changed");

      props.put("key.serializer", "org.apache.kafka.common.serialization-StringSerializer");

      props.put("value.serializer", "org.apache.kafka.common.serializztion-StringSerializer");


      try (Producer<Sring, String> producer = new KafkaProducer<>(props)) {

         try {

            // Start transaction

            producer.initTransactions();

            producer.beginTransaction();


            // Generate messages

            for (int event = 0; event < numberOfEvents; event++) {

               int sensor = rand.nextInt(numberOfSensors);


               // Temperature simulation (from 100º to 300º)

               int intTemp = 100 + rand.nextInt(200);

               double decTemp = rand.nextDouble() + intTemp;


               // Current timestamp in ISO 8601 format

               Date date = new Date(System.currentTimeMillis());

               SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss.SSSXXX");

               sdf.setTimeZone(TimeZone.getTimeZone("CET"));


               // Compose message in CSV format: datetime;sensordId;temperature

               String eventValue = String.format("%s;%d;%.1f", sdf.format(date), sensor, decTemp);


               // Compose the key using the sensorId

               String keyEvent = String.format("tempSensor-%s", sensor);

               

               // Send the event (message)

               producer.send(new ProducerRecord<String, String>(eventTopic, keyEvent, valueEvent));

               logInfo("Event #" + event + " sent: " + valueEvent);

            }  // for

            

            // Commit transaction

            producer.commitTransaction();

         } catch (Exception e) {

            log.error("Error in transaction: " + e.getMessage());


            // Abort Transaction

            producer.abortTransaction();

         }

      } catch (Exception e) {

         log.error("Error sending event: " + e.getMessage());

      }

   }  // main()

}  // class


La idempotencia en el consumer

¿Cómo conseguir la idempotencia en el consumer?

Una vez resuelta la idempotencia en el producer, tal y como hemos desarrollado anteriormente, la parte de la idempotencia en el lado del consumer es mucho más sencilla.

Para ello, vamos a establecer las siguientes propiedades clave:

  • isolation.level="read_committed"

El nivel de aislamiento establecido como "read_committed" indica que el consumer solamente aborde los mensajes cuyo offset esté comprometido, es decir, que tenga la marca committed. Esto solamente ocurrirá cuando la transacción haya terminado exitosamente, por lo que solo aplica a los mensajes de la transacción.

Además, la transacción habrá generado la secuencia correcta de los mensajes para que el consumer los procese en el orden correspondiente.


Planteamiento del código

Para el lado del consumer, vamos a generar una aplicación multihilo que lance en paralelo varios consumers que estén escuchando al mismo tiempo y procesen los mensajes.

Para ello, habrá dos archivos de código:

  • ThreadConsumer.java: Contiene el código de cada hilo consumer.
  • IdempotentConsumer.java: Aplicación principal. Prepara el entorno multihilo, instancia cada hilo y lanza los consumers.


Código del hilo para consumer

El siguiente código corresponde a cada hilo de consumer. En un entorno multihilo, habrá varios procesos ejecutándose al mismo tiempo. Este es el código de cada uno de estos hilos.


import java.time.Duration;

import java.util.Arrays;


import org.apache.kafka.clients.consumer.ConsumerRecord;

import org.apache.kafka.clients.consumer.ConsumerRecords;

import org.apache.kafka.clients.consumer.KafkaConsumer;

import org.apache.kafka.commons.errors.WakeupException;

import org.slf4j.Logger;

import org.slf4j.LoggerFactory;


/**

 * ThreadConsumer 

 * @author: Rafael Hernamperez

 *

 */

public class ThreadConsumer extends Thread {

   private final KafkaConsumer<String, String> consumer;

   private final String eventTopic = "temperature-changed";

   private boolean closed = false;

   private Logger log = LoggerFactory.getLogger(ThreadConsumer.class);


   // Constructor

   public ThreadConsumer(KafkaConsumer<String, String> consumer) {

      this.consumer = consumer;

   }


   // Main thread execution

   @Override

   public void run() {

      consumer.subscribe(Arrays.asList(eventTopic));


      try {

         while(!closed) {

            ConsumerRecords<String, String> events = consumer.poll(Duration.ofMillis(100));


            for (ConsumerRecord<String, String> event : events) {

               log.info("Partition={}, Offset={}, key={}, value={}",

                  event.partition(),

                  event.offset(),

                  event.key(),

                  event.value());

            }  // for

         }  // while

      } catch(WakeupException Exception we) {

         if (!closed) {

            throw we;

         }

      } finally {

         consumer.close();

      }

   }  // run()

} // class


Código principal para el consumer

El código principal del consumer se encargará de preparar un entorno multihilo, en el que varios consumers trabajarán al mismo tiempo para escuchar y leer los mensajes que genere el producer. Al crear un hilo, se le pasará un objeto de tipo KafkaConsumer, inicializado con las propiedades de ese consumer, las cuales ya están preparadas para su funcionamiento con una semántica exactly-once, lo que facilitará la idempotencia.


import java.util.Properties;

import java.util.concurrent.ExecutorService;

import java.util.concurrent.Executors;

import org.apache.kafka.clients.consumer.KafkaConsumer;


/**

 * IdempotentConsumer example

 * @author: Rafael Hernamperez

 *

 */

public class IdempotentProducer {

   public static final int numberOfConsumers = 5;


   public static void main(String[] args) {

      // Consumer properties

      Properties props = new Properties();

      props.put("bootstrap.servers", "localhost:9092");

      props.put("group.id", "temp-event-group");

      props.put("enable.auto.commit", "false");

      props.put("isolation.level", "read_committed");

      props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");

      props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");


      ExecutorService executyor = Executors.newFixedThreadPool(numberOfConsumers);


      // Generate consumers in a multithread environment

      for (int i=0; i<numberOfConsumer; i++) {

         ThreadConsumer consumer = new ThreadConsumer(new KafkaConsumer<>(props));

         executor.execute(consumer);

      }


      // Execution will alive until last thread is terminated

      while(!executor.isTerminated());


   }  // main()

}  // class



Configuraciones clave

A continuación se resumen las propiedades y valores clave para cada caso.


Consumer

  • No-guarantee:
    • enable.auto.commit = true
    • auto.commit.interval.ms = <frecuencia_milisegundos>
  • At-least-once:
    • enable.auto.commit = false
    • O enable.auto.commit=true con auto.commit.interval.ms con valor muy alto.
    • Recomendado utilizar consumer.commitSync() para controlar los commits en los offset de los mensajes.
  • At-most-once:
    • enable.auto.commit = true
    • auto.commit.interval.ms = <frecuencia_milisegundos>
    • No usar el método consumer.commitSync()
  • Exactly-once:
    • enable.auto.commit = false
    • isolation.level = "read_committed"

Producer

  • retries > 0
  • acks = "all"
  • max.in.flight.requests.per.connection <= 5
  • enable.idempotence = true
  • transacional.id = <clave_unica_transaccion>

Para facilitar la idempotencia real tanto en producer como en consumer, realizar el envío de mensajes dentro de una transacción.


Enlaces de referencia






















domingo, 7 de marzo de 2021

Las 21 mejores tipografías para programar

Cuando programas, ¿te aburre la tipografía por defecto de tu IDE o de tu editor de textos?

A continuación comparto una selección de 21 tipografías que te sacarán de la monotonía y te harán disfrutar de la programación.

Esta selección es particular, no un ranking universalmente aceptado. Te invito a dejar un comentario con tu opinión y a aportar otras tipografías que te gusten para esta selección.


CP Mono




Cutive Mono






Exo






Exo 2





Fantasque Sans Mono




JetBrains Mono




Jura



Lekton





Menlo


Monaco



Monolisa




Monospace Typewriter




Overpass Mono



PT Mono


Roboto Mono




SaxMono



SF Mono




Skyhook Mono




Tipografías extra


La lista anterior se quedó pequeña, por lo que se añaden nuevas tipografías.

Anonymous Pro




Courier Prime Code


Everson Mono





Space Mono




Enlaces de interés



lunes, 15 de febrero de 2021

Cómo crear un productor de streaming de eventos para Kafka en Java



Aviso: Este pequeño tutorial no va a realizar una introducción a Kafka. Se asume que el lector ya tiene unas nociones sobre Kafka y quiere comenzar a desarrollar código en Java para enviar eventos a un topic de Kafka. Esto es habitual realizarlo, principalmente, desde microservicios o desde scripts batch.

Nota: El concepto de estado utilizado en este artículo está contextualizado en una EDA (Arquitectura Orientada a Eventos), y representa un estado originado por un evento. Dentro de Kafka se asume este concepto como record o registro, y se refiere al valor que almacenará Kafka en el bus de eventos.

Requisitos

Se asume que ya se dispone de un entorno de Kafka funcionando, y que en dicho entorno existe, al menos, un topic sobre el cual escribir eventos. Dicho entorno puede estar en la propia máquina de desarrollo (localhost), en un servidor dedicado o en un servidor en cloud.

Más adelante (en otro artículo), veremos cómo desarrollar código para suscribirnos a ese topic y poder responder, en consecuencia, a dichos eventos. Ese código corresponderá a la parte del consumidor o suscriptor. En este artículo nos centraremos exclusivamente en la parte del productor o publicador.

En el entorno de desarrollo se recomienda tener lo siguiente:
  • Java JDK versión 11 (he utilizado la versión 15)
  • Gestión de dependencias con Maven
  • El IDE que utilizado ha sido el de Spring Tools, pero con Eclipse o IntelliJ IDEA no debería haber problemas.
  • Dependencias de Kafka en Maven versión 2.7.0

Configuración de Maven


En el archivo pom.xml del proyecto Java, añadir la dependencia encargada de importar las librerías necesarias para trabajar con streams de eventos en Kafka:

<dependencies>
    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
        <version>2.7.0</version>
    </dependency>
</dependencies>

Configuración de las propiedades

El primer paso a realizar en el código, será definir las propiedades para poder configurar la conexión a Kafka. Las más importantes son las siguientes:

// Propiedades del producer
Properties props = new Properties();

// Lista de servidores Kafka a los que conectarse
props.put("bootstrap.servers", "localhost:9092");

// Serializacion de los datos de la clave (key)
props.put("key.serializer", StringSerializer.class.getName());

// Serializacion de los datos del valor (value)
props.put("value.serializer", StringSerializer.class.getName());

Para el ejemplo, usaremos el tipo String, que es el que viene por defecto. Kafka permite definir estructuras de datos mediante JSON y Avro (este es el preferido), que se definen con el objeto Serdes (SERialize/DESerialize).

A continuación se exponen algunas propiedades no tan relevantes ahora (son opcionales), pero que se podrán usar en un futuro para tunear la configuración de la conexión a Kafka:

props.put("acks", "all");
props.put("retries", 0);
props.put("batch.size", 16384);
props.put("linger.ms", 1);
props.put("buffer.memory", 33554432);


Objeto KafkaProducer

El objeto KafkaProducer permite crear un cliente para una conexión a un topic de Kafka, a partir de la información proporcionada en las propiedades descritas anteriormente.

KafkaProducer<String, String> kp = new KafkaProducer<String, String>(props);

Este objeto permite definir el tipo de datos de la clave (primer parámetro) y del valor o estado (segundo parámetro). En nuestro ejemplo utilizaremos el tipo String (si se desea trabajar con tipos customizados, ver los tipos Serdes, y cómo definir tipos en JSON o Avro).

Nota: Este objeto crea un cliente genérico a un servidor de Kafka, por lo que se puede utilizar posteriormente para enviar estados a streams de eventos a diferentes topics.

Preparación del estado (registro) a enviar

Para enviar un estado o registro a Kafka desde el productor, es necesario preparar éste mediante un objeto de tipo ProducerRecord:

ProducerRecord<String, String> pr = new Producer<String, String>(topic, key, estado);

Este objeto permite definir el tipo de datos de la clave (primer parámetro) y del valor o estado (segundo parámetro). En nuestro ejemplo utilizaremos el tipo String (si se desea trabajar con tipos customizados, ver los tipos Serdes, y cómo definir tipos en JSON o Avro).

En la construcción (entre paréntesis) se pasarán los valores correspondientes a:
  • topic: Valor del nombre del topic a usar en Kafka, donde se enviará el evento.
  • key: Valor de la clave (key) del evento o registro.
  • estado: Valor del estado o registro a enviar.
El valor de la clave puede ser opcional en otros contextos de Kafka (se almacenaría como null). Por ello, también permitiría la siguiente sintaxis:

ProducerRecord<String, String> pr = new Producer<String, String>(topic, estado);


Envío del estado a Kafka


Una vez preparado el estado, éste se envía a Kafka a través del objeto KafkaProducer definido al principio, pasándole el objeto ProducerRecord con la información del estado:

kp.send(pr);

El método send() envía el estado o registro al topic especificado.


Cerrar la conexión a Kafka


Cuando nuestro código no va a enviar más estados a Kafka, debemos cerrar el objeto KafkaProducer, para liberar recursos del stream, así como la conexión. Para ello, utilizaremos el método close():

kp.close();


Aplicación de ejemplo de productor Kafka


A continuación os dejo una aplicación completa que hace de productor Kafka.

En este ejemplo, se ejecuta desde la consola de comandos como un script, pero la base os servirá también para microservicios u otro tipo de aplicaciones.

Lo primero que hará será preguntar por el topic al cual queremos enviar los estados. Después, en un bucle, solicitará el valor del estado a enviar. Dicho estado es un texto libre, por lo que se puede introducir cualquier valor, incluso un JSON en formato String.

Este bucle se repetirá hasta que el usuario introduzca el valor 'quit' (sin comillas). En ese momento se cerrará la conexión y terminará la ejecución.


package com.rhernamperez.kafkastreamsdemo;

import java.io.BufferedReader;
import java.io.InputStreamReader;
import java.util.Date;
import java.util.Properties;

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;
//import org.apache.kafka.streams.StreamsConfig;

/**
* Demostracion de un productor Kafka que envia eventos a un stream
* @author rafinguer
*
*/
public class KafkaProducerDemo {

    public static void main(String[] args) {

        // Propiedades del producer
        Properties props = new Properties();

        props.put("bootstrap.servers", "localhost:9092");
        props.put("key.serializer", StringSerializer.class.getName());
        props.put("value.serializer", StringSerializer.class.getName());
        props.put("acks", "all");
        props.put("retries", 0);
        props.put("batch.size", 16384);
        props.put("linger.ms", 1);
        props.put("buffer.memory", 33554432);

        // Creacion del objeto productor
        KafkaProducer<String, String> kp = new KafkaProducer<String, String>(props);
        String estado="", topic="", key = "";

        // Mensaje de bienvenida
        System.out.println("Demostracion de Kafka producer. Introduce datos para cada evento. 'quit' para salir\n");

        // Introduccion del topic por teclado
        BufferedReader br = new BufferedReader(new InputStreamReader(System.in));

        System.out.println("Introduce el nombre del topic: ");

        try {
            topic = br.readLine();
        } catch (Exception e) {
            System.out.println("ERROR > " + e.getMessage());
        }

        // Bucle para introducir estados hasta que se escriba 'quit'
        while(!estado.equals("quit")) {

            try {
                // Lectura de estados por teclado
                System.out.print(">>> ");
                estado = br.readLine();

                if (estado.equals("quit")) continue;

                // Envio del estado. La key sera la fecha y hora actuales
                key = new Date().toString();
                ProducerRecord<String, String> pr = new ProducerRecord<String, String>(topic, key, estado);
kp.send(pr);

                System.out.println("Enviado a topic " + topic + " la clave " + key + " con el estado > " + estado + "\n");
            } catch(Exception e) {
                System.out.println("Error > " + e.getMessage());
            }
        }

        // Cerrar el stream
        kp.close();

        System.out.println("**** FIN ****");
    }

}


Enlaces de interés