Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Ciao, Habr! Oggi costruiremo un sistema che utilizza Spark Streaming per elaborare flussi di messaggi Apache Kafka e registrare il risultato nell'archivio dati cloud di AWS RDS.

Immaginiamo che un'istituzione creditizia ci chieda di elaborare le transazioni in ingresso "al volo" in tutte le sue filiali. Questo può essere fatto per calcolare rapidamente la posizione aperta in valuta per la tesoreria, i limiti o il risultato finanziario delle transazioni, ecc.

Come realizzare questo caso senza ricorrere a magie e incantesimi? Scopriamolo qui sotto! Iniziamo!

Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming
(Fonte immagine)

Introduzione

Senza dubbio, l'elaborazione di grandi volumi di dati in tempo reale offre ampia flessibilità per applicazioni nei sistemi moderni. Una delle combinazioni più popolari per questo è il tandem Apache Kafka e Spark Streaming, dove Kafka genera un flusso di pacchetti di messaggi in ingresso, mentre Spark Streaming elabora questi pacchetti a intervalli di tempo specificati.

Per migliorare l'affidabilità dell'applicazione, utilizzeremo i checkpoint (checkpoints). Con questo meccanismo, quando il modulo Spark Streaming dovrà recuperare dati persi, dovrà solo tornare all'ultimo checkpoint e riprendere i calcoli da lì.

Architettura del sistema in fase di sviluppo

Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Componenti utilizzati:

  • Apache Kafka è un sistema distribuito di messaggistica con pubblicazione e sottoscrizione. È adatto sia per il consumo autonomo che per quello online dei messaggi. Per prevenire la perdita di dati, i messaggi di Kafka vengono salvati su disco e replicati all'interno del cluster. Il sistema Kafka è costruito sulla base del servizio di sincronizzazione ZooKeeper;
  • Apache Spark Streaming — componente Spark per l'elaborazione di dati in streaming. Il modulo Spark Streaming è costruito utilizzando un'architettura a "micro-batch", in cui il flusso di dati è interpretato come una sequenza continua di piccoli pacchetti di dati. Spark Streaming riceve dati da diverse fonti e li aggrega in piccoli pacchetti. Nuovi pacchetti vengono creati a intervalli regolari. All'inizio di ogni intervallo di tempo viene creato un nuovo pacchetto, e tutti i dati ricevuti durante questo intervallo vengono inclusi nel pacchetto. Alla fine dell'intervallo, l'accumulo del pacchetto si arresta. La dimensione dell'intervallo è definita da un parametro noto come intervallo di pacchettizzazione;
  • Apache Spark SQL — combina l'elaborazione relazionale con la programmazione funzionale di Spark. I dati strutturati si riferiscono a dati che hanno uno schema, ossia un insieme uniforme di campi per tutti i record. Spark SQL supporta l'ingresso da molteplici fonti di dati strutturati e, grazie alla presenza di informazioni sullo schema, può estrarre in modo efficiente solo i campi necessari dei record e fornisce interfacce API DataFrame;
  • AWS RDS — è un database relazionale cloud comparativamente economico, un servizio web che semplifica la configurazione, il funzionamento e la scalabilità, gestito direttamente da Amazon.

Installazione e avvio del server Kafka

Prima di utilizzare Kafka, è necessario assicurarsi di avere Java, poiché per il funzionamento è necessario il JVM:

sudo apt-get update 
sudo apt-get install default-jre
java -version

Creiamo un nuovo utente per lavorare con Kafka:

sudo useradd kafka -m
sudo passwd kafka
sudo adduser kafka sudo

Dopo, scarichiamo il pacchetto dal sito ufficiale di Apache Kafka:

wget -P /YOUR_PATH "http://apache-mirror.rbc.ru/pub/apache/kafka/2.2.0/kafka_2.12-2.2.0.tgz"

Estraiamo l'archivio scaricato:

tar -xvzf /YOUR_PATH/kafka_2.12-2.2.0.tgz
ln -s /YOUR_PATH/kafka_2.12-2.2.0 kafka

Il passo successivo è opzionale. Infatti, le impostazioni predefinite non consentono di utilizzare appieno tutte le funzionalità di Apache Kafka. Ad esempio, per eliminare un argomento, una categoria o un gruppo sui quali possono essere pubblicati messaggi. Per modificare ciò, dobbiamo modificare il file di configurazione:

vim ~/kafka/config/server.properties

Aggiungi alla fine del file quanto segue:

delete.topic.enable = true

Prima di avviare il server Kafka, è necessario avviare il server ZooKeeper utilizzando uno script di supporto fornito con la distribuzione di Kafka:

Cd ~/kafka
bin/zookeeper-server-start.sh config/zookeeper.properties

Dopo che ZooKeeper è stato avviato con successo, in un terminale separato avviamo il server Kafka:

bin/kafka-server-start.sh config/server.properties

Creiamo un nuovo argomento chiamato Transaction:

bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 3 --topic transaction

Verifichiamo che l'argomento sia stato creato con il numero corretto di partizioni e replica:

