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.

Studio Projektowe
Wizualne budowanie potoku dekodowania CDR w Studio Projektowym.

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 .

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:

FormatPełna nazwaCel prostym językiem
TAP3Transferred Account Procedure v3Ogólnoświatowy standardowy „plik faktury" zużycia roamingowego, wymieniany między operatorami partiami (codziennie lub co tydzień). DataFlow AI obsługuje TAP3.12.
NRTRDENear Real-Time Roaming Data ExchangeSzybsza, 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.
RAPReturned 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 .cdr i 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, .dat lub .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ść:

  1. Lewy panel boczny → Potoki danych → „+ Nowy Potok danych".
  2. Nazwa: cdr_decode_and_enrich.
  3. Opis: „Dekoduj pliki CDR sieci, wzbogać, sprawdź pod kątem oszustw, załaduj do billingu i analityki."
  4. 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ą.)
  5. 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:

UstawienieWartość AnnyCo oznacza
Połączeniepołączenie GCS cdr-asn1Gdzie mieszkają surowe pliki
filePattern*.cdrOdczytaj każdy plik .cdr w folderze
templateTypeautoPozwól dekoderowi samemu wykryć typ każdego rekordu
templateVersion3GPP-R15Którego standardu układ pól użyć
batchSize10000Dekoduj rekordy partiami po 10 000
skipMalformedtrueJeś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 rekorduCzym jest
CS_VOICE_MO / CS_VOICE_MTNormalne połączenie głosowe — wykonane (MO) lub odebrane (MT)
CS_SMS_MO / CS_SMS_MTWiadomość tekstowa — wysłana lub odebrana
PS_DATASesja danych mobilnych — przenosi APN oraz objętości wysyłania/pobierania
IMS_VOICEPołączenie VoLTE (głos przenoszony przez sieć danych 4G)
ROAMING_IN / ROAMING_OUTZuż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:

  • IMSIInternational 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 telefonu po 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:

KontrolaTyp regułyDlaczego ma znaczenie
IMSI abonenta jest zawsze obecneNOT_NULLRekordu bez abonenta nie da się rozliczyć
Czas trwania połączenia mieści się w rozsądnych granicachRANGE30-godzinne „połączenie" to niemal na pewno usterka lub oszustwo
Numer telefonu pasuje do wzorca E.164REGEXWyłapuje numery, których normalizator nie mógł naprawić
Brak zduplikowanych identyfikatorów rekordówUNIQUEDuplikaty rozliczyłyby klienta dwukrotnie
Objętość roamingu nie jest absurdalnie wysokaSTATISTICALNagły ogromny skok za granicą to klasyczny sygnał oszustwa
Pliki faktycznie dotarły w tej godzinieFRESHNESSCicha 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_IN i ROAMING_OUT są 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

EtapTyp węzłaCel
1connector_source (cdr-asn1)Dekoduj binarne pliki ASN.1 CDR na wiersze
2expressionNormalizuj numery telefonów do E.164, dziel znaczniki czasu
3joiner ×2Wzbogać o taryfę abonenta i lokalizację komórki
4qualityZablokuj błędne dane; oznacz sygnały oszustwa
5routerOddziel rekordy roamingowe od zwykłego zużycia
6connector_sink ×3Załaduj do billingu, analityki i rozliczeń roamingowych
7harmonogram + powiadomieniaUruchamiaj 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.

Poprzednia
Raportowanie BI z SAP HANA