InżynieriaCzas czytania: 4 min

Potoki danych w czasie rzeczywistym: aktualność, przywracanie i raportowanie

Oleksandr Melnychenko··
Na tej stronie

„Czas rzeczywisty” może oznaczać sekundy dla jednego zespołu i minuty dla innego. Przed wyborem platformy strumieniowej ustal, jakie działanie zależy od danych i jak duże opóźnienie aktualizacji zaczyna wprowadzać w błąd. Status wysyłki, rezerwacja zapasu, alarm maszyny i raport miesięczny nie wymagają tej samej ścieżki dostarczania.

Określ decyzję i dopuszczalne opóźnienie

Dyspozytor może potrzebować informacji, które dostawy prawdopodobnie nie zmieszczą się w oknie czasowym. Przydatna odpowiedź zależy od danych zamówień, aktualizacji przewoźnika, znaczników czasu zdarzeń i reguły opóźnienia. System musi pokazywać czas ostatniej aktualizacji źródła; stara lokalizacja nie powinna wyglądać na bieżącą.

Zapisz źródło każdego zdarzenia, osobę odpowiedzialną, oczekiwaną częstotliwość, format i identyfikator. Uwzględnij zachowanie przy urządzeniu offline, ograniczaniu liczby żądań przez API lub dwukrotnym wysłaniu tego samego zdarzenia przez partnera.

Oddziel odbiór danych od logiki biznesowej

Usługa przyjmowania odbiera dane. Etap przetwarzania je weryfikuje, łączy z właściwym rekordem biznesowym i ustala, czy zmieniają bieżący stan. Przechowywanie i raportowanie udostępniają następnie ten stan użytkownikom. Mogą to być oddzielne usługi lub prostsza aplikacja, zależnie od wolumenu, aktualności i potrzeb przywracania.

Broker lub trwała kolejka mogą pomagać, gdy producenci i odbiorcy pracują z różną szybkością. Nie są automatycznym wymaganiem każdej integracji. Wybierz je po oszacowaniu szczytowego wolumenu, potrzeb ponownego przetwarzania i kosztu utrzymania.

Oczekuj spóźnionych i powtórzonych zdarzeń

Sieci, urządzenia i systemy zewnętrzne ulegają awariom. Określ obsługę zduplikowanych komunikatów, zdarzeń przychodzących poza kolejnością i korekt wcześniej przyjętych danych. Oddziel identyfikator zdarzenia od identyfikatora rekordu biznesowego: jedna wysyłka może mieć wiele prawidłowych aktualizacji.

Operacja idempotentna nie powoduje dodatkowego skutku biznesowego przy ponowieniu tej samej operacji. Rozpoznanie duplikatu wymaga trwałego klucza operacji; zapobieganie drugiemu zapisowi wymaga też, aby kontrola duplikatu i zapis nie mogły wejść w wyścig z innym procesem roboczym.

Nie traktuj ustawienia dostarczania kolejki jako obietnicy, że cały proces biznesowy przetworzy każde zdarzenie dokładnie raz. Aktualizacje bazy, wywołania zewnętrznych API i powiadomienia mogą potrzebować własnych reguł deduplikacji lub uzgadniania. Dokumentacja semantyki dostarczania Kafka firmy Confluent wyjaśnia, dlaczego ponowienia mogą powodować wielokrotne dostarczenie i gdzie obowiązują silniejsze gwarancje.

Pokazuj użytkownikowi aktualność i jakość danych

Rozróżniaj czas zdarzenia, czyli moment wystąpienia według źródła, od czasu przetwarzania, gdy etap potoku obsługuje zdarzenie. Spóźnione potwierdzenie dostawy może należeć do wczorajszych sum operacyjnych, choć otrzymano je dzisiaj. Omówienie czasu zdarzeń i spóźnionych danych w Apache Beam wyjaśnia tę różnicę.

Widok operacyjny powinien pokazywać czas obserwacji źródła i ostatniego udanego odświeżenia. Odświeżenie strony nie czyni starej lokalizacji bieżącą. Milczące źródło nie musi być uszkodzone: określ oczekiwany interwał aktualizacji lub oddzielny sygnał stanu. Uzgodnij, co użytkownicy widzą, gdy dane są opóźnione, brakujące lub sprzeczne.

Dla raportowania określ, czy wskaźnik używa czasu zdarzenia czy przetwarzania, co się dzieje, gdy spóźnione zdarzenie zmienia przeszły okres, i jak sumy są sprawdzane ze źródłem. Jeśli zachowujesz obecne BI, określ częstotliwość wymiany i kontrolę zgodności.

Ustal cele, które można przetestować

Mierz opóźnienie od znacznika czasu zdarzenia źródłowego do udostępnienia aktualizacji w widoku operacyjnym. Dla przykładowego zdarzenia zaobserwowanego o 10:00:00, otrzymanego o 10:00:02 i widocznego o 10:00:15 opóźnienie wynosi 15 sekund; 13 sekund przypada po odbiorze. Obliczenie zakłada porównywalne zegary. Rozbieżność zegarów lub niewiarygodny znacznik czasu źródła muszą być zapisane jako ograniczenie pomiaru.

Raportuj medianę i wysoki percentyl, np. p95, dla określonego okresu i obciążenia. p95 równe 15 sekund oznacza, że około 95% zmierzonych aktualizacji zakończyło się w tym czasie; nie mówi nic o zdarzeniach, które nigdy nie dotarły, jeśli tych błędów nie liczysz oddzielnie. Wytyczne Google SRE dotyczące monitoringu wyjaśniają, dlaczego sama średnia może ukrywać wolne żądania.

Sprawdzaj kompletność niezależnie: uzgodnij unikalne oczekiwane zdarzenia źródłowe z przyjętymi, odrzuconymi i nadal oczekującymi dla tego samego momentu odcięcia. Jeśli źródło nie może podać oczekiwanej liczby lub sekwencji, opisz zakres możliwy do zweryfikowania zamiast raportować nieuzasadniony procent kompletności. Testuj zwykłe obciążenie, nagły wzrost, awarię źródła i przywracanie; zachowuj identyfikatory i znaczniki czasu potrzebne do odtworzenia rozbieżności.

Zaplanuj monitoring samego potoku: awarie źródeł, rosnące opóźnienie odbiorców, nieprawidłowe rekordy, powtarzające się ponowienia i raporty tracące zgodność. Wytyczne monitoringu Kinesis firmy AWS ilustrują sygnały strumienia i odbiorców wymagające obserwacji; wybór platformy pozostaje zależny od projektu.

Zaprojektuj pierwsze wydanie wokół jednego przepływu

Połącz jeden strumień przewoźnika z rekordami wysyłek, pokaż dyspozytorom opóźnione lub brakujące aktualizacje i uzgadniaj statusy dzienne. Określ interfejsy, oczekiwany wolumen zdarzeń, czas przechowywania, reguły alertów, procedurę przywracania i zespół reagujący na awarie wymiany danych.

Jeśli zespół przenosi dane między starym a nowym systemem, oddziel też migrację rekordów historycznych od synchronizacji bieżących zmian. Obie potrzebują odpowiedzialności i weryfikacji. Nasz przewodnik modernizacji wyjaśnia przejście, a przykład raportowania ERP pokazuje definiowanie raportu na podstawie rekordów źródłowych.

Aby omówić potok, przygotuj jedną decyzję operacyjną, jej źródła danych i okno czasowe, w którym odpowiedź pozostaje użyteczna. Na tej podstawie możemy określić zakres pierwszego połączonego procesu.