bin/kafka-topics.sh --describe --zookeeper localhost:2181

Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Non tratteremo i dettagli del test del produttore e del consumatore per il nuovo argomento. Maggiori informazioni su come testare l'invio e la ricezione dei messaggi possono essere trovate nella documentazione ufficiale — Invia alcuni messaggi. Passiamo ora alla scrittura del produttore in Python utilizzando l'API KafkaProducer.

Scrittura del produttore

Il produttore genererà dati casuali — 100 messaggi ogni secondo. Per dati casuali intendiamo un dizionario composto da tre campi:

  • Filiale — nome del punto vendita dell'istituto di credito;
  • Currency — valuta della transazione;
  • Amount — importo dell'affare. L'importo sarà un numero positivo se si tratta di un acquisto di valuta da parte della Banca, e negativo se si tratta di una vendita.

Il codice per il produttore appare come segue:

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

Successivamente, utilizzando il metodo send, inviamo un messaggio al server, nel topic a noi necessario, nel 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('--> Il messaggio è stato inviato a un topic: 
            {}, partizione: {}, offset: {}' 
            .format(record_metadata.topic,
                record_metadata.partition,
                record_metadata.offset ))   
                             
except Exception as e:
    print('--> Sembra che si sia verificato un errore: {}'.format(e))

finally:
    producer.flush()

All'avvio dello script, riceviamo i seguenti messaggi nel terminale:

Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Questo significa che tutto funziona come desiderato: il produttore genera e invia messaggi nel topic corretto.
Il passo successivo sarà l'installazione di Spark e l'elaborazione di questo flusso di messaggi.

Installazione di Apache Spark

Apache Spark è una piattaforma di calcolo distribuito versatile e ad alte prestazioni.

In termini di prestazioni, Spark supera le implementazioni più popolari del modello MapReduce, offrendo al contempo supporto per una gamma più ampia di tipi di calcolo, inclusi query interattive ed elaborazione in tempo reale. La velocità gioca un ruolo fondamentale nell'elaborazione di grandi volumi di dati, poiché permette di lavorare in modo interattivo senza dover attendere minuti o ore. Uno dei maggiori punti di forza di Spark, che consente di raggiungere tali velocità elevate, è la capacità di eseguire calcoli in memoria.

Questo framework è scritto in Scala, quindi è necessario installarla per prima cosa:

sudo apt-get install scala

Scarichiamo il pacchetto di Spark dal sito ufficiale:

wget "http://mirror.linux-ia64.org/apache/spark/spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz"

Estraiamo l'archivio:

sudo tar xvf spark-2.4.2/spark-2.4.2-bin-hadoop2.7.tgz -C /usr/local/spark

Aggiungiamo il percorso di Spark al file bash:

vim ~/ .bashrc

Aggiungiamo tramite l'editor le seguenti righe:

SPARK_HOME=/usr/local/spark
export PATH=$SPARK_HOME/bin:$PATH

Eseguiamo il comando sottostante dopo aver modificato il bashrc:

source ~/.bashrc

Distribuzione AWS PostgreSQL

Ora dobbiamo distribuire il database dove caricheremo le informazioni elaborate dai flussi. A tal fine utilizzeremo il servizio AWS RDS.

Accediamo alla console AWS → AWS RDS → Database → Crea database:
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Selezioniamo PostgreSQL e facciamo clic sul pulsante Avanti:
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Poiché questo esempio è trattato esclusivamente a scopo educativo, utilizzeremo un server gratuito "minimalista" (Free Tier):
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Successivamente, spuntiamo la casella nel blocco Free Tier, e dopo di ciò ci verrà automaticamente proposto un'istanza di tipo t2.micro — anche se limitata, è gratuita e si adatta perfettamente alle nostre esigenze:
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Le seguenti sono informazioni molto importanti: nome dell'istanza DB, nome dell'utente master e la sua password. Denominiamo l'istanza: myHabrTest, utente master: habr, password: habr12345 e facciamo clic sul pulsante Avanti:
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Nella pagina successiva troviamo i parametri relativi alla disponibilità del nostro server DB dall'esterno (Accessibilità pubblica) e alla disponibilità delle porte:

Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Creiamo una nuova configurazione per il gruppo di sicurezza VPC, che permetterà l'accesso al nostro server DB dall'esterno attraverso la porta 5432 (PostgreSQL).
Apriamo in una nuova finestra del browser la console AWS nella sezione VPC Dashboard —> Gruppi di Sicurezza —> Crea gruppo di sicurezza:
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Assegniamo un nome al Gruppo di Sicurezza — PostgreSQL, una descrizione, indichiamo a quale VPC questo gruppo deve essere associato e clicchiamo sul pulsante Crea:
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Compiliamo per il nuovo gruppo le regole Inbound per la porta 5432, come mostrato nell'immagine sottostante. Non è necessario specificare manualmente la porta, ma possiamo selezionare PostgreSQL dal menu a discesa Tipo.

