Construirea unui pipeline de analiză a fraudei în timp real
Prezentarea arhitecturii: ingestia a 50K de evenimente/secundă, îmbogățirea cu smart signals și scorarea riscului în sub 10ms cu motorul nostru de streaming.
Procesarea a 50.000 de evenimente de fingerprint pe secundă, îmbogățirea fiecăruia cu smart signals și returnarea unui scor de risc în sub 10 milisecunde necesită o arhitectură de streaming atent proiectată. Acest articol parcurge pipeline-ul nostru, de la ingestie până la decizie.
Stratul de ingestie
Evenimentele sosesc sub formă de cereri HTTPS POST de la agentul nostru JavaScript care rulează în browserele vizitatorilor. Fiecare eveniment conține payload-ul de semnale criptat — de regulă 8-12KB de date comprimate care acoperă peste 130 de semnale de browser. Serverele noastre edge fac terminarea TLS, validează semnătura cererii și transmit payload-ul către pipeline-ul de procesare.
Folosim o implementare multi-regiune, în care serverele edge sunt colocate cu nodurile CDN ale clienților noștri. Astfel, dus-întorsul de rețea rămâne sub 20ms pentru 95% din cereri la nivel global. Serverele edge sunt servicii Go fără stare, care rulează în spatele unui load balancer și scalează orizontal în funcție de volumul de cereri.
Extracția semnalelor
Prima etapă de procesare decriptează și parsează payload-ul de semnale. Fiecare semnal este extras, validat și tipizat. Hash-urile de canvas sunt verificate față de valori imposibile cunoscute (care indică blocarea sau spoofing-ul canvas-ului). Parametrii WebGL sunt validați încrucișat pentru consistență. Proprietățile navigatorului sunt verificate față de combinații valide cunoscute.
Această etapă realizează și normalizarea semnalelor. Șirurile de user agent sunt parsate în componente structurate (browser, versiune, OS, dispozitiv). Dimensiunile ecranului sunt normalizate pentru a ține cont de scalarea DPI. Offset-urile de fus orar sunt validate față de datele de geolocație a IP-ului.
Îmbogățirea cu Smart Signals
Semnalele extrase sunt apoi îmbogățite cu analiza Smart Signals — stratul nostru de inteligență server-side. Aceasta include detectarea incognito (compararea tiparelor de semnale cu semnături cunoscute de navigare privată), detectarea VPN (corelarea datelor IP cu semnalele de fus orar și locale), detectarea tampering-ului de browser (identificarea inconsistențelor care indică spoofing de semnale) și detectarea mașinilor virtuale (recunoașterea profilurilor hardware asociate cu VMware, VirtualBox și VM-uri din cloud).
Fiecare smart signal este calculat independent și produce atât un rezultat boolean, cât și un scor de încredere. Etapa de îmbogățire adaugă 24 de semnale suplimentare fiecărui eveniment, oferind o evaluare cuprinzătoare a amenințărilor care depășește ceea ce poate obține colectarea din partea clientului de una singură.
Motorul de scorare a riscului
Evenimentul îmbogățit este transmis motorului nostru de scorare a riscului — un model de arbore de decizie gradient-boosted, antrenat pe milioane de evenimente etichetate. Modelul ia în calcul toate cele 130+ de semnale brute, 24 de smart signals și câteva caracteristici derivate: metrici de viteză (câte evenimente au venit de la acest dispozitiv în ultimele 5 minute, 1 oră și 24 de ore), tipare de comportament istoric și scoruri de reputație a rețelei.
Modelul produce un scor de risc între 0 și 100, împreună cu principalii factori contributori. Un scor de 85, de exemplu, ar putea fi însoțit de factori precum „VPN detectat”, „mod incognito” și „viteză ridicată — 47 de evenimente în 5 minute”. Această explicabilitate este esențială pentru analiștii de fraudă, care trebuie să înțeleagă de ce a fost semnalat un anumit eveniment.
Stratul de stocare și interogare
Toate evenimentele sunt persistate în ClickHouse — o bază de date coloanară optimizată pentru interogări analitice pe seturi mari de date. ClickHouse gestionează volumul nostru de scriere (50K evenimente/secundă) fără nicio dificultate, iar stocarea coloanară permite interogări analitice sub o secundă pe miliarde de rânduri.
Folosim o strategie de retenție pe mai multe niveluri. Datele „fierbinți” (ultimele 7 zile) sunt stocate pe SSD-uri NVMe pentru un timp de răspuns la interogări sub 100ms. Datele „calde” (7-90 de zile) sunt pe SSD-uri standard. Datele „reci” (90+ de zile) sunt comprimate și mutate în object storage, interogabile, dar cu latență mai mare.
Kafka ca coloană vertebrală
Apache Kafka leagă etapele pipeline-ului între ele. Fiecare etapă citește din și scrie în topicuri Kafka. Stratul de ingestie scrie evenimentele brute. Etapa de extracție a semnalelor citește evenimentele brute și scrie evenimentele extrase. Etapa de îmbogățire cu Smart Signals citește evenimentele extrase și scrie evenimentele îmbogățite. Motorul de scorare a riscului citește evenimentele îmbogățite și scrie evenimentele scorate.
Această arhitectură oferă mai multe avantaje: etapele pot fi scalate independent, defecțiunile dintr-o etapă nu le afectează pe celelalte și putem rerula (replay) evenimentele prin orice etapă pentru depanare sau reprocesare. Grupurile de consumatori Kafka permit procesarea paralelă în cadrul fiecărei etape, iar semantica exactly-once garantează că niciun eveniment nu este procesat de două ori sau pierdut.
Bugetul de latență
Ținta noastră de latență end-to-end este de 10ms, din momentul în care payload-ul de semnale îmbogățit sosește la pipeline-ul de procesare până în momentul în care scorul de risc este returnat. Iată cum se împarte bugetul: extracția semnalelor durează 1-2ms, îmbogățirea cu Smart Signals durează 3-4ms, scorarea riscului durează 2-3ms, iar serializarea și răspunsul durează 1-2ms. Saltul Kafka între etape adaugă sub 1ms în implementarea noastră colocată.
Respectarea consecventă a acestui buget la 50K evenimente/secundă necesită optimizare atentă în fiecare etapă. Folosim pool-uri de memorie prealocate, serializare zero-copy și scrieri în ClickHouse pe loturi. Modelul de scorare a riscului este compilat în cod nativ folosind ONNX Runtime, eliminând overhead-ul interpretorului Python.
Mark a petrecut două săptămâni analizând performanța pipeline-ului înainte de a găsi gâtuirea în stratul nostru de căutare distribuită — un singur mutex serializa căutările pe toate goroutine-urile. După trecerea la un design de blocare sharded, p99 a scăzut de la 48ms la 9ms. Uneori corecția este jenant de simplă odată ce o găsești.