Kompleksowe przypadki użycia
Przykład zastosowania: Przetwarzanie danych CDR i roamingowych telekomu
Ta strona śledzi krok po kroku jedną konkretną osobę, gdy buduje potok danych, który zamienia surowe, zagadkowe pliki produkowane przez sieć telefoniczną w czyste dane, których systemy billingowe, zespoły ds. oszustw i analitycy mogą faktycznie używać. Jest napisana dla czytelników bez doświadczenia w telekomie — każdy element żargonu jest rozpakowywany w momencie, gdy się pojawia.

Poznaj osobę i problem
Budowniczym w tej historii jest Anna Kowalska, Inżynier Danych w Polkomtelu. Wspiera ją Tomasz Wiśniewski, Opiekun Danych, który czuwa nad jakością danych i ochroną danych osobowych, oraz Katarzyna Zielińska, Administrator Platformy, która przygotowuje połączenia.
Zadanie Anny na dziś: wziąć pliki, które sieć mobilna Polkomtela wypluwa co kilka minut, i uczynić je użytecznymi.
By zrozumieć, dlaczego jest to trudne, najpierw musisz zrozumieć, czym te pliki są.
Czym jest CDR?
Za każdym razem, gdy telefon coś robi w sieci — wykonuje połączenie, wysyła SMS, używa danych mobilnych — element sprzętu sieciowego zapisuje maleńki rekord opisujący to. Tym rekordem jest CDR: Call Detail Record (rekord szczegółów połączenia).
Pojedynczy CDR odpowiada na pytania takie jak:
- Kto to zrobił? (identyfikatory abonenta)
- Co zrobił? (połączenie głosowe, SMS, sesja danych)
- Kiedy i jak długo? (czas rozpoczęcia, czas trwania lub zużyte megabajty)
- Gdzie? (która wieża komórkowa, która sieć)
CDR-y są surowcem firmy telekomunikacyjnej. Bez nich nikogo nie można by rozliczyć, żadne oszustwo nie mogłoby zostać wykryte, a żaden menedżer nie mógłby zobaczyć, jak sieć jest używana.
Dlaczego CDR-y są trudne do odczytania
Tu jest haczyk. CDR-y nie są czytelnymi dla człowieka plikami tekstowymi. Są plikami binarnymi zapisanymi w kompaktowym formacie technicznym zwanym ASN.1 (zakodowanym regułami zwanymi BER lub DER). Pomyśl o ASN.1 jak o ciasno zapakowanej walizce: niezwykle oszczędnej miejscowo — co ma znaczenie, gdy sieć produkuje miliony rekordów na godzinę — ale nie możesz jej po prostu otworzyć i odczytać. Potrzebujesz narzędzia, które zna dokładny schemat pakowania, by ją zdekodować.
Schematy pakowania podążają za międzynarodowymi standardami od organu zwanego 3GPP (grupa definiująca technologię sieci mobilnych). DataFlow AI dostarcza szablony dla standardowych typów CDR zdefiniowanych w 3GPP TS 32.298.
Mówiąc prościej
CDR to paragon sieci za jedną akcję telefonu. ASN.1/BER to skrót, w którym napisany jest paragon — genialny do oszczędzania miejsca, bezużyteczny dla człowieka, dopóki coś go nie zdekoduje. Zadaniem DataFlow AI w tym przykładzie zastosowania jest bycie dekoderem, a następnie oczyszczenie i skierowanie tego, co wychodzi.
A co z roamingiem?
Gdy abonent Polkomtela podróżuje za granicą, używa sieci innego kraju. Ta zagraniczna sieć wykonuje pracę, więc Polkomtel musi jej zapłacić — a obaj operatorzy muszą wymienić rekordy, by się rozliczyć. Podobnie, gdy zagraniczny gość używa sieci Polkomtela w Polsce, Polkomtel rozlicza jego operatora.
Ta wymiana używa trzech specjalnych formatów plików, wszystkich również zakodowanych w ASN.1:
| Format | Pełna nazwa | Cel prostym językiem |
|---|---|---|
| TAP3 | Transferred Account Procedure v3 | Ogólnoświatowy standardowy „plik faktury" zużycia roamingowego, wymieniany między operatorami partiami (codziennie lub co tydzień). DataFlow AI obsługuje TAP3.12. |
| NRTRDE | Near Real-Time Roaming Data Exchange | Szybsza, lżejsza wersja TAP3 wysyłana w ciągu godzin, specjalnie po to, by oszustwo można było szybko wyłapać, zanim narośnie ogromny rachunek. |
| RAP | Returned Accounts Procedure | „List reklamacyjny". Jeśli operator otrzyma plik TAP3 z błędami, odsyła plik RAP odrzucający błędne rekordy lub cały plik. |
Mówiąc prościej
Wyobraź sobie roaming jako dwie restauracje, które od czasu do czasu obsługują swoich stałych klientów nawzajem. TAP3 to miesięczny szczegółowy rachunek, który jedna restauracja wysyła drugiej. NRTRDE to szybkie ostrzeżenie tego samego dnia — „twój klient właśnie zamówił dużo, możesz chcieć to sprawdzić" — i tak oszustwo zostaje wcześnie wyłapane. RAP to odpowiedź: „trzy pozycje na twoim rachunku są błędne, za te nie płacimy".
Potrzeba biznesowa
Polkomtel potrzebuje zdekodowanych danych CDR i roamingowych do czterech zadań naraz:
- Taryfikacja — ustalanie ceny każdego połączenia, SMS-a i sesji danych.
- Billing — zamiana otaryfikowanego zużycia w faktury klientów.
- Wykrywanie oszustw — wychwytywanie podejrzanych wzorców (zwłaszcza w roamingu), zanim będą kosztować pieniądze.
- Rozliczenia roamingowe — uzgadnianie tego, co Polkomtel jest winien innym operatorom i co oni są winni Polkomtelowi.
Anna zbuduje jeden potok danych, który dekoduje surowe pliki, czyści je i wzbogaca, uruchamia kontrole jakości ukierunkowane na oszustwa i ładuje wynik, by wszystkie cztery zespoły mogły z niego korzystać.
Kształt potoku danych
Raw CDR files DataFlow AI Design Studio Targets
+---------------+ D +-----------------------------------+ L +-----------+
| gs://.../*.cdr|------>| Normalise -> Enrich -> |------>| Billing |
| (ASN.1/BER) | | Fraud quality checks -> Route | | Analytics |
+---------------+ +-----------------------------------+ | Sub-360 |
Decode Transform Load
- Decode — odczytuje binarne pliki
.cdri zamienia każdy zapakowany rekord w zwykłe kolumny. - Transform — porządkuje identyfikatory, dodaje przydatny kontekst (wzbogacanie).
- Quality — uruchamia kontrole dostrojone do ujawniania sygnałów oszustwa i błędnych danych.
- Load — zapisuje wynik do systemu billingowego, do hurtowni analitycznej i do widoku Subscriber-360.
Krok 1 — Połączenie: gdzie mieszkają surowe pliki
Elementy sieciowe Polkomtela zrzucają swoje pliki CDR do Google Cloud Storage (GCS) — folderu w chmurze — według codziennej konwencji takiej jak gs://polkomtel-cdr-raw/{date}/.
DataFlow AI odczytuje je za pomocą specjalnie zbudowanego konektora zwanego cdr-asn1 („Telecom CDR (ASN.1)"). Nie jest to ogólny czytnik plików — to sam dekoder ASN.1, zbudowany specjalnie dla Polkomtela. Kilka faktów o nim, o których Anna pamięta:
- Jest tylko do odczytu. Dekoduje przychodzące pliki; nigdy nie zapisuje plików CDR z powrotem.
- Akceptuje pliki kończące się na
.cdr,.asn1,.datlub.bin. - Potrafi odczytać cały folder naraz, używając wzorca plików (domyślny wzorzec to
*.cdr). - Potrafi odczytywać z GCS, z Amazon S3 lub z lokalnego dysku.
- Ustala, który przełącznik sieciowy wyprodukował plik, na podstawie nazwy pliku — sprzęt Huawei, Ericsson lub Nokia — więc możliwe są statystyki na poszczególnych dostawców.
Administrator Platformy, Katarzyna, konfiguruje połączenie cdr-asn1 raz, wskazując je na zasobnik GCS. Anna po prostu je wybiera.
Uważaj
Konektor CDR jest tylko do odczytu z założenia. Zdekodowany wynik jest konwertowany na uporządkowany, kolumnowy format zwany Parquet i partycjonowany według daty i godziny. Nie ładujesz danych z powrotem przez ten konektor — jego jedynym kierunkiem jest „dekoduj na wejściu".
Krok 2 — Utwórz potok danych
Anna otwiera Studio Projektowe, wizualny kreator potoków danych typu przeciągnij i upuść:
- Lewy panel boczny → Potoki danych → „+ Nowy Potok danych".
- Nazwa:
cdr_decode_and_enrich. - Opis: „Dekoduj pliki CDR sieci, wzbogać, sprawdź pod kątem oszustw, załaduj do billingu i analityki."
- Tryb: wybiera Batch. Pliki CDR docierają falami, a dekodowanie fali według harmonogramu to zadanie wsadowe. (DataFlow AI potrafi również przetwarzać CDR-y jako ciągły strumień do pracy w czasie rzeczywistym; ten przewodnik buduje wersję wsadową, a notka dalej wyjaśnia opcję strumieniową.)
- Utwórz — otwiera się płótno.
Płótno ma znajome cztery obszary: Paletę Komponentów typów węzłów po lewej, Płótno pośrodku, Inspektor Właściwości po prawej oraz Panel Dolny (logi i podglądy) na dole. Studio automatycznie zapisuje co 30 sekund.
Krok 3 — Decode: odczytaj i rozpakuj pliki CDR
Anna przeciąga węzeł Źródło na płótno i wybiera konektor Telecom CDR (ASN.1). Klika węzeł i konfiguruje go w Inspektorze Właściwości.
Kluczowe ustawienia dla źródła CDR:
| Ustawienie | Wartość Anny | Co oznacza |
|---|---|---|
| Połączenie | połączenie GCS cdr-asn1 | Gdzie mieszkają surowe pliki |
filePattern | *.cdr | Odczytaj każdy plik .cdr w folderze |
templateType | auto | Pozwól dekoderowi samemu wykryć typ każdego rekordu |
templateVersion | 3GPP-R15 | Którego standardu układ pól użyć |
batchSize | 10000 | Dekoduj rekordy partiami po 10 000 |
skipMalformed | true | Jeśli jeden rekord jest uszkodzony, pomiń go i kontynuuj |
Ustawienie templateType: auto warte jest zatrzymania się nad nim. Pojedynczy folder plików CDR zawiera wiele rodzajów rekordów zmieszanych razem — połączenia głosowe, SMS-y, sesje danych, zdarzenia roamingowe. Dekoder bada każdy rekord, odczytuje jego wbudowany kod typu i dopasowuje go do właściwego szablonu — mapy mówiącej „pole 1 to IMSI abonenta, pole 5 to czas trwania" i tak dalej.
DataFlow AI zna dwanaście natywnych typów rekordów CDR. Te, które potok danych Anny zobaczy głównie:
| Typ rekordu | Czym jest |
|---|---|
CS_VOICE_MO / CS_VOICE_MT | Normalne połączenie głosowe — wykonane (MO) lub odebrane (MT) |
CS_SMS_MO / CS_SMS_MT | Wiadomość tekstowa — wysłana lub odebrana |
PS_DATA | Sesja danych mobilnych — przenosi APN oraz objętości wysyłania/pobierania |
IMS_VOICE | Połączenie VoLTE (głos przenoszony przez sieć danych 4G) |
ROAMING_IN / ROAMING_OUT | Zużycie przez gościa w sieci Polkomtela lub przez abonenta Polkomtela za granicą |
Gdy ten węzeł działa, każdy zapakowany rekord binarny staje się wierszem z właściwymi, nazwanymi kolumnami: identyfikatory abonenta, numery dzwoniący i wywoływany, czas rozpoczęcia, czas trwania lub liczby bajtów, lokalizacja komórki, wskaźnik roamingu i powód zakończenia połączenia.
Mówiąc prościej
Dekoder to uniwersalny tłumacz stojący przy drzwiach. Pliki wchodzą, mówiąc trzema różnymi dialektami „binarnego skrótu telekomowego"; wiersze wychodzą, mówiąc prostym, oznaczonym językiem. Od tego węzła dalej reszta potoku danych to zwykła, czytelna praca na danych.
Notka o dziwnych polach identyfikatorów
Dane telekomowe mają własne specjalne typy pól, a dekoder obsługuje je automatycznie. Trzy, o których usłyszysz:
- IMSI — International Mobile Subscriber Identity. Unikatowy numer identyfikujący kartę SIM (a więc abonenta) w dowolnej sieci na całym świecie.
- MSISDN — faktycznie sam numer telefonu (
numer telefonupo polsku). - IMEI — numer identyfikujący fizyczny aparat, niezależnie od tego, która karta SIM jest w nim włożona.
Te są często zapakowane w jeszcze bardziej skompresowany sposób (schemat zwany TBCD). Dekoder je rozpakowuje; późniejsza transformacja porządkuje ich formatowanie.
Krok 4 — Transform: normalizacja i wzbogacanie
Zdekodowane CDR-y są użyteczne, ale wciąż surowe. Anna dodaje dwa węzły transformacji.
4a. Normalise — uczyń numery telefonów spójnymi
Numery telefonów w surowych CDR-ach pojawiają się w wielu kształtach: lokalny prefiks 0, kod międzynarodowy, zapakowane cyfry. Systemy billingowe i ds. oszustw potrzebują ich w jednym spójnym kształcie — międzynarodowym formacie E.164, który dla Polski wygląda jak +48, a po nim dziewięć cyfr.
Anna przeciąga węzeł Expression (węzeł, który oblicza lub przeformatowuje kolumny) i stosuje wbudowaną funkcję platformy normalize_msisdn() do numerów dzwoniącego i wywoływanego. Po tym węźle każdy numer telefonu w zbiorze danych wygląda identycznie pod względem struktury, niezależnie od tego, w jakim kształcie dotarł.
Używa również extract_date() i extract_time(), by podzielić znacznik czasu rekordu na czyste kolumny daty i czasu — przydatne do partycjonowania i raportowania później.
4b. Enrich — dodaj kontekst, którego brakuje surowemu rekordowi
Surowy CDR zna IMSI abonenta, ale nie jego plan taryfowy; zna identyfikator komórki, ale nie miasto, które ta komórka obejmuje. Wzbogacanie oznacza łączenie danych CDR z tabelami referencyjnymi, by każdy rekord niósł kontekst, którego potrzebują zespoły niższego szczebla.
Anna dodaje węzły Join (Join wyrównuje dwa zbiory danych na wspólnej kolumnie), by wprowadzić:
- Taryfę i segment abonenta, połączone z tabeli referencyjnej klientów po IMSI — by taryfikacja mogła zastosować właściwe ceny.
- Region lub miasto lokalizacji komórki, połączone z tabeli referencyjnej sieci — by analitycy mogli zobaczyć zużycie geograficznie.
Po wzbogaceniu każdy wiersz jest samodzielnym, w pełni opisanym zdarzeniem: kto, co, kiedy, gdzie, na jakim planie, w jakim regionie.
Mówiąc prościej
Surowy CDR jest jak zdjęcie bez podpisu. Wzbogacanie pisze podpis — imię osoby, miejsce, datę słowami — wyszukując każdy szczegół na liście referencyjnej i dołączając go. Zespoły ds. oszustw i billingu potrzebują podpisanego zdjęcia, a nie gołego obrazu.
Krok 5 — Kontrole jakości dostrojone do oszustw i poprawności
Tutaj potok danych Anny zarabia na swoje utrzymanie. Zanim jakiekolwiek zdekodowane dane zostaną dopuszczone dalej, węzeł Jakości uruchamia automatyczne kontrole. Niektóre wyłapują zwykłe błędy danych; niektóre są celowo wymierzone we wskaźniki oszustwa.
Przeciąga węzeł Jakości za wzbogacaniem i konfiguruje reguły z dziesięciu typów reguł DataFlow AI:
| Kontrola | Typ reguły | Dlaczego ma znaczenie |
|---|---|---|
| IMSI abonenta jest zawsze obecne | NOT_NULL | Rekordu bez abonenta nie da się rozliczyć |
| Czas trwania połączenia mieści się w rozsądnych granicach | RANGE | 30-godzinne „połączenie" to niemal na pewno usterka lub oszustwo |
| Numer telefonu pasuje do wzorca E.164 | REGEX | Wyłapuje numery, których normalizator nie mógł naprawić |
| Brak zduplikowanych identyfikatorów rekordów | UNIQUE | Duplikaty rozliczyłyby klienta dwukrotnie |
| Objętość roamingu nie jest absurdalnie wysoka | STATISTICAL | Nagły ogromny skok za granicą to klasyczny sygnał oszustwa |
| Pliki faktycznie dotarły w tej godzinie | FRESHNESS | Cicha luka oznacza, że kanał sieci się zaciął |
Każda reguła otrzymuje dotkliwość (Critical, Warning, Info) i ustawienie, czy niepowodzenie blokuje potok danych. Anna ustawia kontrole NOT_NULL i UNIQUE na Critical / blokuj — dane billingowe nigdy nie mogą być błędne — a statystyczną kontrolę skoku roamingu jako Warning, więc podnosi alert dla zespołu ds. oszustw bez zatrzymywania całego ładowania.
Aspekt oszustwa jest tutaj prawdziwy i konkretny. Oszustwo roamingowe działa przez szybkie naliczanie ogromnych opłat za granicą, zanim operator macierzysty zauważy. To właśnie dlatego istnieje format NRTRDE — dostarcza dane roamingowe niemal w czasie rzeczywistym. Statystyczne i zakresowe kontrole Anny na rekordach roamingowych to automatyczne pułapki, które zamieniają te szybkie dane w szybki alert.
Uważaj
Reguły jakości działają przy każdym wykonaniu. Kanały CDR są wysokowolumenowe i bez nadzoru — wadliwa partia z jednego przełącznika lub zdublowany kanał może w ciągu minut wlać miliony błędnych wierszy w kierunku systemu billingowego. Węzeł jakości to brama, która powstrzymuje takie zdarzenia. Tomasz, Opiekun Danych, również obserwuje wskaźnik jakości domeny CDR z Centrum Zarządzania, a wykrywanie anomalii flaguje nietypowe wyniki automatycznie.
Dane osobowe są maskowane również tutaj
CDR-y są pełne danych osobowych — numery telefonów, identyfikatory abonentów, lokalizacje. Skaner PII DataFlow AI klasyfikuje je automatycznie (MSISDN, IMSI/IMEI i lokalizacja są rozpoznawanymi kategoriami). Tam, gdzie odbiorca niższego szczebla nie powinien widzieć surowych identyfikatorów, Anna dodaje węzeł Expression, by je zamaskować, a DataFlow AI stosuje maskowanie oparte na rolach w momencie wyświetlania. Polskie prawo telekomunikacyjne wymaga, by CDR-y były przechowywane przez siedem lat; reguły retencji platformy to wymuszają. Tomasz przegląda to wszystko z Centrum Zarządzania.
Krok 6 — Route i Load: zasilenie czterech odbiorców
Zdekodowane, wzbogacone, sprawdzone pod kątem jakości dane muszą teraz dotrzeć do kilku zespołów naraz. Anna dodaje węzeł Router — węzeł, który wysyła wiersze różnymi ścieżkami w zależności od warunku — oraz trzy węzły Sink.
- Billing — rekordy głosowe, SMS i danych płyną do systemu billingowego/taryfikacji, by klientów można było zafakturować. Anna zapisuje je do sinka bazy danych w trybie Upsert (aktualizuj istniejące rekordy, wstawiaj nowe).
- Analityka — pełna kopia trafia do hurtowni analitycznej (BigQuery, Snowflake lub Teradata) do raportowania o zużyciu sieci i przychodach. Ten sink używa trybu Append, ponieważ historia analityczna tylko rośnie.
- Rozliczenia roamingowe — rekordy
ROAMING_INiROAMING_OUTsą kierowane do dedykowanej tabeli używanej do uzgadniania z przychodzącymi plikami TAP3 i do generowania odrzuceń RAP tam, gdzie plik TAP3 zagranicznego operatora zawiera błędy.
Klika Waliduj, a następnie Uruchom → Uruchom Teraz dla testu. Węzły rozświetlają się, Konsola streamuje liczby zdekodowanych rekordów, a kilka sekund później uruchomienie kończy się na zielono. Anna sprawdza tabelę billingową i widzi czyste, znormalizowane, wzbogacone wiersze. Potok danych działa.
Krok 7 — Aspekt Subscriber-360
Jedną z najcenniejszych rzeczy, które ten potok danych odblokowuje, jest Subscriber-360 — pojedynczy, kompletny widok wszystkiego, co robi jeden abonent.
Same z siebie CDR-y są rozproszone: rekordy głosowe w jednym miejscu, rekordy danych w innym, zdarzenia roamingowe w trzecim, wszystkie kluczowane zagadkowymi identyfikatorami. Gdy potok danych Anny już je zdekodował, znormalizował i wzbogacił — tak by każdy rekord niósł czysty numer telefonu, taryfę, region i datę — wszystkie mogą zostać powiązane z tym samym abonentem.
Załadowane do hurtowni analitycznej, daje to biznesowi 360-stopniowy obraz na klienta: jego nawyki dzwonienia, jego apetyt na dane, jego zachowanie roamingowe, jego typowe lokalizacje. Ten pojedynczy widok zasila:
- Przewidywanie odejść — wychwytywanie klientów skłonnych do odejścia.
- Spersonalizowane oferty — dopasowywanie taryf do rzeczywistego zużycia.
- Profilowanie oszustw — rozpoznawanie, kiedy zachowanie abonenta nagle wygląda inaczej niż jego własna historia.
Nic z tego nie jest możliwe, gdy dane są zamknięte wewnątrz binarnych plików ASN.1. Potok dekodowania i wzbogacania to klucz, który to odblokowuje.
Mówiąc prościej
Subscriber-360 to różnica między pudełkiem po butach pełnym nieposortowanych paragonów a uporządkowaną księgą. Paragony (surowe CDR-y) zawierają wszystko — ale dopiero gdy są zdekodowane, opatrzone datą i oznaczone, możesz otworzyć stronę jednej osoby i zobaczyć jej całą historię na pierwszy rzut oka.
Krok 8 — Potok danych jako YAML
Wszystko, co Anna zbudowała przez przeciąganie węzłów, jest przechowywane jako pojedynczy czytelny dla człowieka plik YAML, automatycznie zatwierdzany do Git przy każdym zapisie. Może przeglądać i edytować go bezpośrednio przez zakładkę YAML; wizualne płótno i YAML pozostają zsynchronizowane. Oto gotowy potok danych:
apiVersion: dataflow.polkomtel.com/v1
kind: Pipeline
metadata:
name: cdr-decode-and-enrich
namespace: network-operations
labels:
domain: cdr
purpose: rating-billing-fraud
annotations:
description: Decode network CDR files, enrich, fraud-check, load to billing and analytics
owner: anna.kowalska@plk.pl
sla: "hourly"
spec:
schedule: "15 * * * *"
timezone: Europe/Warsaw
enabled: true
timeout: 3600
retries: 3
retryDelay: 300
parameters:
- name: cdr_date
type: date
default: "{{today}}"
description: The date-partitioned CDR folder to process
nodes:
- id: src_cdr
type: connector_source
label: Decode CDR files (ASN.1/BER)
config:
connector: cdr-asn1
filePattern: "*.cdr"
templateType: auto
templateVersion: "3GPP-R15"
batchSize: 10000
skipMalformed: true
- id: src_subscriber
type: connector_source
label: Subscriber reference
config:
connector: teradata
table: REFERENCE.SUBSCRIBER_DIM
- id: src_network
type: connector_source
label: Cell-to-region reference
config:
connector: teradata
table: REFERENCE.CELL_LOCATION_DIM
- id: normalise
type: expression
label: Normalise numbers and timestamps
config:
expressions:
- "calling_number = normalize_msisdn(calling_number)"
- "called_number = normalize_msisdn(called_number)"
- "event_date = extract_date(record_timestamp)"
- "event_time = extract_time(record_timestamp)"
- id: enrich_subscriber
type: joiner
label: Add tariff and segment
config:
leftInput: normalise
rightInput: src_subscriber
joinType: LEFT
on: "served_imsi = imsi"
- id: enrich_location
type: joiner
label: Add region and city
config:
leftInput: enrich_subscriber
rightInput: src_network
joinType: LEFT
on: "cell_id"
- id: quality_gate
type: quality
label: Fraud and correctness checks
config:
rules:
- { type: NOT_NULL, column: served_imsi, severity: CRITICAL, blockPipeline: true }
- { type: UNIQUE, column: record_id, severity: CRITICAL, blockPipeline: true }
- { type: RANGE, column: duration_seconds, min: 0, max: 86400, severity: HIGH, blockPipeline: false }
- { type: REGEX, column: calling_number, pattern: "^\\+48[0-9]{9}$", severity: MEDIUM, blockPipeline: false }
- { type: ANOMALY, column: roaming_volume_mb, metric: mean, sigmaThreshold: 3, lookbackDays: 30, severity: HIGH, blockPipeline: false }
- { type: FRESHNESS, timestampColumn: record_timestamp, maxAgeHours: 2, severity: MEDIUM, blockPipeline: false }
- id: split_by_type
type: router
label: Route by record type
config:
routes:
- { condition: "record_type IN ('ROAMING_IN','ROAMING_OUT')", target: load_roaming }
- { condition: "true", target: load_billing }
- id: load_billing
type: connector_sink
label: Billing / rating system
config:
connector: oracle
table: BILLING.RATED_USAGE
writeMode: UPSERT
upsertKeys: [record_id]
batchSize: 10000
- id: load_analytics
type: connector_sink
label: Analytics warehouse (Subscriber-360)
config:
connector: bigquery
table: analytics.cdr.subscriber_events
writeMode: APPEND
batchSize: 20000
- id: load_roaming
type: connector_sink
label: Roaming settlement table
config:
connector: teradata
table: ROAMING.SETTLEMENT_EVENTS
writeMode: UPSERT
upsertKeys: [record_id]
edges:
- { from: src_cdr, to: normalise }
- { from: normalise, to: enrich_subscriber }
- { from: src_subscriber, to: enrich_subscriber }
- { from: enrich_subscriber, to: enrich_location }
- { from: src_network, to: enrich_location }
- { from: enrich_location, to: quality_gate }
- { from: quality_gate, to: split_by_type }
- { from: quality_gate, to: load_analytics }
- { from: split_by_type, to: load_billing }
- { from: split_by_type, to: load_roaming }
notifications:
onSuccess:
- channel: email
message: "CDR decode complete for {{parameters.cdr_date}}"
onFailure:
- channel: pagerduty
message: "FAILED: CDR decode pipeline for {{parameters.cdr_date}}"
Nigdy nie musisz pisać tego ręcznie — Studio Projektowe tworzy to, gdy przeciągasz węzły. Ale ponieważ jest to zwykły tekst w Git, potok danych jest możliwy do przeglądu, porównywalny między wersjami i możliwy do odzyskania.
Krok 9 — Zaplanuj to
Pliki CDR docierają nieustannie, więc Anna planuje potok danych, by uruchamiał się co godzinę. W Ustawieniach Potoku Danych ustawia wyrażenie cron:
15 * * * *
To oznacza „w minucie 15 każdej godziny" — dając każdej godzinnej partii plików czas na wylądowanie, zanim uruchomienie się rozpocznie. Jeśli uruchomienie się nie powiedzie, potok danych ponawia próbę trzy razy, a powiadomienie onFailure wzywa dyżurny zespół przez PagerDuty.
Alternatywa strumieniowa
Dla najszybszego wykrywania oszustw ta sama logika może działać jako potok danych strumieniowy zamiast godzinnego wsadowego. W trybie Streaming konektor CDR podaje rekordy nieprzerwanie do silnika Apache Flink, przetwarzając każdy rekord w ciągu sekund od wylądowania pliku. Potok danych Anny tutaj to godzinna wersja wsadowa; jeśli Polkomtel później będzie potrzebował alertów o oszustwach roamingowych w czasie poniżej sekundy, zespół zbudowałby wariant strumieniowy, używając tych samych kroków dekodowania i wzbogacania.
Co Anna zbudowała, w skrócie
| Etap | Typ węzła | Cel |
|---|---|---|
| 1 | connector_source (cdr-asn1) | Dekoduj binarne pliki ASN.1 CDR na wiersze |
| 2 | expression | Normalizuj numery telefonów do E.164, dziel znaczniki czasu |
| 3 | joiner ×2 | Wzbogać o taryfę abonenta i lokalizację komórki |
| 4 | quality | Zablokuj błędne dane; oznacz sygnały oszustwa |
| 5 | router | Oddziel rekordy roamingowe od zwykłego zużycia |
| 6 | connector_sink ×3 | Załaduj do billingu, analityki i rozliczeń roamingowych |
| 7 | harmonogram + powiadomienia | Uruchamiaj co godzinę, wzywaj przy niepowodzeniu |
Z zagadkowych plików binarnych, których żaden człowiek nie mógłby odczytać, pojedynczy potok danych Anny produkuje teraz czyste dane, które rozliczają klientów, wyłapują oszustwa, rozliczają konta roamingowe z zagranicznymi operatorami i budują kompletny obraz Subscriber-360 — wszystko automatycznie, co godzinę. To przykład zastosowania CDR telekomu, od początku do końca.