A rigor di termini, il valore ::/0 indica l'accessibilità del traffico in ingresso per il server da tutto il mondo, cosa che non è propriamente corretta, ma per la spiegazione dell'esempio ci permettiamo di utilizzare tale approccio:
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Torniamo alla pagina del browser, dove abbiamo aperto "Configura impostazioni avanzate" e selezioniamo nella sezione Gruppi di sicurezza VPC —> Scegli gruppi di sicurezza VPC esistenti —> PostgreSQL:
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Successivamente, nella sezione Opzioni database —> Nome database —> assegniamo il nome — habrDB.

Possiamo lasciare gli altri parametri, a parte disattivare le copie di backup (periodo di conservazione del backup — 0 giorni), il monitoraggio e Performance Insights, impostati su default. Clicchiamo sul pulsante Crea database:
Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Processore di flussi

L'ultimo passo sarà lo sviluppo di un lavoro Spark che elaborerà nuovi dati ricevuti da Kafka ogni due secondi e inserirà il risultato nel database.

Come accennato in precedenza, i checkpoint sono il meccanismo principale in SparkStreaming che deve essere impostato per garantire la resilienza. Utilizzeremo i checkpoint e, in caso di guasto della procedura, il modulo Spark Streaming per il ripristino dei dati persi dovrà semplicemente tornare all'ultimo checkpoint e riprendere i calcoli da lì.

Un checkpoint può essere attivato impostando una directory in un sistema di file affidabile e resistente (ad esempio, HDFS, S3, ecc.) in cui verranno salvate le informazioni del checkpoint. Questo viene fatto ad esempio con:

streamingContext.checkpoint(checkpointDirectory)

Nel nostro esempio utilizzeremo il seguente approccio, ovvero, se checkpointDirectory esiste, il contesto sarà ricreato dai dati del checkpoint. Se la directory non esiste (ossia, viene eseguita per la prima volta), verrà chiamata la funzione functionToCreateContext per creare un nuovo contesto e configurare i DStreams:

from pyspark.streaming import StreamingContext

context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext)

Creiamo un oggetto DirectStream per connetterci al topic «transaction» utilizzando il metodo createDirectStream della libreria 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})

Analizziamo i dati in ingresso in 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")

Utilizzando Spark SQL, facciamo una semplice aggregazione e visualizziamo il risultato nella console:

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

Otteniamo il testo della query e lo eseguiamo tramite Spark SQL:

sql_query = get_sql_query()
testResultDataFrame = spark.sql(sql_query)
testResultDataFrame.show(n=5)

E poi salviamo i dati aggregati in una tabella in AWS RDS. Per salvare i risultati dell'aggregazione in una tabella di database, utilizzeremo il metodo write dell'oggetto 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()

Alcuni dettagli sulla configurazione della connessione a AWS RDS. L'utente e la password sono stati creati nel passaggio "Distribuzione AWS PostgreSQL". Come URL del server del database, bisogna utilizzare l'Endpoint visualizzabile nella sezione Connettività e sicurezza:

Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Per garantire un corretto collegamento tra Spark e Kafka, è necessario eseguire il job tramite spark-submit utilizzando l'artefatto spark-streaming-kafka-0-8_2.11. Inoltre, utilizzeremo anche un artefatto per l'interazione con il database PostgreSQL, che sarà passato tramite —packages.

Per rendere lo script più flessibile, estrarremo anche il nome del server dei messaggi e il topic dal quale vogliamo ricevere i dati come parametri di input.

Quindi, è arrivato il momento di avviare e verificare il funzionamento 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

Tutto ha funzionato! Come si può vedere nell'immagine qui sotto, durante l'esecuzione dell'applicazione i nuovi risultati di aggregazione vengono visualizzati ogni 2 secondi, poiché abbiamo impostato l'intervallo di elaborazione su 2 secondi quando abbiamo creato l'oggetto StreamingContext:

Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Successivamente, facciamo una semplice richiesta al database per verificare la presenza di record nella tabella transaction_flow:

Apache Kafka e l'elaborazione dei flussi di dati con Spark Streaming

Conclusione

In questo articolo è stato presentato un esempio di elaborazione dei dati in tempo reale utilizzando Spark Streaming in combinazione con Apache Kafka e PostgreSQL. Con l'aumento del volume dei dati provenienti da diverse fonti, è difficile sovrastimare il valore pratico di Spark Streaming per la creazione di applicazioni in streaming e applicazioni che operano in tempo reale.

Il codice sorgente completo è disponibile nel mio repository su GitHub.

Sono felice di discutere questo articolo, attendo i vostri commenti e spero in una critica costruttiva da parte di tutti i lettori interessati.

Vi auguro successo!

Ps. Inizialmente era previsto utilizzare un database locale PostgreSQL, ma considerando il mio amore per AWS, ho deciso di spostare il database nel cloud. Nella prossima articolo su questo argomento mostrerò come implementare l'intero sistema descritto sopra in AWS utilizzando AWS Kinesis e AWS EMR. Rimanete sintonizzati!

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, server VPS VDS 🔥 Acquista hosting affidabile per siti web con protezione DDoS, server VPS VDS | ProHoster