Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Ciao, Habr! Oggi costruiremo un sistema che utilizzerà Spark Streaming per elaborare i flussi di messaggi di Apache Kafka e registrare il risultato dell'elaborazione in un database cloud AWS RDS.

Immaginiamo che una certa istituzione creditizia ci affidi il compito di elaborare le transazioni in entrata «al volo» per tutte le sue filiali. Questo potrebbe essere fatto per calcolare in modo operativo la posizione aperta in valuta per il tesoro, limiti o risultati finanziari delle transazioni, ecc.

Come realizzare questo caso senza utilizzare magie e incantesimi — leggiamo qui sotto! Andiamo!

Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming
(Fonte dell'immagine)

Introduzione

Certamente, elaborare un grande insieme di dati in tempo reale offre ampie possibilità di utilizzo nei sistemi moderni. Una delle combinazioni più popolari per questo è il tandem Apache Kafka e Spark Streaming, dove Kafka crea un flusso di pacchetti di messaggi in ingresso, e Spark Streaming elabora questi pacchetti in un intervallo di tempo specificato.

Per aumentare la resilienza dell'applicazione, utilizzeremo i checkpoint. Attraverso questo meccanismo, quando il modulo Spark Streaming avrà bisogno di ripristinare i dati persi, dovrà semplicemente tornare all'ultimo checkpoint e riprendere i calcoli da lì.

Architettura del sistema in fase di sviluppo

Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Componenti utilizzati:

  • Apache Kafka — è un sistema di messaggistica distribuito con pubblicazione e sottoscrizione. Adatto sia per il consumo autonomo che per quello online dei messaggi. Per evitare la perdita di dati, i messaggi Kafka vengono salvati su disco e replicati all'interno del cluster. Il sistema Kafka è costruito sopra il 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 accetta dati da diverse sorgenti e li combina in piccoli pacchetti. Nuovi pacchetti vengono creati a intervalli regolari. All'inizio di ogni intervallo temporale viene creato un nuovo pacchetto e tutti i dati ricevuti durante questo intervallo sono inclusi nel pacchetto. Alla fine dell'intervallo, l'aumento del pacchetto si arresta. La dimensione dell'intervallo è determinata da un parametro chiamato intervallo di batch;
  • 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 tutte le registrazioni. Spark SQL supporta l'immissione da molteplici sorgenti di dati strutturati e, grazie alla presenza di informazioni sullo schema, può estrarre in modo efficiente solo i campi necessari dalle registrazioni, oltre a fornire interfacce API DataFrame;
  • AWS RDS — è un database relazionale cloud relativamente economico, un servizio web che semplifica la configurazione, il funzionamento e la scalabilità, amministrato direttamente da Amazon.

Installazione e avvio del server Kafka

Prima di utilizzare Kafka, è necessario assicurarsi di avere Java, poiché per il funzionamento viene utilizzata la 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

Successivamente, 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 passaggio successivo è facoltativo. Infatti, le impostazioni predefinite non consentono di sfruttare appieno tutte le funzionalità di Apache Kafka. Ad esempio, non è possibile eliminare un argomento, una categoria o un gruppo a cui possono essere pubblicati i messaggi. Per modificare questo, dobbiamo modificare il file di configurazione:

vim ~/kafka/config/server.properties

Aggiungi in fondo al file quanto segue:

delete.topic.enable = true

Prima di avviare il server Kafka, è necessario avviare il server ZooKeeper, utilizzeremo 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 topic chiamato Transaction:

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

Assicuriamoci che il topic con il numero desiderato di partizioni e replica sia stato creato:

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

Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Tralasciamo le fasi di test del produttore e del consumatore per il nuovo topic. Maggiori informazioni su come testare l'invio e la ricezione di messaggi sono disponibili nella documentazione ufficiale — Invia alcuni messaggi. Ma noi passiamo alla scrittura di un produttore in Python utilizzando l'API KafkaProducer.

Scrivere un produttore

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

  • Branch — nome del punto vendita dell'istituto di credito;
  • Valuta — valuta della transazione;
  • Importo — ammontare della transazione. 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 è il seguente:

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 di cui abbiamo bisogno, in 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 otteniamo i seguenti messaggi nel terminale:

Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Questo significa che tutto funziona come desiderato — il produttore genera e invia messaggi al topic di cui abbiamo bisogno.
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 universale e ad alte prestazioni.

In termini di prestazioni, Spark supera le implementazioni più diffuse del modello MapReduce, offrendo allo stesso tempo supporto per una gamma più ampia di tipi di calcolo, compresi query interattive e trattamento in tempo reale. La velocità è cruciale nel trattamento di grandi volumi di dati, poiché è proprio la velocità che consente di lavorare in modalità interattiva, senza aspettare minuti o ore. Uno dei vantaggi principali di Spark, che consente questa alta velocità, è la capacità di eseguire calcoli in memoria.

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

sudo apt-get install scala

Scarichiamo la distribuzione 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"

Decomprimiamo 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 nel 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 qui sotto dopo aver apportato le modifiche a bashrc:

source ~/bashrc

Distribuzione di AWS PostgreSQL

Resta da distribuire un database in cui caricheremo le informazioni elaborate dai flussi. Per questo utilizzeremo il servizio AWS RDS.

Accediamo alla console AWS —> AWS RDS —> Databases —> Crea database:
Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Selezioniamo PostgreSQL e facciamo clic sul pulsante Next:
Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Poiché questo esempio è trattato esclusivamente a scopi educativi, utilizzeremo un server gratuito 'minimale' (Free Tier):
Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Successivamente, spuntiamo la casella nel blocco Free Tier, e dopo di ciò ci verrà automaticamente proposto un'istanza di classe t2.micro — sebbene sia piuttosto modesta, è gratuita e va benissimo per il nostro scopo:
Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Poi ci sono cose molto importanti: il nome dell'istanza del DB, il nome dell'utente master e la sua password. Chiamiamo l'istanza: myHabrTest, utente master: habr, password: habr12345 e facciamo clic sul pulsante Next:
Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Nella pagina successiva si trovano le impostazioni che determinano l'accessibilità del nostro server DB dall'esterno (Accessibilità pubblica) e le impostazioni delle porte:

Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Creiamo una nuova configurazione per il gruppo di sicurezza VPC, che consentirà l'accesso esterno al nostro server DB tramite la porta 5432 (PostgreSQL).
Andiamo in una finestra di browser separata alla console AWS nella sezione VPC Dashboard —> Gruppi di sicurezza —> Crea gruppo di sicurezza:
Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Impostiamo il nome del gruppo di sicurezza - PostgreSQL, forniamo una descrizione, specifichiamo a quale VPC questo gruppo deve essere associato e facciamo clic sul pulsante Crea:
Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

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

Tecnicamente, il valore ::/0 indica la disponibilità del traffico in ingresso per il server da tutto il mondo, il che non è propriamente corretto, ma per questa spiegazione ci permettiamo di adottare questo approccio:
Apache Kafka e l'elaborazione dei dati in streaming 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 dati in streaming con Spark Streaming

Successivamente, nella sezione Opzioni database -> Nome database -> impostiamo il nome - habrDB.

Possiamo lasciare gli altri parametri, a meno che non disattiviamo il backup (periodo di mantenimento del backup - 0 giorni), il monitoraggio e Performance Insights, sui valori predefiniti. Facciamo clic sul pulsante Crea database:
Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Gestore dei flussi

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

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

Il checkpoint può essere attivato impostando una directory in un file system resiliente e sicuro (ad esempio, HDFS, S3, ecc.) dove verranno salvate le informazioni sui checkpoint. Questo avviene, ad esempio, tramite:

streamingContext.checkpoint(checkpointDirectory)

Nel nostro esempio utilizzeremo il seguente approccio, ovvero, se checkpointDirectory esiste, il contesto verrà ricreato dai dati del checkpoint. Se la directory non esiste (cioè 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 con l'obiettivo di connetterci al topic "transaction" utilizzando il metodo createDirectStream della libreria KafkaUtils:

da 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, eseguiamo un semplice raggruppamento e stampiamo 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

Ottenere il testo della query e eseguirlo attraverso Spark SQL:

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

E poi salviamo i dati aggregati ottenuti in una tabella in AWS RDS. Per salvare i risultati dell'aggregazione in una tabella del 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()

Qualche parola sulla configurazione della connessione a AWS RDS. L'utente e la password li abbiamo creati nel passaggio "Distribuzione AWS PostgreSQL". Come URL del server del database, dovremmo utilizzare l'Endpoint, che viene visualizzato nella sezione Connectivity & security:

Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Per una corretta integrazione tra Spark e Kafka, è necessario eseguire il job tramite spark-submit utilizzando l'artefatto spark-streaming-kafka-0-8_2.11. Inoltre, utilizzeremo anche l'artefatto per interagire con il database PostgreSQL, che verrà passato tramite —packages.

Per una maggiore flessibilità dello script, includiamo anche come parametri di ingresso il nome del server di messaggistica e il topic da cui vogliamo ricevere i dati.

Quindi, è giunto il momento di avviare e controllare 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

Ce l'abbiamo fatta! Come mostrato nell'immagine qui sotto — durante l'esecuzione dell'applicazione, i nuovi risultati dell'aggregazione vengono visualizzati ogni 2 secondi, poiché abbiamo impostato l'intervallo di pacchettizzazione a 2 secondi quando abbiamo creato l'oggetto StreamingContext:

Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Successivamente, eseguiamo una semplice query al database per verificare la presenza di registrazioni nella tabella transaction_flow:

Apache Kafka e l'elaborazione dei dati in streaming con Spark Streaming

Conclusione

In questo articolo è stato presentato un esempio di elaborazione continua delle informazioni utilizzando Spark Streaming insieme ad 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 lo puoi trovare nel mio repository su GitHub.

Sono felice di discutere di questo articolo, attendo i tuoi commenti e spero anche in critiche costruttive da parte di tutti i lettori interessati.

Ti auguro successo!

Ps. Inizialmente si prevedeva di utilizzare un database locale PostgreSQL, ma considerando il mio amore per AWS, ho deciso di trasferire il database nel cloud. Nel prossimo articolo su questo tema mostrerò come realizzare l'intero sistema sopra descritto in AWS utilizzando AWS Kinesis e AWS EMR. Rimanete sintonizzati!

Fonte: habr.com

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