Rozszerzanie i współtworzenie
Connector SDK
Connector SDK to framework, na którym zbudowany jest każdy konektor DataFlow AI. Znajduje się w module biblioteki connector-sdk w backend/platform/connector-sdk/, definiuje pojedynczy kontrakt konektora (ConnectorSDK.kt), dostarcza trzy klasy bazowe dla typowych kształtów konektorów, dostarcza 21 produkcyjnych implementacji konektorów — w tym natywny dla telekomunikacji konektor CDR — i odkrywa je w czasie wykonania poprzez ConnectorRegistry. Ta strona jest referencją modelu konektora oraz przewodnikiem krok po kroku po pisaniu własnego.
Model konektora
Konektor łączy krok potoku danych z systemem zewnętrznym: odkrywa schemat, ekstrahuje dane i ładuje dane. Wszystkie konektory implementują interfejs ConnectorSDK, niezależnie od tego, czy mówią w JDBC, API REST, protokole magazynu chmurowego, binarnym formacie pliku czy protokole strumieniowym.
Układ pakietów w com.polkomtel.dataflow.connector jest następujący:
| Pakiet | Zawartość |
|---|---|
connector (root) | Interfejs ConnectorSDK, ConnectorOperation, ConnectionTestResult |
connector.model | ConnectorConfig, Schema, DataStream, DataRecord, DataBatch, LoadResult, ConnectorCapability, StorageType, ConnectorException, EventStream |
connector.base | JDBCConnectorBase, NativeConnectorBase, FileConnectorBase |
connector.registry | ConnectorRegistry, ConnectorAutoConfiguration |
connector.impl.* | 21 implementacji konektorów, po jednym podpakiecie każda |
Interfejs ConnectorSDK
ConnectorSDK rozszerza AutoCloseable. Większość elementów składowych niesie domyślne implementacje, więc konektor nadpisuje tylko to, co faktycznie wspiera.
interface ConnectorSDK : AutoCloseable {
/** Unique identifier for this connector type (e.g. "postgresql", "csv"). */
val connectorType: String get() = javaClass.simpleName
/** Connector instance identifier — defaults to connectorType. */
fun connectorId(): String = connectorType
/** Human-readable name for the UI. */
fun connectorName(): String = connectorType
val displayName: String get() = connectorType
/** Discover the schema of the data source. */
fun discoverSchema(config: ConnectorConfig): Schema
/** Extract data from the source as a DataStream. */
fun extractData(config: ConnectorConfig): DataStream
/** Load data into the target destination. */
fun loadData(data: DataStream, config: ConnectorConfig): LoadResult
/** Validate the connector can reach the configured source/target. */
fun validateConnection(config: ConnectorConfig): Boolean
/** Capabilities supported by this connector. */
fun capabilities(): Set<ConnectorCapability>
/** Operations supported by this connector. */
fun supportedOperations(): Set<ConnectorOperation>
/** Test the connection and return detailed results. */
fun testConnection(config: ConnectorConfig): ConnectionTestResult
/** CDC event stream — throws UnsupportedOperationException unless overridden. */
fun getCDCEvents(config: ConnectorConfig): EventStream
/** Push-down SQL — throws UnsupportedOperationException unless overridden. */
fun pushDownSQL(sql: String, dialect: String, config: ConnectorConfig): ResultSet
/** Release resources. */
override fun close()
}
Dwie metody, które nie są implementowane domyślnie — discoverSchema i extractData — są minimum, jakie konektor musi dostarczyć. loadData domyślnie zapisuje zero rekordów, getCDCEvents i pushDownSQL domyślnie rzucają UnsupportedOperationException, a testConnection domyślnie mierzy czas wywołania validateConnection.
Kontrakt bogaty w wartości domyślne
Interfejs jest celowo bogaty w wartości domyślne. Konektor pliku tylko do odczytu może zaimplementować jedynie discoverSchema i extractData; framework dostarcza rozsądne zachowanie no-op lub rzucające wyjątek dla wszystkiego pozostałego. Zadeklaruj, co wspierasz, poprzez capabilities() i supportedOperations(), aby platforma nigdy nie wywołała metody, której nie zaimplementowałeś.
Odpowiedzialności metod
| Metoda | Odpowiedzialność |
|---|---|
discoverSchema | Inspekcjonuje źródło i zwraca Schema pól, typów i metadanych |
extractData | Zwraca leniwy, zamykalny DataStream rekordów DataRecord odczytanych ze źródła |
loadData | Zapisuje DataStream (lub Iterable<DataRecord>) do celu, zwracając LoadResult |
validateConnection | Tania kontrola osiągalności; domyślnie podpiera testConnection |
testConnection | Szczegółowa sonda łączności — podpiera POST /connections/{id}/test |
getCDCEvents | Zwraca EventStream zdarzeń zmian (tylko konektory CDC) |
pushDownSQL | Wykonuje SQL specyficzny dla dialektu u źródła w celu optymalizacji przepychania |
close | Zwalnia połączenia z puli, klientów HTTP i inne zasoby |
Możliwości i operacje
Konektor ogłasza, co potrafi, poprzez dwa enumy. Platforma konsultuje je przed skierowaniem pracy do konektora.
ConnectorCapability
enum class ConnectorCapability {
SCHEMA_DISCOVERY, // Can discover schema from the source
READ, // Can extract/read data
WRITE, // Can load/write data
COLUMN_PROJECTION, // Can read a subset of columns
PREDICATE_PUSHDOWN, // Can filter at the source
SCHEMA_EVOLUTION, // Supports schema changes
STREAMING, // Supports streaming/incremental reads
PARTITIONING, // Supports partitioned data
COMPRESSION, // Supports compression
CLOUD_STORAGE, // Supports reading from cloud storage
BATCH_WRITE, // Supports batch writing
ENCODING_DETECTION, // Supports encoding detection
MULTI_TABLE // Multiple sheets/tables in one source
}
ConnectorOperation
enum class ConnectorOperation {
EXTRACT,
LOAD,
PUSHDOWN_SQL,
SCHEMA_DISCOVERY,
CDC,
STREAMING,
BULK_COPY,
HEARTBEAT
}
Możliwość opisuje, jakie kształty danych konektor obsługuje; operacja opisuje, jakie czasowniki wspiera. Dwie klasy bazowe wstępnie deklarują rozsądny zestaw operacji: JDBCConnectorBase deklaruje EXTRACT, LOAD, PUSHDOWN_SQL i SCHEMA_DISCOVERY; NativeConnectorBase deklaruje EXTRACT, LOAD i SCHEMA_DISCOVERY.
Schemat konfiguracji
Każda metoda konektora przyjmuje ConnectorConfig — pojedynczą klasę danych niosącą szczegóły połączenia, opcje formatu oraz parametry ekstrakcji/ładowania.
data class ConnectorConfig(
val path: String = "", // file path, gs:// / s3:// / abfss:// URI
val properties: Map<String, String> = emptyMap(), // format-specific options
val columnProjection: List<String>? = null,
val filterPredicate: String? = null,
val maxRecords: Long = -1, // -1 = unlimited
val batchSize: Int = 10_000,
val schemaSampleSize: Int = 100,
val storageType: StorageType? = null,
// ---- JDBC / database properties ----
val host: String = "",
val port: Int = 0,
val database: String = "",
val username: String = "",
val password: String = "",
val schema: String = "",
val query: String = "",
val table: String = "",
val fetchSize: Int = 1000,
val timeoutSeconds: Int = 30,
val ssl: Boolean = false,
val connectionId: String = ""
)
Pola konfiguracji
| Pole | Domyślnie | Cel |
|---|---|---|
path | "" | Ścieżka lub URI dla konektorów plików / magazynu chmurowego |
properties | pusta mapa | Dowolne opcje specyficzne dla konektora |
columnProjection | null | Projektuje podzbiór kolumn, gdy wspierane |
filterPredicate | null | Filtr po stronie źródła dla przepychania predykatów |
maxRecords | -1 | Ogranicza liczbę ekstrahowanych wierszy; -1 oznacza bez limitu |
batchSize | 10000 | Rozmiar wsadu zapisu |
schemaSampleSize | 100 | Wiersze próbkowane do wnioskowania schematu |
host / port / database | "" / 0 / "" | Współrzędne połączenia JDBC |
username / password | "" | Poświadczenia połączenia |
schema / table / query | "" | Co odczytać lub zapisać |
fetchSize | 1000 | Rozmiar pobierania kursora JDBC |
timeoutSeconds | 30 | Limit czasu połączenia |
Akcesory pomocnicze
ConnectorConfig dostarcza typowane akcesory, aby konektory unikały parsowania ciągów znaków:
config.propertyString("delimiter", ",") // String with default
config.propertyInt("maxRows", 1000) // Int with default
config.propertyBoolean("skipMalformed", true)
config.jdbcUrl("jdbc:postgresql://") // prefix + host:port/database
config.effectiveQuery() // configured query or SELECT * FROM table
config.qualifiedTableName() // schema.table when schema is set
LoadResult
loadData zwraca LoadResult opisujący wynik:
data class LoadResult(
val recordsWritten: Long = 0,
val recordsFailed: Long = 0,
val bytesWritten: Long = 0,
val duration: Duration = Duration.ZERO,
val outputPaths: List<String> = emptyList(),
val warnings: List<String> = emptyList(),
val success: Boolean = true,
val rowsRejected: Long = 0,
val errors: List<String> = emptyList(),
val targetTable: String = ""
)
ConnectionTestResult, zwracany przez testConnection, jest mniejszy:
data class ConnectionTestResult(
val success: Boolean,
val latencyMs: Long = 0,
val serverVersion: String? = null,
val message: String? = null
)
Klasy bazowe
Zamiast implementować ConnectorSDK od zera, rozszerz jedną z trzech klas bazowych, które pasują do typowych kształtów konektorów.
| Klasa bazowa | Rozszerz, gdy źródło jest… | Dostarcza |
|---|---|---|
JDBCConnectorBase | Dowolną bazą danych SQL osiągalną przez JDBC | Pulę HikariCP, odkrywanie schematu sterowane metadanymi, ładowanie przez wstawianie wsadowe, przepychanie SQL |
NativeConnectorBase | API REST, kolejką komunikatów lub innym systemem nie-JDBC | Haki cyklu życia connect / disconnect / isConnected |
FileConnectorBase | Formatem pliku (CSV, JSON, XML, Excel, Parquet, CDR) | Abstrakcję magazynu nad lokalnym systemem plików, GCS, S3 i Azure Blob |
JDBCConnectorBase
Konektor JDBC dostarcza tylko trzy rzeczy — klasę sterownika, prefiks URL i domyślny port — i dziedziczy pełny konektor:
abstract class JDBCConnectorBase : ConnectorSDK {
abstract fun driverClassName(): String
abstract fun jdbcPrefix(): String
abstract fun defaultPort(): Int
// discoverSchema, extractData, loadData, pushDownSQL, testConnection are inherited
}
Buduje on HikariDataSource z maximumPoolSize = 10, minimumIdle = 2, idleTimeout = 600000 i maxLifetime = 1800000, kopiując każdy wpis z config.properties jako właściwość źródła danych. discoverSchema przechodzi DatabaseMetaData w poszukiwaniu obiektów TABLE i VIEW, mapując java.sql.Types na enum DataType DataFlow. extractData otwiera strumieniujący ResultSet ze skonfigurowanym fetchSize; loadData uruchamia wsadowe instrukcje INSERT wewnątrz transakcji z wycofaniem przy niepowodzeniu.
NativeConnectorBase
Konektory natywne implementują jawny cykl życia połączenia plus metody pracy do*:
abstract class NativeConnectorBase : ConnectorSDK {
abstract fun connect(config: ConnectorConfig)
abstract fun disconnect()
abstract fun isConnected(): Boolean
abstract fun doDiscoverSchema(config: ConnectorConfig): Schema
abstract fun doExtractData(config: ConnectorConfig): DataStream
abstract fun doLoadData(data: DataStream, config: ConnectorConfig): LoadResult
}
Klasa bazowa opakowuje każdą metodę publiczną, tak aby connect/disconnect zawsze obejmowały pracę, a testConnection mierzy opóźnienie connect.
FileConnectorBase
Konektory plików implementują specyficzne dla formatu inferSchema, readRecords i writeRecords. Klasa bazowa dostarcza StorageAdapter, więc ten sam konektor odczytuje ze ścieżki lokalnej, URI gs://, s3:// lub abfss:// w sposób przezroczysty.
21 wbudowanych konektorów
SDK dostarcza 21 implementacji konektorów w connector.impl.*. Większość jest automatycznie odkrywana przez META-INF/services; DebeziumCdcConnector jest rejestrowany poprzez Spring DI, ponieważ potrzebuje wstrzykiwania zależności.
| Kategoria | Konektory |
|---|---|
| Bazy danych (JDBC) | PostgreSQL, MySQL, MSSQL, Oracle, Sybase, Teradata, SAP HANA |
| Chmurowe hurtownie | BigQuery, Databricks, Snowflake |
| NoSQL | MongoDB |
| Formaty plików | CSV, JSON, XML, Excel, Parquet |
| Streaming i komunikaty | Kafka, Pub/Sub, Event Hubs |
| SaaS / ERP | Salesforce, ServiceNow, SAP, REST API |
| Telekomunikacja | CDR (ASN.1) |
Przekrojowe wsparcie CDC (change-data-capture) jest zapewniane za pośrednictwem Debezium dla binlogu MySQL, śledzenia zmian MSSQL i strumieni zmian MongoDB.
EventHubs i Pub/Sub
Klasy konektorów eventhubs i pubsub istnieją w connector.impl.*, ale są tymczasowo zakomentowane w META-INF/services, dopóki nie zakończy się zadanie przenoszenia — rejestrują się ponownie automatycznie, gdy ta praca zostanie ukończona.
Telekomunikacyjny konektor CDR
Podpakiet cdr/ zawiera konektor zbudowany specjalnie dla telekomunikacyjnych rekordów szczegółów połączeń Polkomtel. CdrConnector rozszerza FileConnectorBase, ma connectorType = "cdr-asn1", displayName = "Telecom CDR (ASN.1)" i deklaruje możliwości SCHEMA_DISCOVERY, READ, STREAMING i CLOUD_STORAGE. Jest tylko do odczytu — writeRecords rzuca wyjątek, kierując wywołujących do CdrParquetConverter w celu uzyskania wyjścia Parquet.
Potok dekodowania jest następujący:
Binary CDR file
→ StorageAdapter.openForRead() (local / GCS / S3 / Azure)
→ Asn1Decoder.decodeStream() (BER/DER TLV parsing)
→ CdrBinaryDecoder.decodeStream() (template-driven field extraction)
→ CdrConnector.cdrRecordToDataRecord
→ DataStream (lazy, closeable sequence)
Dostarcza trzy szablony rekordów roamingowych, walidowane przez RoamingTemplateValidator:
| Szablon | Format telekomunikacyjny |
|---|---|
Tap3Template | TAP3 — Transferred Account Procedure, standard rozliczania użycia roamingu między operatorami |
RapTemplate | RAP — Returned Account Procedure, używany do zwracania odrzuconych rekordów TAP |
NrtrdeTemplate | NRTRDE — Near Real Time Roaming Data Exchange, do szybkiego wykrywania nadużyć w użyciu roamingu |
Konfiguracja konektora CDR jest dostarczana poprzez properties:
| Właściwość | Domyślnie | Opis |
|---|---|---|
templateType | auto | Typ szablonu CDR (np. CS_VOICE_MO lub auto do wykrycia) |
templateVersion | 3GPP-R15 | Wersja szablonu |
batchSize | 10000 | Rozmiar wsadu przetwarzania |
skipMalformed | true | Pomiń zniekształcone rekordy zamiast zawieść |
filePattern | *.cdr | Glob plików dla skanowania katalogu |
sourceSwitch | (wyprowadzony) | Identyfikator przełącznika źródłowego; w przeciwnym razie wyprowadzony z nazwy pliku |
Wspierane rozszerzenia plików to .cdr, .asn1, .dat i .bin. Gdy templateType=auto, konektor odczytuje pierwszą strukturę ASN.1, dopasowuje ją do znanych szablonów i resetuje strumień, tak aby żaden rekord nie został skonsumowany.
Dekodowanie natywne dla telekomunikacji
TAP3, RAP i NRTRDE nie są generycznymi formatami plików — są to konkretne binarne formaty rekordów roamingowych używane przez operatorów mobilnych. Konektor CDR dekoduje je bezpośrednio za pomocą ASN.1, więc Polkomtel może pobierać użycie roamingu i uruchamiać kontrole nadużyć bez osobnego etapu dekodowania.
Cykl życia konektora
Instancja konektora podąża za tym samym cyklem życia za każdym razem, gdy używa jej silnik potoków danych.
| Faza | Co się dzieje |
|---|---|
| Odkrywanie | Przy starcie ConnectorAutoConfiguration scala ścieżki odkrywania ServiceLoader i Spring-DI w ConnectorRegistry |
| Wyszukiwanie | Silnik wywołuje registry.getConnector(connectorId); rejestr wywołuje przechowywaną fabrykę, aby wytworzyć instancję |
| Test | testConnection(config) uruchamia się, gdy użytkownik kliknie Test w kreatorze połączeń lub wywoła POST /connections/{id}/test |
| Schemat | discoverSchema(config) jest wywoływane podczas projektowania potoku danych w celu wypełnienia list pól |
| Uruchomienie | extractData (źródła) lub loadData (cele) jest wywoływane podczas wykonywania potoku danych |
| Zamknięcie | close() zwalnia pule, klientów HTTP i inne zasoby |
DataStream zwracany przez extractData jest leniwy i niesie własne obiekty Closeable — silnik opróżnia go, a następnie zamyka, co z kolei zamyka leżący u podstaw ResultSet, uchwyt pliku lub odpowiedź HTTP.
Budowanie własnego konektora
Silnik potoków danych wspiera podłączalne konektory. Poniższy przykład dodaje generyczny konektor API REST poprzez rozszerzenie NativeConnectorBase.
Krok 1 — Utwórz klasę konektora
package com.polkomtel.dataflow.connector.impl.myrest
import com.polkomtel.dataflow.connector.base.NativeConnectorBase
import com.polkomtel.dataflow.connector.model.*
class MyRestApiConnector : NativeConnectorBase() {
override val connectorType = "my-rest-api"
override val displayName = "My REST API"
private var client: java.net.http.HttpClient? = null
private var baseUrl: String = ""
override fun connect(config: ConnectorConfig) {
baseUrl = config.host.ifBlank {
throw ConnectorException("Base URL (host) is required")
}
client = java.net.http.HttpClient.newBuilder()
.connectTimeout(java.time.Duration.ofSeconds(config.timeoutSeconds.toLong()))
.build()
}
override fun disconnect() { client = null }
override fun isConnected(): Boolean = client != null
override fun capabilities(): Set<ConnectorCapability> = setOf(
ConnectorCapability.SCHEMA_DISCOVERY,
ConnectorCapability.READ,
ConnectorCapability.WRITE
)
override fun doDiscoverSchema(config: ConnectorConfig): Schema {
// Sample the endpoint response and infer fields
return Schema(fields = emptyList())
}
override fun doExtractData(config: ConnectorConfig): DataStream {
// GET config.table (the endpoint path) and map JSON to DataRecords
return DataStream.empty(Schema(fields = emptyList()))
}
override fun doLoadData(data: DataStream, config: ConnectorConfig): LoadResult {
// POST records to the endpoint
return LoadResult(recordsWritten = 0)
}
}
Krok 2 — Zarejestruj konektor
Dla zwykłego konektora z publicznym konstruktorem bezargumentowym dodaj jego w pełni kwalifikowaną nazwę klasy do manifestu ServiceLoader:
# backend/platform/connector-sdk/src/main/resources/
# META-INF/services/com.polkomtel.dataflow.connector.ConnectorSDK
com.polkomtel.dataflow.connector.impl.myrest.MyRestApiConnector
ConnectorRegistry.discoverConnectors() podejmuje go przy starcie za pomocą java.util.ServiceLoader.
Jeśli konektor potrzebuje wstrzykiwania zależności, zamiast tego oznacz go adnotacją @Component — ConnectorAutoConfiguration rejestruje każdy bean ConnectorSDK zarządzany przez Spring, używając samego beana jako wyniku fabryki. Tę ścieżkę wykorzystuje DebeziumCdcConnector.
Krok 3 — Rejestracja programowa
Konektor możesz również zarejestrować ręcznie względem rejestru singletonowego:
val registry = ConnectorRegistry.getInstance()
registry.register("my-rest-api") { MyRestApiConnector() }
// Use it
val connector = registry.getConnector("my-rest-api")
val schema = connector.discoverSchema(config)
Rejestr przechowuje fabrykę (() -> ConnectorSDK), więc każde wywołanie getConnector daje świeżą instancję — co jest ważne, ponieważ konektory przechowują stan dla każdego połączenia, taki jak pula Hikari.
Testowanie konektora
Testy konektorów znajdują się obok ich implementacji i uruchamiają się ze standardowym zadaniem testowym Gradle.
class MyRestApiConnectorTest {
@Test
fun `test connection succeeds with a valid base URL`() {
val connector = MyRestApiConnector()
val result = connector.testConnection(
ConnectorConfig(host = "https://jsonplaceholder.typicode.com")
)
assertTrue(result.success)
}
@Test
fun `capabilities advertise read and write`() {
val caps = MyRestApiConnector().capabilities()
assertTrue(ConnectorCapability.READ in caps)
assertTrue(ConnectorCapability.WRITE in caps)
}
}
Uruchom testy modułu SDK:
./gradlew :platform:connector-sdk:test
| Cel testowy | Polecenie |
|---|---|
| Testy jednostkowe Connector SDK | ./gradlew :platform:connector-sdk:test |
| Wszystkie testy jednostkowe platformy | ./gradlew test |
| Testy integracyjne (wymagany Docker) | ./gradlew integrationTest |
Testuj względem prawdziwych systemów
Test konektora, który mockuje sterownik JDBC lub klienta HTTP, niczego nie dowodzi w kwestii łączności. Skieruj testy testConnection i discoverSchema na prawdziwy kontener bazy danych, publiczne API piaskownicy lub lokalny plik testowy, aby test ćwiczył rzeczywistą ścieżkę dekodowania i połączenia.
Pakowanie
SDK jest modułem biblioteki Gradle. Konektor dodany do connector-sdk jest dostarczany jako część pliku JAR tego modułu; nie jest potrzebne osobne pakowanie wtyczki.
| Zagadnienie | Jak jest obsługiwane |
|---|---|
| Build | ./gradlew :platform:connector-sdk:build wytwarza plik JAR modułu |
| Sterowniki JDBC | Zadeklarowane jako runtimeOnly, więc są na ścieżce klas runtime, ale nie na ścieżce klas kompilacji |
| Wersje zależności | Przypięte w katalogu wersji Gradle gradle/libs.versions.toml |
| Pula połączeń | HikariCP 6.2.1, współdzielony przez wszystkie konektory JDBC poprzez JDBCConnectorBase |
| Manifest odkrywania | META-INF/services/com.polkomtel.dataflow.connector.ConnectorSDK |
Ponieważ sterowniki JDBC są zadeklarowane jako runtimeOnly, dodanie konektora bazy danych oznacza dodanie jednego wpisu katalogu plus jednej linii runtimeOnly(libs....) — sterownik jest rozwiązywany w czasie wykonania i nigdy nie przedostaje się do API czasu kompilacji. Pełny zestaw współrzędnych sterowników, łańcuchów połączeń i strojenia HikariCP znajduje się w Referencji Connector SDK.
Podsumowanie referencyjne
| Temat | Gdzie się znajduje |
|---|---|
| Kontrakt konektora | connector/ConnectorSDK.kt |
| Możliwości / operacje | connector/model/ConnectorCapability.kt, connector/ConnectorOperation.kt |
| Konfiguracja | connector/model/ConnectorConfig.kt (ConnectorConfig, LoadResult) |
| Klasy bazowe | connector/base/ (JDBCConnectorBase, NativeConnectorBase, FileConnectorBase) |
| Rejestr | connector/registry/ (ConnectorRegistry, ConnectorAutoConfiguration) |
| Wbudowane konektory | connector/impl/* — 21 implementacji |
| Telekomunikacyjny CDR | connector/impl/cdr/ — CdrConnector, Asn1Decoder, szablony TAP3/RAP/NRTRDE |
| Manifest odkrywania | META-INF/services/com.polkomtel.dataflow.connector.ConnectorSDK |