Curs 11
Procesarea fluxurilor de
date, Apache Kafka
Prof. Univ. Dr. Diaconita Vlad
Publish–subscribe (Pub-sub)
• Este un pattern in cadrul căruia:
• Furnizorii de mesaje nu trimit mesajele direct unor destinatari ci le publica pe
diverse categorii (topicuri)
• Consumatorii de mesaje se abonează la unul sau mai multe topicuri, primind
mesajele pe măsură ce ele sunt publicate
• Astfel sunt decuplate sistemele care trebuie sa comunice
Temperaturi Constanta
29, 29.5, 29.4, 30
Temperaturi Constanta
Furnizori Consumatori
Cantitati vandute produs A
7, 5, 3
Cantitati vandute produs A
BROKER
Conturi noi
ionel@[Link],maria@[Link]
Tipuri Pub/Sub
• Bazate pe trimitere automata de mesaje (push based – message
queing) – mesajele sunt trimise automat către consumatorii abonați
• Redis, RabitMQ: bazate pe cozi de mesaje (FIFO); pot exista cozi de mesaje
pentru un fiecare consumator sau cozi de mesaje comune; cozile de mesaje
pot avea prioritati diferite
Tipuri Pub/Sub
• Bazate pe extragere de mesaje (pull based – event streaming) –
consumatorul cere mesajele atunci cand le poate procesa
• Apache Kafka -> bazat pe loguri de evenimente; un mesaj poate fi consumat o
data sau de mai multe ori, acesta este păstrat pana când expira conform
politicii de retentie;
Apache Kafka
• O solutie distribuita – poate rula pe un cluster alcatuit din mai multe
noduri
• Mesajele sunt replicate pe mai multe noduri. Un mesaj odata scris, nu
poate fi modificat, poate fi insa ulterior sters sau compactat
• Poate fi vazuta ca o solutie hibrida baza de date distribuita+sistem de
mesagerie
• Util pentru solutii tip ETL (extract, transform, load), CDC (change data
capture), Big Data, agregare de metrici de la diverse locatii
Fluxuri de date
• Kafka propune o abordare orientata pe evenimente
Aplicatie Asigurarea
Mobila securitatii
Platforma de Analiza in
Aplicatie Web
Streaming timp real
API Hadoop
Un produs a fost vizualizat
O comanda noua a fost introdusa
Arhitectura Kafka
• Un broker are rol de nod
Aplicatie 101
Topic 1
Consumator controler
• Un topic este impartit in mai
Broker 1 Broker 2
multe partiții
Topic 1
Aplicatie 102 Topic 1, Partitia 0
Topic 1, Partitia 1
Topic 1, Partitia 0
Topic 1, Partitia 1
Consumator
• Fiecare partiție are un nod
Topic 2, Partitia 0 Topic 2, Partitia 0
lider
Kafka Connect
• Un mesaj este salvat
Topic 1
OSS, Oracle ATP
(commited) când este scris
Aplicatie 103 Cluster Kafka pe nodul lider si replicat pe
nodurile sincronizate* (daca
AWS S3
Broker 3 Broker 4 acks=all)
Sink Connect
Topic 1, Partitia 0
Topic 1, Partitia 0 • Fiecare partiție este
Dispozitiv de Topic 2
Topic 1, Partitia 1
Topic 1, Partitia 1
BD NoSQL
ordonata si fiecare mesaj din
masura
Topic 2, Partitia 0
Topic 3, Partitia 0
partitie primeste un id
incremental (offset)
Partiții
• Partiția este unitatea de baza a paralelismului in Kafka
• Un topic = un log partiționat
• La nivel de producător si broker, se pot scrie pe partiții diferite
in paralel
• La nivel de consumator, Kafka aloca o singura partitie unui fir
de execuție
• Cu cat sunt mai multe partiții, viteza de procesare creste
ID_M_P1 0 1 2 3 4 5 6 7 8 9
ID_M_P2 0 1 2 3 4 5 6 7 Topic 1
ID_M_P3 0 1 2 3 4 5 6 7 8
Partiții
• Ordinea este garantata doar la nivelul unei partitii
• Datele sunt păstrate pentru un timp limitat (e.g., 2 săptămâni)
• Datele odată adăugate într-o partiție, nu pot fi modificate (datele sunt
imuabile)
• Datele sunt împărțite aleatoriu pe partiții, cu excepția situației in care
se furnizează o cheie
Consumarea datelor
• Solutiile clasice de procesarea mesajelor (messaging processing)
presupun in general procesări simple aplicate mesajelor, cel mai
adesea la nivel de mesaj individual
• Pentru a aplica procesari mai complexe (agregari, jonctiuni) se pot
folosi solutii de data streaming: Kafka Streams, Spark Streams
Consumarea datelor
Algoritmi online - Varianța
• Aggr := 0; Medie:=0;
• k :=0
• welford(Stream):
• for citire in Stream:
• k := k+1
• medieVeche := Medie
• Medie := Medie + (citire - Medie)/k
• Aggr := Aggr + (citire -Medie)* (citire - medieVeche)
• return Aggr /(n-1)
for citire in Stream:
k := k+1
medieVeche := Medie
Welford - exemplu Medie := Medie + (citire - Medie)/k
Aggr := Aggr + (citire -Medie)* (citire - medieVeche)
return Aggr /(n-1)
Citiri (dB) k Medie Aggr Varianta Walford Varianta Excel I1 I2 Alarma
0.00000 0.00000
27.00 1 27.00000 0.00000
27.00 2 27.00000 0.00000 0.00000 0.00000 27.00000 27.00000 NU
27.00 3 27.00000 0.00000 0.00000 0.00000 27.00000 27.00000 NU
27.00 4 27.00000 0.00000 0.00000 0.00000 27.00000 27.00000 NU
27.40 5 27.08000 0.12800 0.03200 0.03200 26.89217 27.26783 DA
27.00 6 27.06667 0.13333 0.02667 0.02667 26.89520 27.23813 NU
27.00 7 27.05714 0.13714 0.02286 0.02286 26.89840 27.21589 NU
27.40 8 27.10000 0.24000 0.03429 0.03429 26.90558 27.29442 DA
27.40 9 27.13333 0.32000 0.04000 0.04000 26.92333 27.34333 DA
27.40 10 27.16000 0.38400 0.04267 0.04267 26.94311 27.37689 DA
27.40 11 27.18182 0.43636 0.04364 0.04364 26.96248 27.40116 NU
27.40 12 27.20000 0.48000 0.04364 0.04364 26.98066 27.41934 NU
27.30 13 27.20769 0.48923 0.04077 0.04077 26.99568 27.41970 NU
29.00 14 27.33571 3.47214 0.26709 0.26709 26.79307 27.87836 DA
27.40 15 27.34000 3.47600 0.24829 0.24829 26.81680 27.86320 NU
Prag = Medie ± 105% * Deviatia Standard
Numarat cuvinte
• Toate direcțiile de sănătate publică județene și a municipiului București vor
informa, în cursul zilei de astăzi, Inspectoratele Școlare Județene/al Municipiului
București, respectiv senatele universitare ale instituțiilor de învățământ superior și
Comitetele Județene / al Municipiului București pentru Situații de Urgență
(CJSU/CMBSU) cu privire la situația epidemiologică la nivelul fiecărei localități, iar
până la data de 10 septembrie, în funcție de aceste date, de particularitățile
locale, de infrastructura și resursele umane ale fiecărei unități de învățământ în
parte, consiliul de administrație al unității de învățământ va propune
Inspectoratului Școlar Județean/al Municipiului București aplicarea unuia dintre
scenariile de organizare și desfășurare a cursurilor în unitatea de învățământ,
conform Ordinului ministrului educației și cercetării și al ministrului sănătății
5487/1494/ 01.09.2020.
• Comitetele judeţene/al municipiului Bucureşti pentru situaţii de urgenţă
(CJSU/CMBSU) vor emite hotărârea privind scenariul de funcţionare pentru
fiecare unitate de învăţământ pentru începutul anului şcolar 2020-2021.[…]
Exemplu Kafka + Spark Streams
public class Application { Dataset<String> words = messagesDf
public static void main(String[] args) { .as([Link]())
SparkSession spark = [Link]()
flatMap((FlatMapFunction<String, String>) x -
.appName("Spark SQL Dataframe API") > [Link]([Link](" ")).iterator(), Encoders.
.getOrCreate(); STRING());
Dataset<Row> wordsDf = [Link]();
Dataset<Row> stopWordsDf = [Link]
Dataset<Row> messagesDf = [Link]() aset([Link]([Link]), Enco
.format("kafka")
[Link]()).toDF();
.option("[Link]", "server_ip:port") wordsDf = [Link](stopWordsDf, wordsDf
.col("value").equalTo([Link]("value")
.option("subscribe", "test") ), "leftanti");
.load() wordsDf = [Link]("value").count();
.selectExpr("CAST(value AS STRING)");
[Link](desc("count")).show();
Exemplu Kafka + Spark Streams
Asigurarea securității
• Cu setările implicite, orice utilizator/aplicație poate scrie/citi pe/de pe
orice topic
• Securitea comunicațiilor dintre furnizori, consumatori si Kafka se
poate asigura folosind Transport Layer Security (TLS)
Asigurarea securitatii
• Autentificarea folosind TLS/SSL (2-way authentication) - furnizorii si
consumatorii se pot autentifica pe clusterul Kafka care le verifica
identitatea pe baza certificatelor SSL
• Clientul solicita o resursa protejata folosind protocolul HTTPS
• Serverul returnează certificatul sau public
• Clientul verifica certificatul serverului folosind Certificate authority (CA)
• Daca certificatul serverului este corect, clientul isi trimite propriul certificat
• Serverul verifica certificatul clientului folosind CA
• Daca este validat, comunicatiile viitoare vor fi cripate folosind cheile
schimbate in timpul autentificarii
Asigurarea securitatii - 2-way authentication
Asigurarea securitatii
• Autorizare folosind liste de control (ACL) - dupa ce un client se
autentifica, acesta poate primi drepturi de acces la diverse topicuri:
• bin/[Link] --authorizer [Link] --
authorizer-properties [Link]=localhost:2181 --add --allow-
principal User:Vlad --operation Read --topic comenzi_noi
• bin/[Link] --authorizer [Link] --
authorizer-properties [Link]=localhost:2181 --add --allow-
principal User:Vlad --operation Write --topic comenzi_procesate
Apache Zookeeper
• Folosit împreună cu Kafka are rol in:
• Gestiunea brokerilor din cluster
• Alegerea nodului controller
• Configurarea topicurilor, inclusiv a numărului de partiții, locația acestora
• Menținerea listelor de control pentru acces
Va mulțumesc pentru atenție!