¡Hola, Habr! Hoy construiremos un sistema que utilizará Spark Streaming para procesar flujos de mensajes de Apache Kafka y guardar el resultado en una base de datos en la nube AWS RDS.
Imaginemos que una institución crediticia nos encarga procesar las transacciones entrantes 'en tiempo real' a través de todas sus sucursales. Esto puede hacerse con el fin de calcular de manera rápida la posición abierta en moneda extranjera para la tesorería, los límites o el resultado financiero de las transacciones, etc.
¡Cómo implementar este caso sin usar magia ni hechizos! ¡Leamos más abajo! ¡Vamos!

Introducción
Sin duda, el procesamiento de grandes volúmenes de datos en tiempo real ofrece amplias oportunidades para su uso en sistemas modernos. Una de las combinaciones más populares para esto es el tándem de Apache Kafka y Spark Streaming, donde Kafka crea un flujo de paquetes de mensajes entrantes, y Spark Streaming procesa estos paquetes a través de un intervalo de tiempo determinado.
Para mejorar la resiliencia de la aplicación, usaremos puntos de control — checkpoints. Con este mecanismo, cuando el módulo de Spark Streaming necesite recuperar datos perdidos, solo tendrá que volver al último punto de control y reanudar los cálculos desde allí.
Arquitectura del sistema en desarrollo

Componentes utilizados:
- — es un sistema de mensajería distribuido con publicación y suscripción. Es adecuado tanto para consumo autónomo como en línea de mensajes. Para evitar la pérdida de datos, los mensajes de Kafka se guardan en disco y se replican dentro del clúster. El sistema Kafka se construye sobre el servicio de sincronización ZooKeeper;
- — componente Spark para el procesamiento de datos en tiempo real. El módulo Spark Streaming se construye utilizando una arquitectura de "micro-batch", donde el flujo de datos se interpreta como una secuencia continua de pequeños paquetes de datos. Spark Streaming recibe datos de diferentes fuentes y los agrupa en pequeños paquetes. Se crean nuevos paquetes en intervalos de tiempo regulares. Al inicio de cada intervalo de tiempo se crea un nuevo paquete, y todos los datos que lleguen durante ese intervalo se incluyen en dicho paquete. Al final del intervalo, la acumulación del paquete se detiene. El tamaño del intervalo se determina por un parámetro denominado intervalo de paquetes;
- — combina el procesamiento relacional con la programación funcional de Spark. Por datos estructurados se entienden aquellos que cuentan con un esquema, es decir, un conjunto único de campos para todos los registros. Spark SQL admite la entrada de múltiples fuentes de datos estructurados y, gracias a la información del esquema, puede extraer de manera eficiente solo los campos necesarios de los registros, además de brindar interfaces API de DataFrame;
- — es una base de datos relacional en la nube relativamente económica, un servicio web que facilita la configuración, operación y escalado, administrado directamente por Amazon.
Instalación y ejecución del servidor Kafka
Antes de usar Kafka, es necesario asegurarse de tener Java, ya que se utiliza la JVM:
sudo apt-get update
sudo apt-get install default-jre
java -version
Crearemos un nuevo usuario para trabajar con Kafka:
sudo useradd kafka -m
sudo passwd kafka
sudo adduser kafka sudo
A continuación, descargamos el paquete desde el sitio oficial de Apache Kafka:
wget -P /YOUR_PATH "http://apache-mirror.rbc.ru/pub/apache/kafka/2.2.0/kafka_2.12-2.2.0.tgz"Descomprimimos el archivo descargado:
tar -xvzf /YOUR_PATH/kafka_2.12-2.2.0.tgz
ln -s /YOUR_PATH/kafka_2.12-2.2.0 kafka
El siguiente paso es opcional. Esto se debe a que la configuración predeterminada no permite aprovechar completamente todas las capacidades de Apache Kafka. Por ejemplo, eliminar un tema, categoría o grupo a los cuales se pueden publicar mensajes. Para cambiar esto, editamos el archivo de configuración:
vim ~/kafka/config/server.propertiesAgregue lo siguiente al final del archivo:
delete.topic.enable = trueAntes de iniciar el servidor Kafka, es necesario arrancar el servidor ZooKeeper; utilizaremos un script auxiliar que se proporciona con la distribución de Kafka:
Cd ~/kafka
bin/zookeeper-server-start.sh config/zookeeper.properties
Una vez que ZooKeeper haya iniciado correctamente, en una terminal separada iniciamos el servidor Kafka:
bin/kafka-server-start.sh config/server.propertiesCrearemos un nuevo tema llamado Transaction:
bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic transactionVerificaremos que el tema con el número correcto de particiones y replicación se haya creado:
bin/kafka-topics.sh --describe --zookeeper localhost:2181 
Omitiremos los momentos de prueba del productor y consumidor para el nuevo tema. Más información sobre cómo probar el envío y recepción de mensajes se encuentra en la documentación oficial — . Ahora pasamos a la escritura del productor en Python utilizando la API KafkaProducer.
Escritura del productor
El productor generará datos aleatorios — 100 mensajes cada segundo. Consideraremos que los datos aleatorios son un diccionario compuesto por tres campos:
- Sucursal — el nombre del punto de venta de la institución financiera;
- Moneda — la moneda de la transacción;
- Cantidad — el monto de la transacción. La cantidad será un número positivo si se trata de una compra de divisas por parte del Banco, y negativo si es una venta.
El código para el productor es el siguiente:
from numpy.random import choice, randint
def get_random_value():
new_dict = {}
branch_list = ['Kazan', 'SPB', 'Novosibirsk', 'Surgut']
currency_list = ['RUB', 'USD', 'EUR', 'GBP']
new_dict['branch'] = choice(branch_list)
new_dict['currency'] = choice(currency_list)
new_dict['amount'] = randint(-100, 100)
return new_dict
Luego, utilizando el método send, enviamos el mensaje al servidor, al tema requerido, en formato JSON:
from kafka import KafkaProducer
producer = KafkaProducer(bootstrap_servers=['localhost:9092'],
value_serializer=lambda x:dumps(x).encode('utf-8'),
compression_type='gzip')
my_topic = 'transaction'
data = get_random_value()
try:
future = producer.send(topic = my_topic, value = data)
record_metadata = future.get(timeout=10)
print('--> El mensaje ha sido enviado a un tema:
{}, partición: {}, offset: {}'
.format(record_metadata.topic,
record_metadata.partition,
record_metadata.offset ))
except Exception as e:
print('--> Parece que ocurrió un error: {}'.format(e))
finally:
producer.flush()
Al ejecutar el script, obtenemos en la terminal los siguientes mensajes:

Esto significa que todo funciona como deseábamos — el productor genera y envía mensajes al tema requerido.
El siguiente paso será instalar Spark y procesar este flujo de mensajes.
Instalación de Apache Spark
Apache Spark es una plataforma de computación en clústeres versátil y de alto rendimiento.
En términos de rendimiento, Spark supera las implementaciones populares del modelo MapReduce, a la vez que ofrece soporte para una gama más amplia de tipos de computación, incluidos las consultas interactivas y el procesamiento en tiempo real. La velocidad juega un papel crucial al procesar grandes volúmenes de datos, ya que es la velocidad la que permite trabajar de manera interactiva, sin perder minutos u horas esperando. Una de las principales ventajas de Spark que proporciona una velocidad tan alta es su capacidad para realizar cálculos en memoria.
Este marco está escrito en Scala, por lo que primero debes instalarlo:
sudo apt-get install scalaDescargamos el paquete de Spark desde el sitio web oficial:
wget "http://mirror.linux-ia64.org/apache/spark/spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz"Descomprimimos el archivo:
sudo tar xvf spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz -C /usr/local/sparkAgregamos la ruta a Spark en el archivo bash:
vim ~/ .bashrcAgregamos las siguientes líneas a través del editor:
SPARK_HOME=/usr/local/spark
export PATH=$SPARK_HOME/bin:$PATH
Ejecutamos el siguiente comando después de hacer cambios en el bashrc:
source ~/ .bashrcImplementación de AWS PostgreSQL
Ahora queda implementar la base de datos donde vamos a cargar la información procesada de los flujos. Para ello, utilizaremos el servicio AWS RDS.
Entramos en la consola de AWS -> AWS RDS -> Bases de datos -> Crear base de datos:

Seleccionamos PostgreSQL y hacemos clic en el botón Siguiente:

Dado que este ejemplo se analiza exclusivamente con fines educativos, utilizaremos un servidor gratuito "en su versión mínima" (Free Tier):

A continuación, marcamos la casilla en el bloque Free Tier, y después se nos ofrecerá automáticamente una instancia de clase t2.micro — aunque es débil, es gratuita y se adapta perfectamente a nuestra tarea:

Después siguen cosas muy importantes: el nombre de la instancia de la base de datos, el nombre del usuario maestro y su contraseña. Llamaremos a la instancia: myHabrTest, usuario maestro: habr, contraseña: habr12345 y hacemos clic en el botón Siguiente:

En la siguiente página están los parámetros que responden a la accesibilidad de nuestro servidor de base de datos desde el exterior (Accesibilidad pública) y la disponibilidad de puertos:

Vamos a crear una nueva configuración para el grupo de seguridad de VPC, que permitirá acceder a nuestro servidor de base de datos desde el exterior a través del puerto 5432 (PostgreSQL).
Pasemos en una ventana de navegador separada a la consola de AWS en la sección VPC Dashboard -> Grupos de seguridad -> Crear grupo de seguridad:

Establecemos el nombre para el grupo de seguridad — PostgreSQL, descripción, especificamos a qué VPC debe asociarse este grupo y hacemos clic en el botón Crear:

Llenamos las reglas de entrada para el puerto 5432 del grupo recién creado, como se muestra en la imagen a continuación. No es necesario indicar el puerto manualmente, se puede seleccionar PostgreSQL de la lista desplegable Tipo.
Técnicamente, el valor ::/0 significa que el tráfico de entrada es accesible para el servidor desde todo el mundo, lo cual no es del todo correcto canonícamente, pero para fines del ejemplo nos permitiremos usar este enfoque:

Regresamos a la página del navegador donde tenemos abierto 'Configurar opciones avanzadas' y elegimos en la sección grupos de seguridad de VPC —> Elegir grupos de seguridad de VPC existentes —> PostgreSQL:

Luego, en la sección Opciones de base de datos —> Nombre de la base de datos —> establecemos el nombre — habrDB.
Los demás parámetros, salvo tal vez la desactivación de la copia de seguridad (período de retención de copias de seguridad — 0 días), monitoreo e Insights de rendimiento, podemos dejarlos por defecto. Hacemos clic en el botón Crear base de datos:

Manejador de flujos
La etapa final será el desarrollo de un trabajo de Spark que manejará nuevos datos provenientes de Kafka cada dos segundos y almacenará el resultado en la base de datos.
Como se mencionó anteriormente, los puntos de control (checkpoints) son el mecanismo principal en Spark Streaming, que debe configurarse para garantizar la resiliencia. Usaremos puntos de control y, en caso de que falle el procedimiento, el módulo de Spark Streaming deberá restaurar los datos perdidos regresando al último punto de control y reanudando los cálculos desde ahí.
El punto de control se puede habilitar estableciendo un directorio en un sistema de archivos resiliente y confiable (por ejemplo, HDFS, S3, etc.) donde se guardará la información del punto de control. Esto se hace, por ejemplo:
streamingContext.checkpoint(checkpointDirectory)En nuestro ejemplo, utilizaremos el siguiente enfoque, a saber, que si checkpointDirectory existe, el contexto se recreará a partir de los datos del punto de control. Si el directorio no existe (es decir, se ejecuta por primera vez), se llama a la función functionToCreateContext para crear un nuevo contexto y configurar DStreams:
from pyspark.streaming import StreamingContext
context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext)
Creamos un objeto DirectStream con el objetivo de conectarnos al tópico 'transaction' utilizando el método createDirectStream de la biblioteca KafkaUtils:
from pyspark.streaming.kafka import KafkaUtils
sc = SparkContext(conf=conf)
ssc = StreamingContext(sc, 2)
broker_list = 'localhost:9092'
topic = 'transaction'
directKafkaStream = KafkaUtils.createDirectStream(ssc,
[topic],
{"metadata.broker.list": broker_list})
Analizamos los datos entrantes en formato JSON:
rowRdd = rdd.map(lambda w: Row(branch=w['branch'],
currency=w['currency'],
amount=w['amount']))
testDataFrame = spark.createDataFrame(rowRdd)
testDataFrame.createOrReplaceTempView("treasury_stream")
Usando Spark SQL, realizamos una simple agrupación y mostramos el resultado en la consola:
select
from_unixtime(unix_timestamp()) as curr_time,
t.branch as branch_name,
t.currency as currency_code,
sum(amount) as batch_value
from treasury_stream t
group by
t.branch,
t.currency
Obtenemos el texto de la consulta y la ejecutamos a través de Spark SQL:
sql_query = get_sql_query()
testResultDataFrame = spark.sql(sql_query)
testResultDataFrame.show(n=5)
Y luego guardamos los datos agregados obtenidos en una tabla en AWS RDS. Para guardar los resultados de la agregación en una tabla de base de datos, utilizaremos el método write del objeto DataFrame:
testResultDataFrame.write
.format("jdbc")
.mode("append")
.option("driver", 'org.postgresql.Driver')
.option("url","jdbc:postgresql://myhabrtest.ciny8bykwxeg.us-east-1.rds.amazonaws.com:5432/habrDB")
.option("dbtable", "transaction_flow")
.option("user", "habr")
.option("password", "habr12345")
.save()
Algunas palabras sobre la configuración de la conexión a AWS RDS. El usuario y la contraseña fueron creados en el paso "Despliegue de AWS PostgreSQL". Como URL del servidor de bases de datos, se debe utilizar el Endpoint que se muestra en la sección Connectivity & security:
Para una correcta vinculación entre Spark y Kafka, se debe ejecutar la tarea a través de spark-submit utilizando el artefacto spark-streaming-kafka-0-8_2.11. Además, aplicaremos también el artefacto para interactuar con la base de datos PostgreSQL, los pasaremos a través de —packages.
Para flexibilidad del script, también extraeremos como parámetros de entrada el nombre del servidor de mensajes y el tema desde el cual queremos recibir datos.
Así que ha llegado el momento de iniciar y verificar el funcionamiento del sistema:
spark-submit
--packages org.apache.spark:spark-streaming-kafka-0-8_2.11:2.0.2,
org.postgresql:postgresql:9.4.1207
spark_job.py localhost:9092 transaction
¡Todo salió bien! Como se puede ver en la imagen a continuación, durante el funcionamiento de la aplicación, se muestran nuevos resultados de agregación cada 2 segundos, porque establecimos el intervalo de agrupamiento en 2 segundos al crear el objeto StreamingContext:

Luego, hacemos una sencilla consulta a la base de datos para verificar la existencia de registros en la tabla transaction_flow:

Conclusión
En este artículo se revisó un ejemplo de procesamiento de información en tiempo real utilizando Spark Streaming en combinación con Apache Kafka y PostgreSQL. Con el aumento de los volúmenes de datos procedentes de diversas fuentes, es difícil sobrestimar el valor práctico de Spark Streaming para la creación de aplicaciones en tiempo real.
El código fuente completo lo puedes encontrar en mi repositorio en .
Estoy ansioso por discutir este artículo, espero tus comentarios y también espero una crítica constructiva de todos los lectores interesados.
¡Te deseo éxito!
Ps. Inicialmente se planeaba utilizar una base de datos local de PostgreSQL, pero considerando mi amor por AWS, decidí trasladar la base de datos a la nube. En el próximo artículo sobre este tema, mostraré cómo implementar completamente el sistema descrito anteriormente en AWS utilizando AWS Kinesis y AWS EMR. ¡Estén atentos a las novedades!
Fuente: habr.com

