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:

PakietZawartość
connector (root)Interfejs ConnectorSDK, ConnectorOperation, ConnectionTestResult
connector.modelConnectorConfig, Schema, DataStream, DataRecord, DataBatch, LoadResult, ConnectorCapability, StorageType, ConnectorException, EventStream
connector.baseJDBCConnectorBase, NativeConnectorBase, FileConnectorBase
connector.registryConnectorRegistry, 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

MetodaOdpowiedzialność
discoverSchemaInspekcjonuje źródło i zwraca Schema pól, typów i metadanych
extractDataZwraca leniwy, zamykalny DataStream rekordów DataRecord odczytanych ze źródła
loadDataZapisuje DataStream (lub Iterable<DataRecord>) do celu, zwracając LoadResult
validateConnectionTania kontrola osiągalności; domyślnie podpiera testConnection
testConnectionSzczegółowa sonda łączności — podpiera POST /connections/{id}/test
getCDCEventsZwraca EventStream zdarzeń zmian (tylko konektory CDC)
pushDownSQLWykonuje SQL specyficzny dla dialektu u źródła w celu optymalizacji przepychania
closeZwalnia 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

PoleDomyślnieCel
path""Ścieżka lub URI dla konektorów plików / magazynu chmurowego
propertiespusta mapaDowolne opcje specyficzne dla konektora
columnProjectionnullProjektuje podzbiór kolumn, gdy wspierane
filterPredicatenullFiltr po stronie źródła dla przepychania predykatów
maxRecords-1Ogranicza liczbę ekstrahowanych wierszy; -1 oznacza bez limitu
batchSize10000Rozmiar wsadu zapisu
schemaSampleSize100Wiersze 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ć
fetchSize1000Rozmiar pobierania kursora JDBC
timeoutSeconds30Limit 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 bazowaRozszerz, gdy źródło jest…Dostarcza
JDBCConnectorBaseDowolną bazą danych SQL osiągalną przez JDBCPulę HikariCP, odkrywanie schematu sterowane metadanymi, ładowanie przez wstawianie wsadowe, przepychanie SQL
NativeConnectorBaseAPI REST, kolejką komunikatów lub innym systemem nie-JDBCHaki cyklu życia connect / disconnect / isConnected
FileConnectorBaseFormatem 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.

KategoriaKonektory
Bazy danych (JDBC)PostgreSQL, MySQL, MSSQL, Oracle, Sybase, Teradata, SAP HANA
Chmurowe hurtownieBigQuery, Databricks, Snowflake
NoSQLMongoDB
Formaty plikówCSV, JSON, XML, Excel, Parquet
Streaming i komunikatyKafka, Pub/Sub, Event Hubs
SaaS / ERPSalesforce, ServiceNow, SAP, REST API
TelekomunikacjaCDR (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:

SzablonFormat telekomunikacyjny
Tap3TemplateTAP3 — Transferred Account Procedure, standard rozliczania użycia roamingu między operatorami
RapTemplateRAP — Returned Account Procedure, używany do zwracania odrzuconych rekordów TAP
NrtrdeTemplateNRTRDE — 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ślnieOpis
templateTypeautoTyp szablonu CDR (np. CS_VOICE_MO lub auto do wykrycia)
templateVersion3GPP-R15Wersja szablonu
batchSize10000Rozmiar wsadu przetwarzania
skipMalformedtruePomiń zniekształcone rekordy zamiast zawieść
filePattern*.cdrGlob 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.

FazaCo się dzieje
OdkrywaniePrzy starcie ConnectorAutoConfiguration scala ścieżki odkrywania ServiceLoader i Spring-DI w ConnectorRegistry
WyszukiwanieSilnik wywołuje registry.getConnector(connectorId); rejestr wywołuje przechowywaną fabrykę, aby wytworzyć instancję
TesttestConnection(config) uruchamia się, gdy użytkownik kliknie Test w kreatorze połączeń lub wywoła POST /connections/{id}/test
SchematdiscoverSchema(config) jest wywoływane podczas projektowania potoku danych w celu wypełnienia list pól
UruchomienieextractData (źródła) lub loadData (cele) jest wywoływane podczas wykonywania potoku danych
Zamknięcieclose() 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ą @ComponentConnectorAutoConfiguration 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 testowyPolecenie
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.

ZagadnienieJak jest obsługiwane
Build./gradlew :platform:connector-sdk:build wytwarza plik JAR modułu
Sterowniki JDBCZadeklarowane jako runtimeOnly, więc są na ścieżce klas runtime, ale nie na ścieżce klas kompilacji
Wersje zależnościPrzypię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 odkrywaniaMETA-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

TematGdzie się znajduje
Kontrakt konektoraconnector/ConnectorSDK.kt
Możliwości / operacjeconnector/model/ConnectorCapability.kt, connector/ConnectorOperation.kt
Konfiguracjaconnector/model/ConnectorConfig.kt (ConnectorConfig, LoadResult)
Klasy bazoweconnector/base/ (JDBCConnectorBase, NativeConnectorBase, FileConnectorBase)
Rejestrconnector/registry/ (ConnectorRegistry, ConnectorAutoConfiguration)
Wbudowane konektoryconnector/impl/* — 21 implementacji
Telekomunikacyjny CDRconnector/impl/cdr/CdrConnector, Asn1Decoder, szablony TAP3/RAP/NRTRDE
Manifest odkrywaniaMETA-INF/services/com.polkomtel.dataflow.connector.ConnectorSDK
Poprzednia
Playbooki i runbooki