Kafka Connect 커넥터 개발 가이드
Kafka Connect 커넥터 개발 가이드 (Connector Development Guide)
이 페이지는 Kafka와 다른 시스템 사이에서 데이터를 옮기는 새 커넥터를 직접 개발할 때 필요한 핵심 개념과 코드 작성법을 다뤄요. 간단한 파일 소스 커넥터 예제를 따라가며 SourceConnector/SourceTask/SinkTask 구현, 오프셋 재개, 정확히 한 번 지원 같은 부분을 익힐 수 있어요.
출처: 문서
본문
이 가이드는 개발자가 Kafka와 다른 시스템 사이에서 데이터를 옮기기 위해 Kafka Connect용 새 커넥터를 작성하는 방법을 설명합니다. 몇 가지 핵심 개념을 간단히 검토한 다음, 간단한 커넥터를 만드는 방법을 설명합니다.
핵심 개념과 API (Core Concepts and APIs)
커넥터와 태스크 (Connectors and Tasks)
Kafka와 다른 시스템 사이에서 데이터를 복사하려면 사용자는 데이터를 가져오거나 보낼 시스템에 대한 커넥터(Connector)를 생성합니다. 커넥터에는 두 가지 종류가 있습니다. SourceConnector는 다른 시스템에서 데이터를 가져오고(예: JDBCSourceConnector는 관계형 데이터베이스를 Kafka로 가져옴), SinkConnector는 데이터를 내보냅니다(예: HDFSSinkConnector는 Kafka 토픽의 내용을 HDFS 파일로 내보냄).
커넥터 자체는 데이터 복사를 수행하지 않습니다. 커넥터의 구성이 복사할 데이터를 설명하고, 커넥터는 그 작업을 워커에 분산될 수 있는 태스크(Task) 집합으로 나누는 일을 담당합니다. 이 태스크에도 두 가지 대응하는 종류가 있습니다. SourceTask와 SinkTask입니다.
할당(assignment)을 받으면 각 태스크는 데이터의 자신의 부분집합을 Kafka로/에서 복사해야 합니다. Kafka Connect에서 이러한 할당은 항상 일관된 스키마를 가진 레코드로 구성된 입·출력 스트림 집합으로 표현할 수 있어야 합니다. 때로는 이 매핑이 명확합니다. 로그 파일 집합의 각 파일은 스트림으로 간주할 수 있으며, 각 파싱된 줄이 동일한 스키마를 사용하는 레코드를 형성하고 오프셋은 파일의 바이트 오프셋으로 저장됩니다. 다른 경우에는 이 모델에 매핑하기 위해 더 많은 노력이 필요할 수 있습니다. JDBC 커넥터는 각 테이블을 스트림으로 매핑할 수 있지만 오프셋은 덜 분명합니다. 한 가지 가능한 매핑은 타임스탬프 열을 사용해 새 데이터를 점진적으로 반환하는 쿼리를 생성하고, 마지막으로 쿼리된 타임스탬프를 오프셋으로 사용하는 것입니다.
스트림과 레코드 (Streams and Records)
각 스트림은 키-값 레코드의 시퀀스여야 합니다. 키와 값 모두 복잡한 구조를 가질 수 있습니다. 많은 기본 타입이 제공되지만, 배열, 객체, 중첩 데이터 구조도 표현할 수 있습니다. 런타임 데이터 형식은 특정 직렬화 형식을 가정하지 않습니다. 이 변환은 프레임워크가 내부적으로 처리합니다.
키와 값 외에도 (소스가 생성하고 싱크에 전달되는) 레코드에는 연관된 스트림 ID와 오프셋이 있습니다. 이들은 프레임워크가 처리된 데이터의 오프셋을 주기적으로 커밋하는 데 사용되므로, 장애 발생 시 마지막으로 커밋된 오프셋부터 처리를 재개해 불필요한 재처리와 이벤트 중복을 피할 수 있습니다.
동적 커넥터 (Dynamic Connectors)
모든 작업이 정적인 것은 아니므로 커넥터 구현은 재구성이 필요할 수 있는 외부 시스템의 변경을 모니터링할 책임도 있습니다. 예를 들어 JDBCSourceConnector 예제에서 커넥터는 각 태스크에 테이블 집합을 할당할 수 있습니다. 새 테이블이 생성되면 커넥터는 이를 발견해 구성을 업데이트함으로써 새 테이블을 태스크 중 하나에 할당해야 합니다. 재구성이 필요한 변경(또는 태스크 수의 변경)을 감지하면 프레임워크에 알리고, 프레임워크는 해당 태스크들을 업데이트합니다.
간단한 커넥터 개발 (Developing a Simple Connector)
커넥터를 개발하려면 Connector와 Task 두 인터페이스만 구현하면 됩니다. 간단한 예제가 Kafka 소스 코드의 파일 패키지에 포함되어 있습니다. 이 커넥터는 standalone 모드에서 사용하기 위한 것으로, 파일의 각 줄을 읽어 레코드로 내보내는 SourceConnector/SourceTask 구현과 각 레코드를 파일에 쓰는 SinkConnector/SinkTask 구현을 담고 있습니다.
이 섹션의 나머지 부분에서는 커넥터 생성의 핵심 단계를 보여주는 일부 코드를 살펴볼 것입니다. 하지만 간결함을 위해 많은 세부 사항이 생략되었으므로 개발자는 전체 예제 소스 코드도 참조해야 합니다.
커넥터 예제 (Connector Example)
SourceConnector를 간단한 예제로 다루겠습니다. SinkConnector 구현도 매우 유사합니다. 패키지와 클래스 이름을 선택하세요. 이 예제는 FileStreamSourceConnector를 사용하지만, 해당 위치에 자신의 클래스 이름을 대체하세요. 런타임에 플러그인을 발견 가능하게 만들기 위해 META-INF/services/org.apache.kafka.connect.source.SourceConnector의 리소스에 정규화된 클래스 이름을 한 줄에 담은 ServiceLoader 매니페스트를 추가하세요.
com.example.FileStreamSourceConnector
SourceConnector를 상속하는 클래스를 만들고, 태스크에 전달될 구성 정보(데이터를 보낼 토픽, 선택적으로 읽을 파일 이름과 최대 배치 크기)를 저장할 필드를 추가하세요.
package com.example;
public class FileStreamSourceConnector extends SourceConnector {
private Map<String, String> props;
가장 채우기 쉬운 메서드는 taskClass()로, 실제로 데이터를 읽기 위해 워커 프로세스에서 인스턴스화해야 하는 클래스를 정의합니다.
@Override
public Class<? extends Task> taskClass() {
return FileStreamSourceTask.class;
}
FileStreamSourceTask 클래스는 아래에서 정의하겠습니다. 다음으로 표준 라이프사이클 메서드인 start()와 stop()을 추가합니다.
@Override
public void start(Map<String, String> props) {
// Initialization logic and setting up of resources can take place in this method.
// This connector doesn't need to do any of that, but we do log a helpful message to the user.
this.props = props;
AbstractConfig config = new AbstractConfig(CONFIG_DEF, props);
String filename = config.getString(FILE_CONFIG);
filename = (filename == null || filename.isEmpty()) ? "standard input" : config.getString(FILE_CONFIG);
log.info("Starting file source connector reading from {}", filename);
}
@Override
public void stop() {
// Nothing to do since no background monitoring is required.
}
마지막으로 구현의 실제 핵심은 taskConfigs()에 있습니다. 여기서는 단일 파일만 처리하므로 maxTasks 인자에 따라 더 많은 태스크를 생성할 수 있더라도, 항목이 하나뿐인 목록을 반환합니다.
@Override
public List<Map<String, String>> taskConfigs(int maxTasks) {
// Note that the task configs could contain configs additional to or different from the connector configs if needed. For instance,
// if different tasks have different responsibilities, or if different tasks are meant to process different subsets of the source data stream).
ArrayList<Map<String, String>> configs = new ArrayList<>();
// Only one input stream makes sense.
configs.add(props);
return configs;
}
여러 태스크가 있어도 이 메서드 구현은 보통 꽤 간단합니다. 입력 태스크 수를 결정해야 하며, 이는 데이터를 가져오는 원격 서비스에 연락해야 할 수도 있고, 그다음 작업을 나눕니다. 태스크 간 작업 분할 패턴 중 일부는 매우 흔하므로 ConnectorUtils에 이를 단순화하는 유틸리티가 제공됩니다.
이 간단한 예제에는 동적 입력이 포함되지 않는다는 점에 유의하세요. 태스크 구성 업데이트를 트리거하는 방법은 다음 섹션의 설명을 참고하세요.
태스크 예제 - 소스 태스크 (Task Example - Source Task)
다음으로 대응하는 SourceTask의 구현을 설명하겠습니다. 구현은 짧지만 이 가이드에서 완전히 다루기에는 너무 깁니다. 구현의 대부분을 설명하는 데 의사코드를 사용하겠지만, 전체 예제는 소스 코드를 참조할 수 있습니다.
커넥터와 마찬가지로 적절한 기본 Task 클래스를 상속하는 클래스를 만들어야 합니다. 여기에도 표준 라이프사이클 메서드가 있습니다.
public class FileStreamSourceTask extends SourceTask {
private String filename;
private InputStream stream;
private String topic;
private int batchSize;
@Override
public void start(Map<String, String> props) {
filename = props.get(FileStreamSourceConnector.FILE_CONFIG);
stream = openOrThrowError(filename);
topic = props.get(FileStreamSourceConnector.TOPIC_CONFIG);
batchSize = props.get(FileStreamSourceConnector.TASK_BATCH_SIZE_CONFIG);
}
@Override
public synchronized void stop() {
stream.close();
}
}
이들은 약간 단순화된 버전이지만, 이 메서드들이 비교적 간단해야 하며 수행해야 하는 유일한 작업은 리소스 할당 또는 해제라는 것을 보여줍니다. 이 구현에 대해 두 가지 주의할 점이 있습니다. 첫째, start() 메서드는 아직 이전 오프셋에서 재개하는 것을 처리하지 않으며, 이는 이후 섹션에서 다룹니다. 둘째, stop() 메서드는 synchronized입니다. 이는 SourceTask에 전용 스레드가 주어져 무기한 차단할 수 있으므로, Worker의 다른 스레드에서 호출하여 중지해야 하기 때문에 필요합니다.
다음으로 태스크의 주요 기능인 poll() 메서드를 구현합니다. 이 메서드는 입력 시스템에서 이벤트를 가져와 List<SourceRecord>를 반환합니다.
@Override
public List<SourceRecord> poll() throws InterruptedException {
try {
ArrayList<SourceRecord> records = new ArrayList<>();
while (streamValid(stream) && records.isEmpty()) {
LineAndOffset line = readToNextLine(stream);
if (line != null) {
Map<String, Object> sourcePartition = Collections.singletonMap("filename", filename);
Map<String, Object> sourceOffset = Collections.singletonMap("position", streamOffset);
records.add(new SourceRecord(sourcePartition, sourceOffset, topic, Schema.STRING_SCHEMA, line));
if (records.size() >= batchSize) {
return records;
}
} else {
Thread.sleep(1);
}
}
return records;
} catch (IOException e) {
// Underlying stream was killed, probably as a result of calling stop. Allow to return
// null, and driving thread will handle any shutdown if necessary.
}
return null;
}
다시 말하지만 몇 가지 세부 사항을 생략했지만, 중요한 단계를 볼 수 있습니다. poll() 메서드는 반복적으로 호출될 것이며, 각 호출에서 파일에서 레코드를 읽으려고 루프를 돕니다. 읽는 각 줄에 대해 파일 오프셋도 추적합니다. 이 정보를 사용해 네 가지 정보로 출력 SourceRecord를 만듭니다. 소스 파티션(하나뿐이며, 읽는 단일 파일), 소스 오프셋(파일의 바이트 오프셋), 출력 토픽 이름, 출력 값(줄이며, 이 값이 항상 문자열임을 나타내는 스키마 포함)입니다. SourceRecord 생성자의 다른 변형은 특정 출력 파티션, 키, 헤더도 포함할 수 있습니다.
이 구현은 일반 Java InputStream 인터페이스를 사용하며, 데이터가 없으면 sleep할 수 있습니다. Kafka Connect가 각 태스크에 전용 스레드를 제공하므로 이는 허용 가능합니다. 태스크 구현은 기본 poll() 인터페이스를 준수해야 하지만 구현 방식에는 많은 유연성이 있습니다. 이 경우 NIO 기반 구현이 더 효율적이겠지만, 이 간단한 방식은 작동하고 구현이 빠르며 이전 버전의 Java와도 호환됩니다.
예제에서 사용되지는 않지만 SourceTask는 소스 시스템에서 오프셋을 커밋하는 두 API commit과 commitRecord도 제공합니다. 이 API는 메시지에 대한 확인(acknowledgement) 메커니즘이 있는 소스 시스템을 위해 제공됩니다. 이 메서드들을 오버라이드하면 소스 커넥터가 Kafka에 기록된 후 소스 시스템의 메시지를 일괄 또는 개별로 확인할 수 있습니다. commit API는 poll이 반환한 오프셋까지 소스 시스템에 오프셋을 저장합니다. 이 API의 구현은 커밋이 완료될 때까지 차단해야 합니다. commitRecord API는 각 SourceRecord가 Kafka에 기록된 후 소스 시스템에 오프셋을 저장합니다. Kafka Connect가 오프셋을 자동으로 기록하므로 SourceTask가 이를 구현할 필요는 없습니다. 소스 시스템에서 메시지를 확인해야 하는 커넥터의 경우 일반적으로 두 API 중 하나만 필요합니다.
싱크 태스크 (Sink Tasks)
이전 섹션에서는 간단한 SourceTask를 구현하는 방법을 설명했습니다. SourceConnector와 SinkConnector와 달리 SourceTask와 SinkTask는 매우 다른 인터페이스를 가집니다. SourceTask는 pull 인터페이스를 사용하고 SinkTask는 push 인터페이스를 사용하기 때문입니다. 둘 다 공통 라이프사이클 메서드를 공유하지만 SinkTask 인터페이스는 꽤 다릅니다.
public abstract class SinkTask implements Task {
public void initialize(SinkTaskContext context) {
this.context = context;
}
public abstract void put(Collection<SinkRecord> records);
public void flush(Map<TopicPartition, OffsetAndMetadata> currentOffsets) {
}
}
SinkTask 문서에 전체 세부 사항이 포함되어 있지만, 이 인터페이스는 SourceTask만큼 거의 간단합니다. put() 메서드가 구현의 대부분을 담아야 하며, SinkRecord 집합을 받아 필요한 변환을 수행하고 대상 시스템에 저장합니다. 이 메서드는 반환하기 전에 데이터가 대상 시스템에 완전히 기록되었는지 확인할 필요가 없습니다. 사실 많은 경우 내부 버퍼링이 유용해서 전체 배치의 레코드를 한 번에 보내 다운스트림 데이터 저장소에 이벤트를 삽입하는 오버헤드를 줄일 수 있습니다. SinkRecord는 본질적으로 SourceRecord와 동일한 정보(Kafka 토픽, 파티션, 오프셋, 이벤트 키·값, 선택적 헤더)를 담고 있습니다.
flush() 메서드는 오프셋 커밋 과정에서 사용되며, 태스크가 장애에서 복구하고 안전한 지점에서 재개해 어떤 이벤트도 놓치지 않게 해줍니다. 이 메서드는 미결 데이터를 대상 시스템에 밀어 넣은 다음 쓰기가 확인될 때까지 차단해야 합니다. offsets 파라미터는 종종 무시할 수 있지만, 구현이 오프셋 정보를 대상 저장소에 저장해 정확히 한 번 전달을 제공하려는 경우에 유용합니다. 예를 들어 HDFS 커넥터는 이를 수행하고 원자적 이동 연산을 사용해 flush() 연산이 데이터와 오프셋을 HDFS의 최종 위치에 원자적으로 커밋하도록 보장할 수 있습니다.
오류 레코드 리포터 (Errant Record Reporter)
커넥터에 대해 오류 보고가 활성화되면 커넥터는 ErrantRecordReporter를 사용해 싱크 커넥터로 보내진 개별 레코드의 문제를 보고할 수 있습니다. 다음 예제는 커넥터의 SinkTask 하위 클래스가 ErrantRecordReporter를 얻고 사용하는 방법을 보여주며, DLQ가 활성화되지 않았거나 이 리포터 기능이 없는 이전 Connect 런타임에 커넥터가 설치된 경우 null 리포터를 안전하게 처리합니다.
private ErrantRecordReporter reporter;
@Override
public void start(Map<String, String> props) {
...
try {
reporter = context.errantRecordReporter(); // may be null if DLQ not enabled
} catch (NoSuchMethodException | NoClassDefFoundError e) {
// Will occur in Connect runtimes earlier than 2.6
reporter = null;
}
}
@Override
public void put(Collection<SinkRecord> records) {
for (SinkRecord record: records) {
try {
// attempt to process and send record to data sink
process(record);
} catch(Exception e) {
if (reporter != null) {
// Send errant record to error reporter
reporter.report(record, e);
} else {
// There's no error reporter, so fail
throw new ConnectException("Failed on record", e);
}
}
}
}
이전 오프셋에서 재개 (Resuming from Previous Offsets)
SourceTask 구현은 각 레코드에 스트림 ID(입력 파일 이름)와 오프셋(파일 내 위치)을 포함했습니다. 프레임워크는 이를 사용해 오프셋을 주기적으로 커밋하므로, 장애 발생 시 태스크가 복구되고 재처리·중복될 수 있는 이벤트 수를 최소화할 수 있습니다(Kafka Connect가 정상 종료된 경우, 예: standalone 모드 또는 작업 재구성으로 인한 경우에는 가장 최근 오프셋에서 재개). 이 커밋 과정은 프레임워크가 완전히 자동화하지만, 그 위치에서 재개하기 위해 입력 스트림의 올바른 위치로 seek하는 방법은 커넥터만 압니다.
시작 시 올바르게 재개하려면 태스크는 initialize() 메서드로 전달된 SourceContext를 사용해 오프셋 데이터에 접근할 수 있습니다. initialize()에서 오프셋(존재하는 경우)을 읽고 그 위치로 seek하는 코드를 조금 추가할 수 있습니다.
stream = new FileInputStream(filename);
Map<String, Object> offset = context.offsetStorageReader().offset(Collections.singletonMap(FILENAME_FIELD, filename));
if (offset != null) {
Long lastRecordedOffset = (Long) offset.get("position");
if (lastRecordedOffset != null)
seekToOffset(stream, lastRecordedOffset);
}
물론 각 입력 스트림에 대해 많은 키를 읽어야 할 수도 있습니다. OffsetStorageReader 인터페이스는 대량 읽기를 수행해 모든 오프셋을 효율적으로 로드한 다음, 각 입력 스트림을 적절한 위치로 seek해 적용할 수도 있게 해줍니다.
정확히 한 번 소스 커넥터 (Exactly-once source connectors)
정확히 한 번 지원 (Supporting exactly-once)
KIP-618의 통과로 Kafka Connect는 3.3.0 버전부터 정확히 한 번 소스 커넥터를 지원합니다. 소스 커넥터가 이 지원을 활용하려면 내보내는 각 레코드에 대해 의미 있는 소스 오프셋을 제공할 수 있어야 하고, 메시지를 버리거나 중복하지 않고 외부 시스템에서 해당 오프셋 중 어느 것에 해당하는 정확한 위치에서 소비를 재개할 수 있어야 합니다.
트랜잭션 경계 정의 (Defining transaction boundaries)
기본적으로 Kafka Connect 프레임워크는 소스 태스크가 poll 메서드에서 반환하는 각 레코드 배치에 대해 새 Kafka 트랜잭션을 만들고 커밋합니다. 하지만 커넥터가 자신의 트랜잭션 경계를 정의할 수도 있으며, 사용자는 커넥터 구성에서 transaction.boundary 속성을 connector로 설정해 이를 활성화할 수 있습니다.
활성화되면 커넥터의 태스크는 SourceTaskContext에서 TransactionContext에 접근할 수 있으며, 이를 사용해 트랜잭션을 언제 중단(abort)하고 커밋할지 제어할 수 있습니다.
예를 들어 적어도 10개 레코드마다 트랜잭션을 커밋하려면:
private int recordsSent;
@Override
public void start(Map<String, String> props) {
this.recordsSent = 0;
}
@Override
public List<SourceRecord> poll() {
List<SourceRecord> records = fetchRecords();
boolean shouldCommit = false;
for (SourceRecord record : records) {
if (++this.recordsSent >= 10) {
shouldCommit = true;
}
}
if (shouldCommit) {
this.recordsSent = 0;
this.context.transactionContext().commitTransaction();
}
return records;
}
또는 정확히 10번째 레코드마다 트랜잭션을 커밋하려면:
private int recordsSent;
@Override
public void start(Map<String, String> props) {
this.recordsSent = 0;
}
@Override
public List<SourceRecord> poll() {
List<SourceRecord> records = fetchRecords();
for (SourceRecord record : records) {
if (++this.recordsSent % 10 == 0) {
this.context.transactionContext().commitTransaction(record);
}
}
return records;
}
대부분의 커넥터는 자신의 트랜잭션 경계를 정의할 필요가 없습니다. 하지만 소스 시스템의 파일·객체가 여러 소스 레코드로 나뉘어져 있지만 원자적으로 전달되어야 하는 경우 유용할 수 있습니다. 또한 각 소스 레코드에 고유한 소스 오프셋을 부여하는 것이 불가능하고, 주어진 오프셋의 모든 레코드가 단일 트랜잭션 내에서 전달되는 경우에도 유용할 수 있습니다.
사용자가 커넥터 구성에서 커넥터 정의 트랜잭션 경계를 활성화하지 않았다면 context.transactionContext()가 반환하는 TransactionContext는 null이 된다는 점에 유의하세요.
검증 API (Validation APIs)
소스 커넥터 개발자가 구현할 수 있는 몇 가지 추가 사전 점검 검증 API가 있습니다.
일부 사용자는 커넥터에 정확히 한 번 의미론을 요구할 수 있습니다. 이 경우 커넥터 구성에서 exactly.once.support 속성을 required로 설정할 수 있습니다. 이렇게 되면 Kafka Connect 프레임워크는 커넥터에 지정된 구성으로 정확히 한 번 의미론을 제공할 수 있는지 묻습니다. 이는 커넥터에서 exactlyOnceSupport 메서드를 호출함으로써 수행됩니다.
커넥터가 정확히 한 번 의미론을 지원하지 않는다면, 정확히 한 번 의미론을 제공할 수 없다는 것을 사용자에게 확실히 알리기 위해 이 메서드를 구현해야 합니다.
@Override
public ExactlyOnceSupport exactlyOnceSupport(Map<String, String> props) {
// This connector cannot provide exactly-once semantics under any conditions
return ExactlyOnceSupport.UNSUPPORTED;
}
그렇지 않으면 커넥터는 구성을 검사하고, 정확히 한 번 의미론을 제공할 수 있다면 ExactlyOnceSupport.SUPPORTED를 반환해야 합니다.
@Override
public ExactlyOnceSupport exactlyOnceSupport(Map<String, String> props) {
// This connector can always provide exactly-once semantics
return ExactlyOnceSupport.SUPPORTED;
}
또한 사용자가 커넥터를 자신의 트랜잭션 경계를 정의하도록 구성했다면, Kafka Connect 프레임워크는 canDefineTransactionBoundaries 메서드를 사용해 커넥터가 지정된 구성으로 자신의 트랜잭션 경계를 정의할 수 있는지 물을 것입니다.
@Override
public ConnectorTransactionBoundaries canDefineTransactionBoundaries(Map<String, String> props) {
// This connector can always define its own transaction boundaries
return ConnectorTransactionBoundaries.SUPPORTED;
}
이 메서드는 어떤 경우에는 자신의 트랜잭션 경계를 정의할 수 있는 커넥터에 대해서만 구현해야 합니다. 커넥터가 자신의 트랜잭션 경계를 정의할 수 없는 경우 이 메서드를 구현할 필요가 없습니다.
동적 입·출력 스트림 (Dynamic Input/Output Streams)
Kafka Connect는 각 테이블을 개별적으로 복사하는 많은 작업을 만드는 대신 전체 데이터베이스를 복사하는 것과 같은 대량 데이터 복사 작업을 정의하기 위한 것입니다. 이 설계의 한 결과는 커넥터의 입력 또는 출력 스트림 집합이 시간에 따라 달라질 수 있다는 것입니다.
소스 커넥터는 변경(예: 데이터베이스의 테이블 추가/삭제)을 위해 소스 시스템을 모니터링해야 합니다. 변경을 감지하면 ConnectorContext 객체를 통해 프레임워크에 재구성이 필요하다고 알려야 합니다. 예를 들어 SourceConnector에서:
if (inputsChanged())
this.context.requestTaskReconfiguration();
프레임워크는 신속하게 새 구성 정보를 요청하고 태스크를 업데이트하며, 재구성 전에 진행 상황을 정상적으로 커밋할 수 있게 해줍니다. SourceConnector에서 이 모니터링은 현재 커넥터 구현에 맡겨진다는 점에 유의하세요. 이 모니터링을 수행하려면 추가 스레드가 필요하다면 커넥터가 스스로 할당해야 합니다.
이상적으로 변경 모니터링 코드는 Connector에 격리되어 태스크가 이를 신경 쓸 필요가 없어야 합니다. 하지만 변경은 태스크에도 영향을 줄 수 있으며, 가장 흔한 경우는 입력 시스템에서 입력 스트림 중 하나가 파괴될 때입니다(예: 데이터베이스에서 테이블이 삭제된 경우). 커넥터가 변경을 폴링해야 하기 때문에 종종 커넥터보다 태스크가 문제를 먼저 만나게 되며, 태스크는 이후 오류를 처리해야 합니다. 다행히 이는 보통 적절한 예외를 잡아 처리하는 것으로 간단히 해결할 수 있습니다.
SinkConnector는 보통 스트림 추가만 처리하면 되며, 이는 출력의 새 항목(예: 새 데이터베이스 테이블)으로 이어질 수 있습니다. 프레임워크는 정규식 구독 때문에 입력 토픽 집합이 변하는 것 같은 Kafka 입력의 모든 변경을 관리합니다. SinkTask는 새 입력 스트림을 예상해야 하며, 이는 다운스트림 시스템에 새 테이블 같은 새 리소스를 만들어야 할 수 있습니다. 이러한 경우 가장 까다로운 상황은 여러 SinkTask가 새 입력 스트림을 처음 보고 동시에 새 리소스를 만들려고 할 때의 충돌일 수 있습니다. 반면 SinkConnector는 일반적으로 동적 스트림 집합 처리를 위한 특별한 코드가 필요하지 않습니다.
구성 검증 (Configuration Validation)
Kafka Connect를 사용하면 커넥터를 제출해 실행하기 전에 커넥터 구성을 검증할 수 있으며, 오류와 권장 값에 대한 피드백을 제공할 수 있습니다. 이를 활용하려면 커넥터 개발자는 config() 구현을 제공해 구성 정의를 프레임워크에 노출해야 합니다.
FileStreamSourceConnector의 다음 코드는 구성을 정의하고 프레임워크에 노출합니다.
static final ConfigDef CONFIG_DEF = new ConfigDef()
.define(FILE_CONFIG, Type.STRING, null, Importance.HIGH, "Source filename. If not specified, the standard input will be used")
.define(TOPIC_CONFIG, Type.STRING, ConfigDef.NO_DEFAULT_VALUE, new ConfigDef.NonEmptyString(), Importance.HIGH, "The topic to publish data to")
.define(TASK_BATCH_SIZE_CONFIG, Type.INT, DEFAULT_TASK_BATCH_SIZE, Importance.LOW,
"The maximum number of records the source task can read from the file each time it is polled");
public ConfigDef config() {
return CONFIG_DEF;
}
ConfigDef 클래스는 예상 구성 집합을 지정하는 데 사용됩니다. 각 구성에 대해 이름, 타입, 기본값, 문서, 그룹 정보, 그룹 내 순서, 구성 값의 폭, UI에 표시하기에 적합한 이름을 지정할 수 있습니다. 또한 Validator 클래스를 오버라이드해 단일 구성 검증에 사용되는 특별한 검증 로직을 제공할 수 있습니다. 게다가 구성 간에는 의존성이 있을 수 있습니다. 예를 들어 어떤 구성의 유효 값과 가시성은 다른 구성의 값에 따라 달라질 수 있습니다. 이를 처리하기 위해 ConfigDef는 구성의 의존(dependent) 항목을 지정하고, 현재 구성 값을 고려해 유효 값을 얻고 구성의 가시성을 설정하는 Recommender 구현을 제공할 수 있게 해줍니다.
또한 Connector의 validate() 메서드는 허용된 구성 목록과 각 구성에 대한 구성 오류 및 권장 값을 반환하는 기본 검증 구현을 제공합니다. 하지만 기본 구현은 구성 검증에 권장 값을 사용하지 않습니다. 커스텀 구성 검증을 위해 권장 값을 사용할 수 있는 기본 구현의 오버라이드를 제공할 수 있습니다.
스키마 작업 (Working with Schemas)
FileStream 커넥터는 간단하므로 좋은 예제이지만, 데이터 구조도 사소합니다. 각 줄이 그저 문자열이기 때문입니다. 거의 모든 실제 커넥터는 더 복잡한 데이터 형식의 스키마가 필요할 것입니다.
더 복잡한 데이터를 만들려면 Kafka Connect 데이터 API로 작업해야 합니다. 대부분의 구조화된 레코드는 기본 타입 외에 Schema와 Struct 두 클래스와 상호작용해야 합니다.
API 문서에 완전한 참조가 있지만, 여기 Schema와 Struct를 만드는 간단한 예제가 있습니다.
Schema schema = SchemaBuilder.struct().name(NAME)
.field("name", Schema.STRING_SCHEMA)
.field("age", Schema.INT_SCHEMA)
.field("admin", SchemaBuilder.bool().defaultValue(false).build())
.build();
Struct struct = new Struct(schema)
.put("name", "Barbara Liskov")
.put("age", 75);
소스 커넥터를 구현한다면 스키마를 언제 어떻게 만들지 결정해야 합니다. 가능하면 스키마를 다시 계산하는 것을 최대한 피해야 합니다. 예를 들어 커넥터가 고정 스키마를 가질 것이 보장된다면 정적으로 만들고 단일 인스턴스를 재사용하세요.
하지만 많은 커넥터는 동적 스키마를 가집니다. 간단한 예가 데이터베이스 커넥터입니다. 단일 테이블만 고려해도 스키마는 커넥터 전체에 대해 미리 정의되지 않습니다(테이블마다 다르기 때문). 또한 사용자가 ALTER TABLE 명령을 실행할 수 있으므로 커넥터 수명 동안 단일 테이블에 대해 고정되어 있지 않을 수도 있습니다. 커넥터는 이러한 변경을 감지하고 적절히 대응할 수 있어야 합니다.
싱크 커넥터는 보통 더 간단한데, 데이터를 소비하므로 스키마를 만들 필요가 없기 때문입니다. 하지만 받은 스키마가 예상 형식인지 검증하는 데는 똑같이 주의를 기울여야 합니다. 스키마가 일치하지 않을 때(보통 업스트림 프로듀서가 대상 시스템으로 올바르게 변환할 수 없는 잘못된 데이터를 생성하고 있음을 나타냄) 싱크 커넥터는 시스템에 이 오류를 나타내기 위해 예외를 던져야 합니다.
더 알아보기 (Learn more)
SourceTask.pull()과SinkTask.put()인터페이스가 커넥터 구현의 핵심이에요.- 오프셋 커밋은 프레임워크가 자동화하지만, 올바른 위치로 seek하는 건 커넥터 책임이라는 걸 기억하세요.
- 정확히 한 번 지원은
exactlyOnceSupport와canDefineTransactionBoundaries같은 검증 API로 관리돼요.