ETL w DataFlow

Katalog konektorów

Konektor to wtyczka, która uczy DataFlow AI Platform, jak komunikować się z jednym rodzajem systemu zewnętrznego — bazą danych, chmurową hurtownią danych, magistralą strumieniową, magazynem plików, aplikacją SaaS lub formatem rekordów telekomunikacyjnych. Każdy potok danych odczytuje z konektora i zapisuje do konektora. Ta strona jest kompletną referencją wszystkich 21 wbudowanych konektorów dostarczanych z platformą, pogrupowanych według kategorii, z parametrami konfiguracji, metodami uwierzytelniania i możliwościami każdego z nich.

Rynek konektorów
Rynek konektorów, na którym przegląda się i dodaje konektory.

Jak działają konektory

Każdy konektor implementuje ten sam interfejs Kotlin, ConnectorSDK. Ten współdzielony kontrakt oznacza, że konektor PostgreSQL i konektor Salesforce są konfigurowane i obsługiwane w ten sam sposób, mimo że jeden mówi w JDBC, a drugi w API REST. Platforma odkrywa konektory przy starcie poprzez ConnectorRegistry i udostępnia je na Rynku konektorów.

Współdzielony kontrakt

Każdy konektor dostarcza ten sam zestaw operacji:

MetodaCo robi
discoverSchemaInspekcjonuje źródło i zwraca jego tabele, kolumny i typy
extractDataOdczytuje rekordy z systemu (strona odczytu)
loadDataZapisuje rekordy do systemu (strona zapisu)
testConnectionSprawdza łączność, zwracając opóźnienie i wersję serwera
capabilitiesDeklaruje, co konektor potrafi zrobić
getCDCEventsStrumieniuje zdarzenia zmian, tam gdzie wspierane jest CDC (Change Data Capture)

Operacje i możliwości

Konektory deklarują dwa zestawy faktów o sobie. Operacje to wysokopoziomowe akcje, które konektor wspiera: EXTRACT (odczyt), LOAD (zapis), PUSHDOWN_SQL (uruchamianie SQL wewnątrz źródła), SCHEMA_DISCOVERY, CDC, STREAMING, BULK_COPY i HEARTBEAT. Możliwości to bardziej szczegółowe cechy, których silnik używa do optymalizacji: COLUMN_PROJECTION (odczyt tylko potrzebnych kolumn), PREDICATE_PUSHDOWN (filtrowanie u źródła), SCHEMA_EVOLUTION, PARTITIONING, COMPRESSION, CLOUD_STORAGE, BATCH_WRITE, ENCODING_DETECTION i MULTI_TABLE.

Klasy bazowe

Konektory są budowane na jednej z trzech klas bazowych, w zależności od tego, jak komunikują się ze swoim systemem:

Klasa bazowaUżywana przezUwagi
JDBCConnectorBase9 relacyjnych baz danychUżywa puli połączeń HikariCP (domyślnie maks. 10 połączeń, min. bezczynnych 2). Odkrywanie schematu przez JDBC DatabaseMetaData. Domyślnie transakcyjne wsadowe ładowanie przez wstawianie.
NativeConnectorBaseMongoDB, BigQuery, Kafka, Salesforce, ServiceNow, SAP ERP, REST, CDCUżywa natywnego SDK każdego dostawcy zamiast JDBC. Zarządza własnym cyklem życia połącz/rozłącz.
FileConnectorBaseCSV, Excel, JSON, Parquet, XML, CDRUżywa abstrakcji magazynu, która odczytuje z dysku lokalnego, Google Cloud Storage (gs://), Amazon S3 (s3://) lub Azure Blob.

21 aktywnych konektorów

Platforma rejestruje 21 aktywnych konektorów. Dwa kolejne konektory strumieniowe — Azure Event Hubs i Google Cloud Pub/Sub — istnieją w kodzie źródłowym, ale są obecnie wyłączone w rejestrze konektorów („tymczasowo wyłączone — przenoszenie w toku"), więc nie są opisane poniżej jako aktywne konektory.


Konektory relacyjnych baz danych

Tych osiem konektorów łączy potoki danych z tradycyjnymi bazami danych SQL. Wszystkie działają na JDBC, współdzielą to samo strojenie puli połączeń, wspierają odczyt, zapis, ładowanie wsadowe, odkrywanie schematu i przepychanie SQL — a niektóre z nich również Change Data Capture.

PostgreSQL — postgresql

Łączy się z bazą danych PostgreSQL, szeroko używaną relacyjną bazą danych typu open source. Używaj jej jako źródła lub celu dla danych transakcyjnych oraz jako źródła CDC do strumieniowania zmian na poziomie wierszy.

  • Operacje: odczyt, zapis, przepychanie SQL, odkrywanie schematu, kopiowanie zbiorcze, CDC.
  • Uwierzytelnianie: password (md5 lub scram-sha-256), certificate (certyfikat klienta SSL) lub gcp_iam (krótkotrwałe tokeny OAuth2 dla Google Cloud SQL).
  • Odczyt za pomocą standardowego strumieniowania JDBC lub szybkiego protokołu COPY TO STDOUT.
  • Zapis domyślnie za pomocą protokołu ładowania zbiorczego COPY lub upsertów INSERT ... ON CONFLICT.
  • CDC: oparty na dzienniku Write-Ahead Log (WAL) z użyciem slotów replikacji logicznej. Wspierane są wtyczki pgoutput, wal2json i test_decoding.
  • Ograniczenia: CDC wymaga uprzednio utworzonego slotu replikacji (pg.cdc_slot_name).
ParametrTypWymaganyOpis
hoststringtakNazwa hosta serwera (domyślnie localhost)
portintniePort (domyślnie 5432)
databasestringtakNazwa bazy danych (domyślnie postgres)
username / passwordstringtakPoświadczenia logowania
schemastringnieNazwa schematu (domyślnie public)
pg.auth_modestringniepassword, certificate lub gcp_iam
pg.sslmodestringnieverify-full gdy SSL jest włączone, w przeciwnym razie prefer
pg.use_copyboolnieUżyj ładowania zbiorczego COPY (domyślnie true)
pg.cdc_slot_namestringdla CDCNazwa slotu replikacji logicznej
pg.cdc_pluginstringnieWtyczka CDC (domyślnie test_decoding)
connector: postgresql
config:
  host: pg-prod.polkomtel.internal
  port: 5432
  database: billing
  schema: public
  username: dataflow
  password: ${PG_PASSWORD}
  pg.use_copy: true

MySQL — mysql

Łączy się z bazą danych MySQL (wersje 5.7 i 8.0+), inną popularną relacyjną bazą danych typu open source. Odpowiednia dla obciążeń transakcyjnych i CDC opartego na dzienniku binarnym.

  • Operacje: odczyt, zapis, CDC, przepychanie SQL, odkrywanie schematu, kopiowanie zbiorcze.
  • Uwierzytelnianie: mysql_native_password lub caching_sha2_password (domyślne dla MySQL 8), z trybami SSL od DISABLED do VERIFY_IDENTITY.
  • Zapis w czterech trybach wstawiania: STANDARD, IGNORE, REPLACE i UPSERT (ON DUPLICATE KEY UPDATE), plus zbiorcza ścieżka LOAD DATA LOCAL INFILE.
  • CDC: odczytuje dziennik binarny (binlog), adresowany przez plik/pozycję lub przez GTID.
  • Ograniczenia: CDC wymaga unikatowego mysql.cdc.serverId.
ParametrTypWymaganyOpis
host / portstring / inthost takAdres serwera (domyślny port 3306)
database / username / passwordstringtakSzczegóły połączenia
mysql.sslModestringniePoziom egzekwowania SSL
mysql.insertModestringnieSTANDARD, IGNORE, REPLACE, UPSERT
mysql.streamingResultboolnieStrumieniuj duże zbiory wyników
mysql.cdc.serverIdlongdla CDCUnikatowy identyfikator serwera repliki
mysql.cdc.useGtidboolnieŚledź zmiany przez GTID zamiast pozycji
connector: mysql
config:
  host: mysql-crm.polkomtel.internal
  database: crm
  username: dataflow
  password: ${MYSQL_PASSWORD}
  mysql.insertMode: UPSERT
  mysql.sslMode: REQUIRED

Oracle Database — oracle

Łączy się z Oracle Database, korporacyjną relacyjną bazą danych powszechną w dużych systemach back-office operatorów telekomunikacyjnych. Używaj jej dla źródeł transakcyjnych i ładowań o wysokiej przepustowości.

  • Operacje: odczyt, zapis, przepychanie SQL, odkrywanie schematu, kopiowanie zbiorcze.
  • Uwierzytelnianie: hasło, Oracle Wallet / mTLS (dla Autonomous Database) lub Kerberos.
  • Tryby połączenia: service_name, sid lub pełny deskryptor tns.
  • Zapis za pomocą wstawiania wsadowego, ładowania bezpośrednią ścieżką (INSERT /*+ APPEND */) lub MERGE.
  • Ograniczenia: brak wsparcia CDC.
ParametrTypWymaganyOpis
host / portstring / inthost takAdres serwera (domyślny port 1521)
database / username / passwordstringtakSzczegóły połączenia
oracle.connection_typestringnieservice_name (domyślnie), sid lub tns
oracle.wallet_locationstringdla uwierzytelniania walletŚcieżka do Oracle Wallet
oracle.direct_pathboolnieUżyj wskazówki bezpośredniej ścieżki APPEND
oracle.batch_sizeintnieRozmiar wsadu wstawiania (domyślnie 10000)
connector: oracle
config:
  host: oracle-ar.polkomtel.internal
  port: 1521
  oracle.connection_type: service_name
  database: ARPROD
  username: dataflow
  password: ${ORA_PASSWORD}
  oracle.direct_path: true

Microsoft SQL Server — mssql

Łączy się z Microsoft SQL Server lub Azure SQL Database. Używaj go dla danych transakcyjnych oraz jako źródła CDC poprzez natywne tabele zmian SQL Server.

  • Operacje: odczyt, zapis, CDC, przepychanie SQL, odkrywanie schematu, kopiowanie zbiorcze.
  • Uwierzytelnianie: uwierzytelnianie SQL, uwierzytelnianie zintegrowane Windows oraz pięć trybów Azure Active Directory (ActiveDirectoryPassword, ActiveDirectoryIntegrated, ActiveDirectoryMSI, ActiveDirectoryServicePrincipal, ActiveDirectoryInteractive).
  • Zapis za pomocą SQLServerBulkCopy z opcjami rozmiaru wsadu, blokowania tabel i wstawiania tożsamości.
  • CDC: odczytuje natywne tabele zmian CDC; baza danych i tabele docelowe muszą mieć najpierw włączone CDC.
  • Możliwości: kolumny Always Encrypted oraz przełączanie awaryjne AlwaysOn w wielu podsieciach.
ParametrTypWymaganyOpis
host / portstring / inthost takAdres serwera (domyślny port 1433)
database / username / passwordstringtakSzczegóły połączenia
bulkCopy.enabledboolnieUżyj SQLServerBulkCopy do ładowań
cdc.captureInstancestringdla CDCNazwa instancji przechwytywania
cdc.fromLsnstringnieLog Sequence Number, od którego rozpocząć CDC
connector: mssql
config:
  host: sqlsrv-ops.polkomtel.internal
  database: Operations
  username: dataflow
  password: ${MSSQL_PASSWORD}
  bulkCopy.enabled: true

Teradata — teradata

Łączy się z Teradata, masowo równoległą hurtownią danych używaną do analityki telekomunikacyjnej na dużą skalę. Używaj jej dla odczytów o dużym wolumenie i ładowań zbiorczych.

  • Operacje: odczyt, zapis, kopiowanie zbiorcze, przepychanie SQL, odkrywanie schematu.
  • Uwierzytelnianie (auth_method): TD2 (natywne, domyślne), LDAP lub KRB5 (Kerberos).
  • Tryby ładowania: fastload (wstawienia o wysokiej przepustowości do pustych tabel, domyślne), multiload (upsert MERGE), stream (wiersz po wierszu o niskim opóźnieniu) lub batch.
  • Możliwości: odczyty ACCESS LOCK, analiza skosu AMP oraz tagowanie sesji query-band.
  • Ograniczenia: multiload i stream wymagają key_columns.
ParametrTypWymaganyOpis
host / username / passwordstringtakSzczegóły połączenia (domyślny port 1025)
auth_methodstringnieTD2, LDAP lub KRB5
load_modestringniefastload, multiload, stream, batch
key_columnsstringdla multiload/streamKolumny klucza oddzielone przecinkami
query_bandstringnieEtykieta tagowania sesji
connector: teradata
config:
  host: teradata-dw.polkomtel.internal
  database: ANALYTICS
  username: dataflow
  password: ${TD_PASSWORD}
  auth_method: LDAP
  load_mode: fastload

SAP HANA — sap-hana

Łączy się z SAP HANA, kolumnową bazą danych SAP działającą w pamięci (lokalnie lub HANA Cloud). Używaj jej dla źródeł analitycznych i jako celu ładowania.

  • Operacje: odczyt, zapis, przepychanie SQL, odkrywanie schematu, kopiowanie zbiorcze.
  • Uwierzytelnianie (authentication): PASSWORD (domyślne), X509 (certyfikat klienta), KERBEROS, SAML lub JWT.
  • Zapis za pomocą UPSERT, zoptymalizowanego wstawiania wsadowego lub po stronie serwera IMPORT FROM / EXPORT INTO CSV.
  • Możliwości: odczyt widoków obliczeniowych SAP HANA; filtrowanie tylko dla magazynu kolumnowego.
  • Ograniczenia: brak wsparcia CDC.
ParametrTypWymaganyOpis
host / portstring / inthost takAdres serwera (najemca 30015, HANA Cloud 443)
database / username / passwordstringtakSzczegóły połączenia
authenticationstringniePASSWORD, X509, KERBEROS, SAML, JWT
tokenstringdla SAML/JWTAsercja bearer lub token
includeCalculationViewsboolnieOdkryj widoki obliczeniowe _SYS_BIC
connector: sap-hana
config:
  host: hana.polkomtel.internal
  port: 30015
  database: HXE
  username: DATAFLOW
  password: ${HANA_PASSWORD}
  authentication: PASSWORD

Sybase ASE — sybase

Łączy się z Sybase Adaptive Server Enterprise (ASE), starszą korporacyjną relacyjną bazą danych, za pośrednictwem sterownika open source jTDS. Używaj go do ekstrakcji ze starszych systemów telekomunikacyjnych.

  • Operacje: odczyt, zapis, przepychanie SQL, odkrywanie schematu, kopiowanie zbiorcze.
  • Uwierzytelnianie: natywna nazwa użytkownika i hasło Sybase.
  • Zapis za pomocą wsadowych przygotowanych instrukcji z automatycznym ponawianiem przy zakleszczeniu lub eksportu w stylu BCP-out rozdzielanego tabulatorami.
  • Ograniczenia: brak wsparcia CDC.
ParametrTypWymaganyOpis
host / portstring / inthost takAdres serwera (domyślny port 5000)
database / username / passwordstringtakSzczegóły połączenia
sybase.use_bcpboolnieUżyj wsadowych ładowań kopiowania zbiorczego (domyślnie true)
connector: sybase
config:
  host: sybase-legacy.polkomtel.internal
  port: 5000
  database: legacy_billing
  username: dataflow
  password: ${SYBASE_PASSWORD}
  sybase.use_bcp: true

Konektor bazy danych NoSQL

MongoDB — mongodb

Łączy się z MongoDB, zorientowaną na dokumenty bazą danych NoSQL, która przechowuje elastyczne dokumenty podobne do JSON zamiast wierszy o stałej strukturze tabel. Używaj go dla danych częściowo ustrukturyzowanych i jako źródła CDC poprzez strumienie zmian. W przeciwieństwie do konektorów relacyjnych, MongoDB używa swojego natywnego sterownika zamiast JDBC.

  • Operacje / możliwości: odkrywanie schematu, odczyt, zapis, projekcja kolumn, przepychanie predykatów, zapis wsadowy, streaming, CDC.
  • Uwierzytelnianie: SCRAM-SHA-256 (domyślne), SCRAM-SHA-1, certyfikat klienta X.509, AWS IAM lub LDAP.
  • Odczyt za pomocą zapytań find() lub pełnych potoków agregacji. Schemat jest wnioskowany przez próbkowanie dokumentów (domyślnie 100), ponieważ MongoDB nie ma stałego schematu.
  • Zapis w pięciu trybach: insert, upsert, replace, update i bulk.
  • CDC: strumienie zmian za pomocą API watch(), z tokenami wznowienia zapewniającymi odporność na awarie.
ParametrTypWymaganyOpis
uristringtakŁańcuch połączenia mongodb:// lub mongodb+srv://
database / collectionstringtakDocelowa baza danych i kolekcja
mongodb.writeModestringnieinsert, upsert, replace, update, bulk
schemaSampleSizeintnieDokumenty do próbkowania dla schematu (domyślnie 100)
mongodb.aggregationPipelinestringniePotok agregacji dla odczytów
connector: mongodb
config:
  uri: mongodb+srv://cluster.polkomtel.mongodb.net
  database: customer360
  collection: events
  mongodb.writeMode: upsert
  schemaSampleSize: 200

Konektory chmurowych hurtowni danych

Te konektory łączą potoki danych z hostowanymi w chmurze hurtowniami analitycznymi — usługami zarządzanymi zbudowanymi do zapytań na dużą skalę.

Snowflake — snowflake

Łączy się z Snowflake Data Cloud, w pełni zarządzaną chmurową hurtownią danych. Używaj jej jako źródła analitycznego o dużym wolumenie lub celu ładowania.

  • Operacje: odczyt, zapis, kopiowanie zbiorcze, przepychanie SQL, odkrywanie schematu.
  • Uwierzytelnianie (auth_type): password (domyślne), key_pair (JWT z parą kluczy RSA), oauth lub external_browser (interaktywne SSO).
  • Odczyt za pomocą strumieniowania JDBC lub zbiorczego wyładowania COPY INTO z etapu.
  • Zapis za pomocą etapowanych ładowań PUT + COPY INTO lub wstawiania wsadowego JDBC.
  • Możliwości: zarządzanie automatycznym wstrzymywaniem/wznawianiem hurtowni oraz tagowanie zapytań.
ParametrTypWymaganyOpis
accountstringtakIdentyfikator konta Snowflake
database / schema / warehousestringtakObiekty docelowe
rolestringnieRola Snowflake do przyjęcia
username / passwordstringdla uwierzytelniania hasłemPoświadczenia logowania
auth_typestringniepassword, key_pair, oauth, external_browser
stage_namestringnieEtap dla ładowania/wyładowania zbiorczego
connector: snowflake
config:
  account: polkomtel-eu
  database: ANALYTICS
  schema: PUBLIC
  warehouse: LOAD_WH
  username: DATAFLOW
  password: ${SF_PASSWORD}
  auth_type: password

Databricks — databricks

Łączy się z Databricks Lakehouse, zunifikowaną platformą analityczną zbudowaną na Delta Lake. Używaj go do analityki lakehouse, ze wsparciem dla podróży w czasie Delta.

  • Operacje: odczyt, zapis, kopiowanie zbiorcze, przepychanie SQL, odkrywanie schematu.
  • Uwierzytelnianie (auth_type): pat (Personal Access Token, domyślne), oauth (machine-to-machine) lub azure_ad.
  • Obliczenia (compute_type): SQL Warehouse lub klaster All-Purpose.
  • Odczyt wspiera podróż w czasie Delta Lake — odczyt tabeli według przeszłej wersji lub znacznika czasu.
  • Zapis za pomocą COPY INTO z magazynu chmurowego, upsertów Delta MERGE INTO lub wsadu JDBC.
  • Możliwości: trójpoziomowa przestrzeń nazw Unity Catalog oraz utrzymanie tabel Delta (OPTIMIZE, VACUUM).
ParametrTypWymaganyOpis
hoststringtakHost obszaru roboczego Databricks
httpPathstringtakŚcieżka HTTP SQL Warehouse lub klastra
auth_typestringniepat, oauth, azure_ad
compute_typestringniesql_warehouse lub cluster
extractAtVersionintnieWersja Delta do odczytu (podróż w czasie)
merge_keysstringnieKlucze dla upsertów Delta MERGE
connector: databricks
config:
  host: adb-123.4.azuredatabricks.net
  httpPath: /sql/1.0/warehouses/abc123
  auth_type: pat
  token: ${DATABRICKS_TOKEN}
  compute_type: sql_warehouse

Google BigQuery — bigquery

Łączy się z Google BigQuery, bezserwerową hurtownią danych Google Cloud. Używaj go do analityki na dużą skalę i jako celu ładowania. Używa natywnego SDK BigQuery zamiast JDBC.

  • Operacje / możliwości: odczyt, zapis, zapis wsadowy, streaming, partycjonowanie, projekcja kolumn, odkrywanie schematu, kopiowanie zbiorcze.
  • Uwierzytelnianie: Application Default Credentials (Workload Identity) lub klucz JSON konta usługi.
  • Odczyt za pomocą zapytań Standard SQL lub Legacy SQL, ze wsparciem dla zapytań parametryzowanych.
  • Zapis za pomocą wsadowych zadań ładowania (CSV/JSON/Parquet/Avro) lub wstawień strumieniowych (do 10 000 wierszy na żądanie).
  • Możliwości: tabele partycjonowane i klastrowane, tabele zmaterializowane i zewnętrzne.
  • Ograniczenia: brak przepychania SQL ani CDC.
ParametrTypWymaganyOpis
projectIdstringtakIdentyfikator projektu Google Cloud
datasetIdstringtakZbiór danych BigQuery
locationstringnieRegion zbioru danych (domyślnie US)
serviceAccountKeyPathstringnieŚcieżka do klucza JSON konta usługi
useStreamingboolnieUżyj wstawień strumieniowych zamiast zadań ładowania
partitionFieldstringnieKolumna do partycjonowania tabeli docelowej
connector: bigquery
config:
  projectId: polkomtel-data
  datasetId: telecom_analytics
  location: EU
  serviceAccountKeyPath: /secrets/bq-sa.json
  useStreaming: false

Konektory strumieniowe

Konektory strumieniowe przenoszą dane jako ciągły przepływ zdarzeń, a nie jako ograniczone wsady.

Apache Kafka — kafka

Łączy się z Apache Kafka, rozproszoną platformą strumieniowania zdarzeń. Używaj go do konsumowania zdarzeń z tematów (odczyt) lub publikowania zdarzeń do tematów (zapis).

  • Operacje: odczyt (konsument), zapis (producent), odkrywanie schematu.
  • Uwierzytelnianie (security.protocol): PLAINTEXT (domyślne), SASL_PLAINTEXT, SASL_SSL lub SSL (wzajemne TLS). Mechanizmy SASL obejmują PLAIN, SCRAM-SHA-256/512 i OAUTHBEARER.
  • Odczyt: subskrybuje tematy, ręcznie zatwierdza przesunięcia (offsety) i automatycznie wstrzykuje kolumny metadanych (_topic, _partition, _offset, _timestamp, _key). Niepowodzenia deserializacji są kierowane do tematu martwych listów.
  • Zapis: produkuje rekordy z konfigurowalnymi potwierdzeniami, wsadowaniem i opcjonalnym dostarczaniem idempotentnym.
  • Odkrywanie schematu: z Confluent Schema Registry lub wnioskowane z próbkowanych komunikatów JSON.
  • Ograniczenia: Kafka nie udostępnia CDC bazy danych — do tego użyj konektora cdc-debezium.
ParametrTypWymaganyOpis
bootstrap.serversstringtakAdresy brokerów Kafka (domyślny port 9092)
topics / topicstringtakTemat(y) do konsumowania lub produkowania
security.protocolstringnieProtokół bezpieczeństwa połączenia
sasl.mechanismstringnieMechanizm uwierzytelniania SASL
schema.registry.urlstringnieURL Confluent Schema Registry
serde.formatstringnieFormat serializacji (domyślnie JSON)
connector: kafka
config:
  bootstrap.servers: kafka-1.polkomtel.internal:9092
  topics: cdr-events
  security.protocol: SASL_SSL
  sasl.mechanism: SCRAM-SHA-256
  serde.format: JSON

Change Data Capture — cdc-debezium

Łączy się ze źródłową bazą danych, aby przechwytywać zdarzenia Change Data Capture (CDC) — ciągły strumień każdego wstawienia, aktualizacji i usunięcia wiersza. Zbudowany na Debezium, zamienia zwykłą bazę danych w źródło zdarzeń czasu rzeczywistego bez zmiany własnego obciążenia bazy danych.

  • Wspierane źródłowe bazy danych: Oracle (dzienniki redo LogMiner), PostgreSQL (replikacja logiczna WAL), SQL Server (natywne tabele CDC), MySQL (binlog) i MongoDB (strumienie zmian).
  • Możliwości: odkrywanie schematu za pomocą migawki, wsadowa i ciągła ekstrakcja zdarzeń zmian, monitorowanie opóźnienia i kondycji oraz maskowanie kolumn dla PII.
  • Ograniczenia: tylko do odczytu — konektory CDC przechwytują zmiany, nie zapisują.
ParametrTypWymaganyOpis
cdc.database.typestringtakoracle, postgresql, sqlserver, mysql, mongodb
cdc.connection.urlstringtakURL źródłowej bazy danych
cdc.username / cdc.passwordstringtakPoświadczenia źródła
cdc.tablesstringtakTabele, z których przechwytywać zmiany
cdc.snapshot.modestringnieinitial, schema_only, never, when_needed
cdc.start.fromstringniebeginning, latest lub znacznik czasu
cdc.column.masksstringnieKolumny do maskowania dla PII
connector: cdc-debezium
config:
  cdc.database.type: postgresql
  cdc.connection.url: jdbc:postgresql://pg-prod:5432/billing
  cdc.username: debezium
  cdc.password: ${CDC_PASSWORD}
  cdc.tables: public.invoices,public.payments
  cdc.snapshot.mode: initial

Konektory formatów plików

Tych pięć konektorów odczytuje i zapisuje pliki. Wszystkie korzystają ze współdzielonej abstrakcji magazynu plików, więc ten sam konektor działa względem dysku lokalnego, Google Cloud Storage, Amazon S3 lub Azure Blob — backend magazynu jest wykrywany na podstawie prefiksu ścieżki (gs://, s3://).

CSV — csv

Odczytuje i zapisuje pliki CSV (comma-separated values) oraz inne pliki tekstowe z separatorami. Używaj go dla najpowszechniejszego formatu wymiany danych tabelarycznych.

  • Możliwości: odkrywanie schematu, odczyt, zapis, projekcja kolumn, wykrywanie kodowania, magazyn chmurowy, zapis wsadowy.
  • Funkcje: pełne wsparcie RFC 4180 (pola cytowane, wartości wielowierszowe); automatyczne wykrywanie kodowania znaków i separatora; wnioskowanie typów (integer, decimal, date, timestamp, boolean, string); dekompresja gzip i bzip2.
ParametrTypWymaganyOpis
pathstringtakŚcieżka pliku lub katalogu
delimiterstringnieSeparator pól (wykrywany automatycznie, jeśli pominięty)
hasHeaderboolniePierwszy wiersz jest nagłówkiem
encodingstringnieKodowanie znaków (wykrywane automatycznie, jeśli pominięte)
dialectstringnierfc4180, excel, excel_eu, tsv, pipe
errorModestringnieFAIL lub SKIP przy nieprawidłowych wierszach
connector: csv
config:
  path: gs://polkomtel-imports/subscribers.csv
  delimiter: ","
  hasHeader: true
  errorMode: SKIP

Microsoft Excel — excel

Odczytuje i zapisuje skoroszyty Microsoft Excel (.xlsx oraz starsze .xls). Używaj go dla arkuszy kalkulacyjnych wymienianych z zespołami biznesowymi.

  • Możliwości: odkrywanie schematu, odczyt, zapis, projekcja kolumn, multi-table, zapis wsadowy.
  • Funkcje: wybór arkusza według nazwy lub indeksu; rozwiązywanie scalonych komórek; odczyty nazwanych zakresów; ewaluacja formuł z fallbackiem na buforowany wynik; zapisy strumieniowe dla dużych plików.
  • Ograniczenia: pliki chronione hasłem są odrzucane.
ParametrTypWymaganyOpis
pathstringtakŚcieżka pliku skoroszytu
sheetName / sheetIndexstring / intnieKtóry arkusz odczytać
headerRowintnieIndeks wiersza nagłówka
evaluateFormulasboolniePrzeliczaj formuły przy odczycie
namedRangestringnieOdczytaj tylko nazwany zakres
connector: excel
config:
  path: /data/reports/monthly.xlsx
  sheetName: Subscribers
  headerRow: 1
  evaluateFormulas: true

JSON — json

Odczytuje i zapisuje pliki JSON. Używaj go dla zagnieżdżonych, częściowo ustrukturyzowanych dokumentów. Schemat jest wnioskowany automatycznie z zawartości pliku.

  • Możliwości: odkrywanie schematu, odczyt, zapis, magazyn chmurowy.
ParametrTypWymaganyOpis
pathstringtakŚcieżka pliku lub katalogu
connector: json
config:
  path: s3://polkomtel-events/2026/05/events.json

Parquet — parquet

Odczytuje i zapisuje pliki Apache Parquet — skompresowany, kolumnowy format magazynowania zbudowany do analityki. Używaj go do wydajnej wymiany dużych zbiorów danych między systemami analitycznymi.

  • Możliwości: odkrywanie schematu, odczyt, zapis, magazyn chmurowy. Schemat jest mapowany bezpośrednio z wbudowanego schematu pliku Parquet.
ParametrTypWymaganyOpis
pathstringtakŚcieżka pliku lub katalogu
connector: parquet
config:
  path: gs://polkomtel-lake/cdr/voice/part-0001.parquet

XML — xml

Odczytuje i zapisuje pliki XML. Używaj go dla dokumentów hierarchicznych i eksportów ze starszych systemów. Schemat jest mapowany ze struktury XML.

  • Możliwości: odkrywanie schematu, odczyt, zapis, magazyn chmurowy.
ParametrTypWymaganyOpis
pathstringtakŚcieżka pliku lub katalogu
connector: xml
config:
  path: /data/exports/inventory.xml

Konektory SaaS i aplikacji

Te konektory łączą potoki danych z aplikacjami chmurowymi poprzez ich API REST lub OData.

Salesforce CRM — salesforce

Łączy się z Salesforce, chmurową platformą CRM, poprzez jej API REST i SOQL. Używaj go do ekstrakcji danych klientów i sprzedaży lub do ładowania rekordów z powrotem.

  • Operacje / możliwości: odczyt, zapis, odkrywanie schematu, kopiowanie zbiorcze, zapis wsadowy, projekcja kolumn.
  • Uwierzytelnianie (auth_type): password (OAuth2 nazwa użytkownika-hasło, domyślne), jwt_bearer (serwer-serwer) lub refresh_token.
  • Odczyt za pomocą zapytań SOQL z paginacją kursorową. Schemat jest odkrywany poprzez API describe.
  • Zapis za pomocą wywołań REST pojedynczego rekordu lub Bulk API 2.0 dla dużych wolumenów.
  • Możliwości: przechodzenie relacji, pola polimorficzne i złożone; respektuje limity szybkości API Salesforce.
ParametrTypWymaganyOpis
loginUrlstringnielogin.salesforce.com lub test.salesforce.com
auth_typestringniepassword, jwt_bearer, refresh_token
clientId / clientSecretstringtakPoświadczenia Connected App
username / passwordstringdla uwierzytelniania hasłemPoświadczenia logowania
securityTokenstringdla uwierzytelniania hasłemToken bezpieczeństwa Salesforce
private_key_pathstringdla jwt_bearerŚcieżka klucza prywatnego RSA
connector: salesforce
config:
  loginUrl: login.salesforce.com
  auth_type: password
  clientId: ${SF_CLIENT_ID}
  clientSecret: ${SF_CLIENT_SECRET}
  username: integration@polkomtel.com
  password: ${SF_PASSWORD}
  securityToken: ${SF_TOKEN}

SAP ERP / S/4HANA — sap-erp

Łączy się z SAP ERP lub S/4HANA poprzez API OData (v2 lub v4). Używaj go do ekstrakcji danych planowania zasobów przedsiębiorstwa, takich jak finanse, materiały i zamówienia.

  • Operacje / możliwości: odczyt, zapis, odkrywanie schematu, kopiowanie zbiorcze, przepychanie predykatów, zapis wsadowy, multi-table.
  • Uwierzytelnianie (authType): BASIC, OAUTH2 lub PRINCIPAL_PROPAGATION (SSO przez SAP Cloud Connector).
  • Odczyt: schemat z dokumentu OData $metadata; ekstrakcja z $select, $filter i stronicowaniem sterowanym przez serwer.
  • Zapis za pomocą tworzenia/aktualizacji OData lub żądania wieloczęściowego $batch.
  • Specyfika SAP: pobranie tokenu CSRF, nagłówki sap-client i sap-language.
ParametrTypWymaganyOpis
host / portstring / inthost takSerwer SAP (domyślny port 443)
serviceNamestringtakNazwa usługi OData
authTypestringnieBASIC, OAUTH2, PRINCIPAL_PROPAGATION
sapClientstringnieNumer klienta SAP (domyślnie 100)
odataVersionintnieWersja OData (domyślnie 2)
entitySetstringnieZestaw encji do ekstrakcji
useBatchboolnieUżyj $batch do zapisów
connector: sap-erp
config:
  host: s4hana.polkomtel.internal
  serviceName: API_SALES_ORDER_SRV
  authType: BASIC
  username: DATAFLOW
  password: ${SAP_PASSWORD}
  sapClient: "100"
  odataVersion: 4

ServiceNow — servicenow

Łączy się z ServiceNow, chmurową platformą zarządzania usługami IT (ITSM), poprzez jej API REST Table. Używaj go do ekstrakcji incydentów, zasobów i innych rekordów ITSM.

  • Operacje / możliwości: odczyt, zapis, odkrywanie schematu, przepychanie predykatów, zapis wsadowy, projekcja kolumn.
  • Uwierzytelnianie (authType): basic (domyślne) lub oauth2 (poświadczenia klienta).
  • Odczyt: schemat z sys_dictionary; stronicowana ekstrakcja z zakodowanymi zapytaniami i wyborem pól.
  • Zapis za pomocą API Table lub API Import Set dla ładowań zbiorczych.
ParametrTypWymaganyOpis
instanceUrlstringtakURL instancji ServiceNow
authTypestringniebasic lub oauth2
clientId / clientSecretstringdla oauth2Poświadczenia OAuth2
sysparmQuerystringnieFiltr zakodowanego zapytania
sysparmFieldsstringniePola do zwrócenia
pageSizeintnieRozmiar strony (domyślnie 1000, maks. 10000)
importTablestringnieTabela etapowania import-set dla ładowań zbiorczych
connector: servicenow
config:
  instanceUrl: https://polkomtel.service-now.com
  authType: basic
  username: dataflow
  password: ${SNOW_PASSWORD}
  sysparmQuery: active=true
  pageSize: 1000

Generyczne API REST — rest_api

Łączy się z dowolnym punktem końcowym HTTP REST. Używaj go, gdy system nie ma dedykowanego konektora, ale udostępnia API REST.

  • Operacje / możliwości: odczyt, zapis, odkrywanie schematu, zapis wsadowy.
  • Uwierzytelnianie (authType): none, basic, bearer, oauth2_cc (poświadczenia klienta OAuth2), api_key lub custom_header.
  • Odczyt: konfigurowalna metoda HTTP, formaty odpowiedzi (json, xml, csv, ndjson) oraz sześć strategii paginacji (none, offset, cursor, page_number, link_header, keyset).
  • Zapis: wsadowe POST/PUT/PATCH tablic JSON.
  • Niezawodność: ograniczanie szybkości metodą token bucket, ponawianie z wykładniczym wycofaniem przy 429/5xx oraz wyłącznik bezpieczeństwa.
ParametrTypWymaganyOpis
baseUrlstringtakBazowy URL API
pathstringnieŚcieżka punktu końcowego
methodstringnieMetoda HTTP
authTypestringnieSchemat uwierzytelniania
responseFormatstringniejson, xml, csv, ndjson
responsePathstringnieŚcieżka kropkowa do tablicy danych
paginationTypestringnieStrategia paginacji
rateLimitPerSecondintnieLimit szybkości żądań (domyślnie 10)
connector: rest_api
config:
  baseUrl: https://api.partner.example.com
  path: /v1/usage
  method: GET
  authType: bearer
  responseFormat: json
  responsePath: data.items
  paginationType: cursor

Telekomunikacyjny konektor CDR — cdr-asn1

Konektor CDR to specjalnie zbudowany telekomunikacyjny konektor platformy. Odczytuje on rekordy szczegółów połączeń (CDR — Call Detail Records) — binarne rekordy, które sieć mobilna generuje dla każdego połączenia, SMS-a, sesji danych i zdarzenia roamingu. CDR-y są surowcem dla rozliczeń, zapewnienia przychodów i wykrywania nadużyć u operatora telekomunikacyjnego takiego jak Polkomtel.

CDR-y nie są przechowywane jako zwykły tekst. Są zakodowane w ASN.1 (Abstract Syntax Notation One), binarnym standardzie zdefiniowanym przez organy telekomunikacyjne 3GPP i ETSI. Konektor CDR dekoduje ASN.1 bezpośrednio, więc potoki danych mogą pobierać rekordy sieciowe bez osobnego etapu dekodowania.

  • Klasa bazowa / status: rozszerza FileConnectorBase; tylko do odczytu — nie może zapisywać plików CDR (wyjście Parquet jest wytwarzane przez osobny konwerter).
  • Możliwości: odkrywanie schematu, odczyt, streaming, magazyn chmurowy.
  • Obsługa plików: odczytuje pojedyncze pliki lub całe katalogi za pomocą wzorca glob (domyślnie *.cdr); wspierane rozszerzenia .cdr, .asn1, .dat, .bin.
  • Magazyn chmurowy: odczytuje z GCS, S3 lub dysku lokalnego — konwencją Polkomtel jest gs://polkomtel-cdr-raw/{date}/.

Potok dekodowania ASN.1

Plik CDR jest dekodowany etapami:

EtapCo się dzieje
Odczyt z magazynuAdapter magazynu strumieniuje surowy plik z GCS, S3 lub dysku lokalnego
Dekodowanie ASN.1Asn1Decoder parsuje kodowanie binarne BER/DER do struktur tag-długość-wartość
Dekodowanie binarneCdrBinaryDecoder stosuje szablon do wyodrębnienia nazwanych pól z tych struktur
Mapowanie rekordówZdekodowane pola stają się standardowymi wierszami DataRecord w DataStream

Typy danych telekomunikacyjnych są mapowane na typy platformy — w szczególności TBCD_STRING (kodowanie TBCD używane dla identyfikatorów numerów telefonów IMSI i MSISDN), BCD_TIMESTAMP i ASN1_TIMESTAMP, LOCATION_INFO (lokalizacja komórki sieciowej: MCC, MNC, LAC, Cell ID), IP_ADDRESS oraz ENUMERATED (kody przyczyn).

Atrybucja przełącznika źródłowego i znaczniki wodne

Konektor wyprowadza pochodzący przełącznik sieciowy z nazwy pliku (huaweihuawei-mss, ericssonericsson-msc, nokianokia-msc), aby można było śledzić metryki per dostawca. Klasyfikuje również terminowość każdego rekordu — ON_TIME, LATE lub DROPPED — za pomocą śledzenia znaczników wodnych, dołączając wynik jako metadane _lateness obok _recordType, _timestamp i _sourceSwitch.

Natywne typy rekordów CDR

Konektor dostarcza szablony dla 12 natywnych typów rekordów 3GPP-R15:

KodTypOpis
0CS_VOICE_MOPołączenie głosowe z komutacją łączy, zainicjowane przez telefon
1CS_VOICE_MTPołączenie głosowe z komutacją łączy, zakończone na telefonie
2CS_SMS_MOSMS z komutacją łączy, zainicjowany przez telefon
3CS_SMS_MTSMS z komutacją łączy, zakończony na telefonie
4PS_DATADane z komutacją pakietów (GPRS/LTE — APN, wolumeny)
5IMS_VOICEGłos VoLTE/IMS
6IMS_VIDEOPołączenie wideo IMS
7MMS_MOMMS, zainicjowany przez telefon
8MMS_MTMMS, zakończony na telefonie
9ROAMING_INRoaming przychodzący
10ROAMING_OUTRoaming wychodzący
11SUPPLEMENTARY_SERVICEUsługa uzupełniająca

Konektor automatycznie wykrywa, który szablon zastosować: odczytuje kod typu rekordu (znacznik ASN.1 0), wraca do heurystyk wzorca znaczników, a na końcu do wykrywania formatu roamingu.

Parametry konfiguracji

ParametrTypWymaganyOpis
path / filePatternstringpath takPlik lub katalog; glob (domyślnie *.cdr)
templateTypestringnieSzablon rekordu (domyślnie auto dla autowykrywania)
templateVersionstringnieWersja standardu (domyślnie 3GPP-R15)
batchSizeintnieRekordów na wsad (domyślnie 10000)
skipMalformedboolniePomiń niedekodowalne rekordy (domyślnie true)
sourceSwitchstringnieNadpisz nazwę przełącznika wyprowadzoną z nazwy pliku
connector: cdr-asn1
config:
  path: gs://polkomtel-cdr-raw/2026-05-20/
  filePattern: "*.cdr"
  templateType: auto
  templateVersion: 3GPP-R15
  batchSize: 10000
  skipMalformed: true

Formaty rozliczeń roamingowych — TAP3, NRTRDE, RAP

Gdy abonent Polkomtel korzysta z sieci innego operatora za granicą — lub zagraniczny abonent roamuje w sieci Polkomtel — obaj operatorzy muszą wymienić rekordy użycia i rozliczyć opłaty. Branża używa do tego trzech formatów CDR ustandaryzowanych przez GSMA, a konektor cdr-asn1 dekoduje wszystkie trzy. Gdy widzi plik roamingowy, konektor próbuje najpierw TAP3 (najpowszechniejszy), potem NRTRDE, a następnie RAP, używając znaczników klasy APPLICATION ASN.1 i wzorców pól do zidentyfikowania formatu.

TAP3 — Transferred Account Procedure

TAP3 to ogólnoświatowy standard wymiany CDR rozliczeń roamingowych między operatorami — format, którego operatorzy używają do wzajemnego rozliczania użycia roamingu. Jest wymieniany w partiach, zazwyczaj codziennie lub cotygodniowo.

  • Standard: GSMA TD.57 / 3GPP TS 32.298. Ciąg wersji TAP3.12.
  • Typy rekordów: 7 — MOC (Mobile Originated Call), MTC (Mobile Terminated Call), GPRS (dane pakietowe), SMS_MO, SMS_MT, SCF (usługi uzupełniające) i VAS (usługa o wartości dodanej).
  • Kodowanie: ASN.1 z 69 stałymi znacznikami kontekstowymi obejmującymi identyfikatory abonentów, strony połączenia, lokalizację, użycie usług podstawowych, znaczniki czasu, informacje o opłatach, informacje podatkowe i pola na poziomie pliku.
  • Szczegóły: kwoty opłat są przenoszone w SDR (tysięcznych) z kursami wymiany i blokami podatkowymi; rekord Mobile Originated Call ma około 33 pól.
  • Mapowanie: MOC mapuje się na ROAMING_OUT, MTC na ROAMING_IN, GPRS na PS_DATA, SMS na typy SMS z komutacją łączy, a SCF/VAS na SUPPLEMENTARY_SERVICE.

NRTRDE — Near Real-Time Roaming Data Exchange

NRTRDE to wymiana użycia roamingu w czasie zbliżonym do rzeczywistego — dostarczana w ciągu godzin, a nie codziennie — zaprojektowana do szybkiego wykrywania nadużyć roamingowych i egzekwowania limitów kredytowych.

  • Standard: GSMA PRD IR.35. Wersja NRTRDE-2.1. Kodowanie ASN.1 BER.
  • Typy rekordów: 5 — VOICE, DATA, SMS, SUPPLEMENTARY i CONTENT.
  • Kodowanie: 32 znaczniki kontekstowe, w tym dedykowane pola nadużyćFRAUD_INDICATOR (0 = normalny, 1 = podejrzewany, 2 = potwierdzony), HIGH_USAGE_INDICATOR, ROAMING_STATUS i CAMEL_INDICATOR.
  • Dlaczego różni się od TAP3: NRTRDE jest celowo uproszczonym formatem — minimalnym, ale wystarczającym zestawem pól zoptymalizowanym pod kątem szybkości wykrywania nadużyć i kontroli kredytowej, a nie pełnego rozliczenia.

RAP — Returned Accounts Procedure

RAP obsługuje spory rozliczeniowe. Gdy operator macierzysty otrzymuje plik TAP i znajduje w nim błędy, generuje rekordy RAP, aby odrzucić konkretne zdarzenia — lub cały plik — i zwrócić je do nadawcy.

  • Standard: GSMA PRD BA.12. Wersja RAP-1.5.
  • Typy rekordów: 4 — CALL_EVENT_REJECTION (pojedyncze nieprawidłowe zdarzenie), FILE_REJECTION (cały plik z błędami krytycznymi), AUDIT_CONTROL_REJECTION (błędy metadanych pliku) i FILE_NOTIFICATION (brakujący lub zduplikowany plik).
  • Kodowanie: 28 znaczników obejmujących odniesienie do oryginalnego pliku TAP, szczegóły odrzucenia, znaczniki czasu, odniesienia do oryginalnych rekordów i kwoty rozliczeniowe.
  • Kody odrzucenia: 20 standardowych kodów błędów GSMA — na przykład UNKNOWN_SUBSCRIBER, MISSING_MANDATORY_FIELD, INVALID_FIELD_VALUE, DUPLICATE_RECORD i CHARGING_ERROR — każdy z czytelnym dla człowieka opisem.
  • Mapowanie: wszystkie rekordy RAP mapują się na ROAMING_IN.

Dekodowanie natywne dla telekomunikacji

TAP3, NRTRDE i RAP nie są generycznymi formatami plików — są to konkretne binarne formaty rekordów roamingowych nakazane przez GSMA dla operatorów mobilnych. Ponieważ konektor CDR dekoduje je bezpośrednio za pomocą ASN.1, Polkomtel może pobierać użycie roamingu, uruchamiać kontrole nadużyć i przetwarzać spory rozliczeniowe bez osobnego, specyficznego dla dostawcy narzędzia dekodującego.

connector: cdr-asn1
config:
  path: gs://polkomtel-cdr-raw/roaming/2026-05-20/
  filePattern: "*.dat"
  templateType: auto
  skipMalformed: true

Podsumowanie możliwości

Poniższa tabela podsumowuje, co potrafi każdy aktywny konektor. CDC = Change Data Capture; Pushdown = przepychanie SQL do źródła.

KonektorOdczytZapisCDCZbiorczyOdkrywanie schematuPushdown
postgresql
mysql
oracle
mssql
teradata
sap-hana
sybase
mongodb
snowflake
databricks
bigquery
kafka
cdc-debezium
csv excel json parquet xml
cdr-asn1
salesforce
sap-erp
servicenow
rest_api

Uwaga

Konektory Azure Event Hubs (azure-eventhubs) i Google Cloud Pub/Sub (google-pubsub) istnieją w kodzie źródłowym, ale są obecnie wyłączone w rejestrze konektorów. Nie są dostępne do użytku, dopóki przenoszenie nie zostanie ukończone.

Poprzednia
Jak DataFlow realizuje ETL