Kafka Connect는 Kafka를 다양한 외부 데이터 소스 및 싱크와 통합하여 데이터 파이프라인을 구축하는 강력한 프레임워크입니다. 주로 데이터의 시작점 또는 끝점으로 사용되거나, Kafka를 데이터 전송의 중간 버퍼로 활용하는 시나리오에 적용됩니다.
Kafka Connect는 크게 두 가지 유형의 커넥터로 나뉩니다:
- Source 커넥터: 외부 시스템의 데이터를 Kafka 토픽으로 가져옵니다.
- Sink 커넥터: Kafka 토픽의 데이터를 외부 시스템으로 내보냅니다.
이러한 커넥터는 Kafka와 함께 배포되므로 별도의 설치가 필요하지 않습니다.
Kafka Connect의 주요 특징
- 표준화된 프레임워크: Kafka와 다른 시스템 간의 통합을 위한 표준을 제공하여 커넥터 개발, 배포 및 관리를 간소화합니다.
- 단일 모드 및 분산 모드 지원: 소규모 개발 및 테스트부터 대규모 프로덕션 환경까지 유연하게 확장 가능합니다.
- REST API: RESTful API를 통해 커넥터의 제출, 관리 및 모니터링을 수행할 수 있습니다.
- 자동 오프셋 관리: 커넥터는 데이터 처리의 진행 상황을 자동으로 추적하고 관리합니다.
- 분산 및 확장성: Kafka의 기존 그룹 관리 프로토콜을 기반으로 하여, 더 많은 커넥터 인스턴스를 추가하여 수평적 확장이 가능합니다.
- 스트림 및 배치 통합: Kafka의 능력을 활용하여 데이터 스트림과 배치 시스템을 연결하는 이상적인 솔루션입니다.
Kafka Connect 핵심 개념
- 커넥터 인스턴스: 데이터 흐름의 방향을 결정하며, 데이터의 복제 원천과 목적지를 지정합니다. Kafka와 외부 시스템 간의 논리적 처리를 담당합니다.
- 태스크 수: 분산 모드에서 커넥터 인스턴스는 작업을 여러 개의 태스크로 분할하여 여러 워커 스레드에 분산시킬 수 있습니다. 태스크는 상태 정보를 유지하지 않으며, 상태는 특정 Kafka 토픽(예:
offset.storage.topic,status.storage.topic)에 저장됩니다. - 워커 스레드: 커넥터 인스턴스와 태스크는 실제 실행을 위해 워커 스레드를 필요로 하며, 단일 모드와 분산 모드로 구성됩니다.
- 변환기 (Converter): 데이터를 Kafka Connect의 내부 형식과 바이트 형식 간에 상호 변환하는 역할을 합니다.
Kafka Connect 사용 예제
1. 단일 모드 (Standalone Mode)
단일 모드는 개발 및 테스트 환경에 적합합니다.
단일 모드 설정 파일
config/connect-standalone.properties 파일에 다음과 같이 설정합니다.
# Kafka 브로커 주소 bootstrap.servers=kafka1:9092,kafka2:9092,kafka3:9092 # 키-값 변환기 설정 (JSON 사용) key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter # 스키마 포함 여부 key.converter.schemas.enable=true value.converter.schemas.enable=true # 오프셋 저장 파일 경로 offset.storage.file.filename=/tmp/connect.offsets # 오프셋 플러시 간격 (밀리초) offset.flush.interval.ms=10000
파일 데이터를 Kafka 토픽으로 가져오기 (Source)
config/connect-file-source.properties 파일을 다음과 같이 설정합니다.
# 커넥터 이름 name=local-file-source # 커넥터 클래스 지정 connector.class=FileStreamSource # 최대 태스크 수 tasks.max=1 # 읽어올 파일 경로 file=/tmp/test.txt # 데이터를 쓸 Kafka 토픽 topic=connect_test
/tmp/test.txt 파일을 생성하고 데이터를 추가합니다.
kafka hadoop kafka-connect
단일 모드 커넥터를 실행하여 데이터를 Kafka로 로드합니다.
connect-standalone.sh config/connect-standalone.properties config/connect-file-source.properties
Kafka 콘솔 소비자로 데이터를 확인합니다.
kafka-console-consumer.sh --bootstrap-server kafka1:9092 --topic connect_test --from-beginning
결과:
{"schema":{"type":"string","optional":false},"payload":"kafka"}
{"schema":{"type":"string","optional":false},"payload":"hadoop"}
{"schema":{"type":"string","optional":false},"payload":"kafka-connect"}
/tmp/test.txt 파일에 데이터를 추가하면 소비자는 실시간으로 새로운 데이터를 받을 수 있습니다.
echo java >> /tmp/test.txt echo python >> /tmp/test.txt
Kafka 토픽 데이터를 파일로 내보내기 (Sink)
config/connect-file-sink.properties 파일을 다음과 같이 설정합니다.
# 커넥터 이름 name=local-file-sink # 커넥터 클래스 지정 connector.class=FileStreamSink # 최대 태스크 수 tasks.max=1 # 데이터를 쓸 파일 경로 file=/tmp/sink.txt # 데이터를 가져올 Kafka 토픽 topics=connect_test
단일 모드 커넥터를 실행하여 데이터를 Kafka에서 파일로 내보냅니다.
connect-standalone.sh config/connect-standalone.properties config/connect-file-sink.properties
/tmp/sink.txt 파일 내용을 확인합니다.
python kafka hadoop kafka-connect java
2. 분산 모드 (Distributed Mode)
분산 모드는 프로덕션 환경에 적합하며, 자동 로드 밸런싱, 복원력 및 확장성을 제공합니다.
분산 모드에서는 오프셋, 설정, 상태 정보가 Kafka 토픽에 저장됩니다. 따라서 사전에 해당 토픽들을 생성하는 것이 좋습니다.
커넥터 관련 토픽 생성
# 오프셋 저장 토픽 kafka-topics.sh --create --bootstrap-server kafka1:9092 --replication-factor 3 --partitions 1 --topic connect-offsets # 설정 저장 토픽 kafka-topics.sh --create --bootstrap-server kafka1:9092 --replication-factor 3 --partitions 6 --topic connect-configs # 상태 저장 토픽 kafka-topics.sh --create --bootstrap-server kafka1:9092 --replication-factor 3 --partitions 6 --topic connect-status
분산 모드 설정 파일
config/connect-distributed.properties 파일을 다음과 같이 설정합니다.
# Kafka 브로커 주소 bootstrap.servers=kafka1:9092,kafka2:9092,kafka3:9092 # 커넥터 그룹 ID group.id=connect-cluster # 키-값 변환기 설정 key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=true value.converter.schemas.enable=true # 오프셋, 설정, 상태 저장 토픽 지정 offset.storage.topic=connect-offsets config.storage.topic=connect-configs status.storage.topic=connect-status # 오프셋 플러시 간격 offset.flush.interval.ms=10000
분산 모드 커넥터를 시작합니다.
connect-distributed.sh config/connect-distributed.properties
커넥터 버전 정보를 확인합니다.
curl http://kafka1:8083
결과:
{"version":"2.7.0","commit":"448719dc99a19793","kafka_cluster_id":"wp8iI172SaqLHqNvEh3T-w"}
설치된 커넥터 플러그인 목록을 조회합니다.
curl http://kafka1:8083/connector-plugins -s | jq
결과 (일부):
[
{
"class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
"type": "sink",
"version": "2.7.0"
},
{
"class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
"type": "source",
"version": "2.7.0"
}
// ... 기타 플러그인
]
커넥터 REST API
Kafka Connect는 포트 8083에서 REST API를 제공하여 커넥터를 관리합니다. 주요 API는 다음과 같습니다:
GET /connectors: 활성 커넥터 목록 조회POST /connectors: 새 커넥터 생성 (요청 본문에 커넥터 이름과 설정 포함)GET /connectors/{name}: 특정 커넥터 정보 조회GET /connectors/{name}/config: 특정 커넥터 설정 조회PUT /connectors/{name}/config: 특정 커넥터 설정 업데이트GET /connectors/{name}/status: 특정 커넥터 상태 조회GET /connectors/{name}/tasks: 특정 커넥터의 태스크 목록 조회GET /connectors/{name}/tasks/{taskId}/status: 특정 태스크 상태 조회PUT /connectors/{name}/pause: 커넥터 및 태스크 일시 중지PUT /connectors/{name}/resume: 일시 중지된 커넥터 재개POST /connectors/{name}/restart: 커넥터 재시작POST /connectors/{name}/tasks/{taskId}/restart: 개별 태스크 재시작DELETE /connectors/{name}: 커넥터 삭제GET /connector-plugins: 설치된 커넥터 플러그인 목록 조회PUT /connector-plugins/{connector-type}/config/validate: 커넥터 설정 유효성 검사
파일 데이터를 Kafka 토픽으로 가져오기 (Source - REST API)
REST API를 사용하여 커넥터를 생성합니다. (예: Chrome 확장 프로그램 API Tester 사용)
요청 URL: http://kafka1:8083/connectors
요청 Body:
{
"name": "distributed-console-source",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
"tasks.max": "1",
"topic": "distributed_connect_test",
"file": "/tmp/distributed_test.txt"
}
}
생성된 커넥터 확인:
curl http://kafka1:8083/connectors -s | jq
결과:
[ "distributed-console-source" ]
/tmp/distributed_test.txt 파일을 생성하고 데이터를 추가합니다.
distributed_kafka kafka hadoop
Kafka 콘솔 소비자로 데이터를 확인합니다.
kafka-console-consumer.sh --bootstrap-server kafka1:9092 --topic distributed_connect_test --from-beginning
결과:
{"schema":{"type":"string","optional":false},"payload":"distributed_kafka"}
{"schema":{"type":"string","optional":false},"payload":"kafka"}
{"schema":{"type":"string","optional":false},"payload":"hadoop"}
Kafka 토픽 데이터를 파일로 내보내기 (Sink - REST API)
REST API를 사용하여 Sink 커넥터를 생성합니다.
요청 URL: http://kafka1:8083/connectors
요청 Body:
{
"name": "distributed-console-sink",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSinkConnector",
"tasks.max": "1",
"topics": "distributed_connect_test",
"file": "/tmp/distributed_sink.txt"
}
}
현재 커넥터 목록 확인:
curl http://kafka1:8083/connectors -s | jq
결과:
[ "distributed-console-sink", "distributed-console-source" ]
/tmp/distributed_sink.txt 파일 내용을 확인합니다.
distributed_kafka kafka hadoop
사용자 정의 Kafka 커넥터 개발
사용자 정의 커넥터는 Source 커넥터와 Sink 커넥터로 개발할 수 있습니다.
Source 커넥터 개발
외부 시스템의 데이터를 Kafka로 가져오는 커넥터를 개발합니다. SourceConnector 및 SourceTask 클래스를 구현해야 합니다.
SourceConnector 구현
package book_8;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import org.apache.kafka.common.config.ConfigDef;
import org.apache.kafka.common.config.ConfigDef.Importance;
import org.apache.kafka.common.config.ConfigDef.Type;
import org.apache.kafka.common.utils.AppInfoParser;
import org.apache.kafka.connect.connector.Task;
import org.apache.kafka.connect.errors.ConnectException;
import org.apache.kafka.connect.source.SourceConnector;
public class CustomFileSourceConnector extends SourceConnector {
public static final String TOPIC_CONFIG = "topic";
public static final String FILE_CONFIG = "file";
private static final ConfigDef CONFIG_DEF = new ConfigDef()
.define(FILE_CONFIG, Type.STRING, Importance.HIGH, "Source filename.")
.define(TOPIC_CONFIG, Type.STRING, Importance.HIGH, "The topic to publish data to");
private String filename;
private String topic;
@Override
public String version() {
return AppInfoParser.getVersion();
}
@Override
public void start(Map props) {
filename = props.get(FILE_CONFIG);
topic = props.get(TOPIC_CONFIG);
if (topic == null || topic.isEmpty()) {
throw new ConnectException("FileStreamSourceConnector configuration must include 'topic' setting");
}
if (topic.contains(",")) {
throw new ConnectException("FileStreamSourceConnector should only have a single topic when used as a source.");
}
}
@Override
public Class extends Task> taskClass() {
return CustomFileSourceTask.class;
}
@Override
public List
SourceTask 구현
package book_8;
import java.io.BufferedReader;
import java.io.FileInputStream;
import java.io.FileNotFoundException;
import java.io.IOException;
import java.io.InputStream;
import java.io.InputStreamReader;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import org.apache.kafka.connect.data.Schema;
import org.apache.kafka.connect.errors.ConnectException;
import org.apache.kafka.connect.source.SourceRecord;
import org.apache.kafka.connect.source.SourceTask;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class CustomFileSourceTask extends SourceTask {
private static final Logger LOG = LoggerFactory.getLogger(CustomFileSourceTask.class);
public static final String FILENAME_FIELD = "filename";
public static final String POSITION_FIELD = "position";
private static final Schema VALUE_SCHEMA = Schema.STRING_SCHEMA;
private String filename;
private InputStream stream = null;
private BufferedReader reader = null;
private char[] buffer = new char[1024];
private int offset = 0;
private String topic = null;
private Long streamOffset;
@Override
public String version() {
return new CustomFileSourceConnector().version();
}
@Override
public void start(Map props) {
filename = props.get(CustomFileSourceConnector.FILE_CONFIG);
topic = props.get(CustomFileSourceConnector.TOPIC_CONFIG);
if (filename == null || filename.isEmpty()) {
stream = System.in;
streamOffset = null;
reader = new BufferedReader(new InputStreamReader(stream, StandardCharsets.UTF_8));
} else {
try {
Map offset = context.offsetStorageReader().offset(Collections.singletonMap(FILENAME_FIELD, filename));
if (offset != null) {
Object lastRecordedOffset = offset.get(POSITION_FIELD);
if (!(lastRecordedOffset instanceof Long)) {
throw new ConnectException("Offset position is the incorrect type");
}
streamOffset = (Long) lastRecordedOffset;
} else {
streamOffset = 0L;
}
} catch (Exception e) {
LOG.error("Error retrieving offset for file {}: {}", filename, e.getMessage());
streamOffset = 0L; // Fallback to start if offset retrieval fails
}
}
if (topic == null) {
throw new ConnectException("FileStreamSourceTask config missing topic setting");
}
}
@Override
public List<SourceRecord> poll() throws InterruptedException {
if (stream == null) {
try {
stream = new FileInputStream(filename);
if (streamOffset != null && streamOffset > 0) {
LOG.debug("Seeking to file offset {}", streamOffset);
long skipped = stream.skip(streamOffset);
if (skipped != streamOffset) {
LOG.warn("Skipped only {} bytes, expected {}", skipped, streamOffset);
}
}
reader = new BufferedReader(new InputStreamReader(stream, StandardCharsets.UTF_8));
LOG.debug("Opened {} for reading", logFilename());
} catch (FileNotFoundException e) {
LOG.warn("File {} not found, waiting...", logFilename());
synchronized (this) {
this.wait(1000);
}
return null;
} catch (IOException e) {
LOG.error("Error opening file {}: {}", logFilename(), e.getMessage());
throw new ConnectException(e);
}
}
try {
final BufferedReader currentReader;
synchronized (this) {
currentReader = reader;
}
if (currentReader == null) return null;
ArrayList<SourceRecord> records = new ArrayList<>();
int nread = 0;
while (currentReader.ready()) {
nread = currentReader.read(buffer, offset, buffer.length - offset);
if (nread > 0) {
offset += nread;
if (offset == buffer.length) {
char[] newbuf = new char[buffer.length * 2];
System.arraycopy(buffer, 0, newbuf, 0, buffer.length);
buffer = newbuf;
}
String line;
do {
line = extractLine();
if (line != null) {
if (records.isEmpty()) {
LOG.debug("Reading line from {}: '{}'", logFilename(), line);
}
records.add(new SourceRecord(offsetKey(filename), offsetValue(streamOffset), topic, null, null, VALUE_SCHEMA, line, System.currentTimeMillis()));
}
} while (line != null);
}
}
if (nread <= 0) {
synchronized (this) {
this.wait(1000); // Wait a bit if no data is ready
}
}
if (!records.isEmpty()) {
return records;
} else {
return null; // Return null if no records were produced in this poll
}
} catch (IOException e) {
LOG.error("Error reading from file {}: {}", logFilename(), e.getMessage());
throw new ConnectException(e);
}
}
private String extractLine() {
int until = -1, newStart = -1;
for (int i = 0; i < offset; i++) {
if (buffer[i] == '\n') {
until = i;
newStart = i + 1;
break;
} else if (buffer[i] == '\r') {
if (i + 1 >= offset) return null;
until = i;
newStart = (buffer[i + 1] == '\n') ? i + 2 : i + 1;
break;
}
}
if (until != -1) {
String result = new String(buffer, 0, until);
int remaining = buffer.length - newStart;
if (remaining > 0) {
System.arraycopy(buffer, newStart, buffer, 0, remaining);
}
offset = remaining;
if (streamOffset != null) {
streamOffset += newStart;
}
return result;
} else {
return null;
}
}
@Override
public void stop() {
LOG.trace("Stopping task");
synchronized (this) {
try {
if (stream != null && stream != System.in) {
stream.close();
LOG.debug("Closed input stream for {}", logFilename());
}
} catch (IOException e) {
LOG.error("Failed to close stream: {}", e.getMessage());
}
this.notify();
}
}
private Map offsetKey(String filename) {
return Collections.singletonMap(FILENAME_FIELD, filename);
}
private Map offsetValue(Long pos) {
return Collections.singletonMap(POSITION_FIELD, pos);
}
private String logFilename() {
return filename == null ? "stdin" : filename;
}
}
Sink 커넥터 개발
Kafka 토픽의 데이터를 외부 시스템으로 내보내는 커넥터를 개발합니다. SinkConnector 및 SinkTask 클래스를 구현해야 합니다.
SinkConnector 구현
package book_8;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import org.apache.kafka.common.config.ConfigDef;
import org.apache.kafka.common.config.ConfigDef.Importance;
import org.apache.kafka.common.config.ConfigDef.Type;
import org.apache.kafka.common.utils.AppInfoParser;
import org.apache.kafka.connect.connector.Task;
import org.apache.kafka.connect.sink.SinkConnector;
public class CustomFileSinkConnector extends SinkConnector {
public static final String FILE_CONFIG = "file";
private static final ConfigDef CONFIG_DEF = new ConfigDef()
.define(FILE_CONFIG, Type.STRING, Importance.HIGH, "Destination filename.");
private String filename;
@Override
public String version() {
return AppInfoParser.getVersion();
}
@Override
public void start(Map props) {
filename = props.get(FILE_CONFIG);
}
@Override
public Class extends Task> taskClass() {
return CustomFileSinkTask.class;
}
@Override
public List> taskConfigs(int maxTasks) {
List> configs = new ArrayList<>();
for (int i = 0; i < maxTasks; i++) {
Map config = new HashMap<>();
if (filename != null) {
config.put(FILE_CONFIG, filename);
}
configs.add(config);
}
return configs;
}
@Override
public void stop() {
}
@Override
public ConfigDef config() {
return CONFIG_DEF;
}
}
SinkTask 구현
package book_8;
import java.io.FileNotFoundException;
import java.io.FileOutputStream;
import java.io.PrintStream;
import java.io.UnsupportedEncodingException;
import java.nio.charset.StandardCharsets;
import java.util.Collection;
import java.util.Map;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.connect.errors.ConnectException;
import org.apache.kafka.connect.sink.SinkRecord;
import org.apache.kafka.connect.sink.SinkTask;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class CustomFileSinkTask extends SinkTask {
private static final Logger LOG = LoggerFactory.getLogger(CustomFileSinkTask.class);
private String filename;
private PrintStream outputStream;
// Constructor for testing
public CustomFileSinkTask(PrintStream outputStream) {
this.filename = null; // Indicate stdout
this.outputStream = outputStream;
}
public CustomFileSinkTask() {
// Default constructor
}
@Override
public String version() {
return new CustomFileSinkConnector().version();
}
@Override
public void start(Map props) {
filename = props.get(CustomFileSinkConnector.FILE_CONFIG);
if (filename == null) {
outputStream = System.out;
} else {
try {
// Append mode true, using UTF-8 encoding
outputStream = new PrintStream(new FileOutputStream(filename, true), false, StandardCharsets.UTF_8.name());
} catch (FileNotFoundException | UnsupportedEncodingException e) {
throw new ConnectException("Couldn't find or create file for FileStreamSinkTask: " + filename, e);
}
}
LOG.info("Sink task started for file: {}", logFilename());
}
@Override
public void put(Collection<SinkRecord> sinkRecords) {
for (SinkRecord record : sinkRecords) {
if (record.value() != null) {
LOG.trace("Writing line to {}: {}", logFilename(), record.value());
outputStream.println(record.value().toString());
} else {
LOG.trace("Skipping null value record from topic {} partition {} offset {}",
record.topic(), record.kafkaPartition(), record.kafkaOffset());
}
}
}
@Override
public void flush(Map offsets) {
LOG.trace("Flushing output stream for {}", logFilename());
outputStream.flush();
}
@Override
public void stop() {
LOG.info("Stopping task for file: {}", logFilename());
if (outputStream != null && outputStream != System.out) {
outputStream.close();
}
}
private String logFilename() {
return filename == null ? "stdout" : filename;
}
}
빌드 및 배포
개발한 커넥터 코드를 JAR 파일로 빌드합니다. 이 JAR 파일을 Kafka Connect 워커 노드의 plugin.path에 지정된 디렉토리(일반적으로 Kafka 설치 경로의 libs 또는 커스텀 플러그인 디렉토리)에 복사합니다. 그런 다음 Kafka Connect 클러스터를 재시작합니다.
재시작 후, REST API를 통해 커넥터 플러그인 목록을 조회하여 사용자 정의 커넥터가 성공적으로 로드되었는지 확인할 수 있습니다.
curl http://localhost:8083/connector-plugins -s | jq
결과 (사용자 정의 커넥터 포함):
[
{
"class": "book_8.CustomFileSinkConnector",
"type": "sink",
"version": "2.7.0"
},
{
"class": "book_8.CustomFileSourceConnector",
"type": "source",
"version": "2.7.0"
},
// ... 기존 커넥터들
]
파일 데이터를 Kafka 토픽으로 가져오기 (사용자 정의 Source)
REST API를 사용하여 사용자 정의 Source 커넥터를 생성합니다.
요청 URL: http://kafka1:8083/connectors
요청 Body:
{
"name": "customer-distributed-source",
"config": {
"connector.class": "book_8.CustomFileSourceConnector",
"tasks.max": "1",
"topic": "customer_distributed_connect_test",
"file": "/tmp/customer_distributed_test.txt"
}
}
생성된 커넥터 확인:
curl http://kafka1:8083/connectors -s | jq
결과:
[ "customer-distributed-source", "distributed-console-sink", "distributed-console-source" ]
/tmp/customer_distributed_test.txt 파일에 데이터를 추가합니다.
echo kubernetes >> /tmp/customer_distributed_test.txt echo netty >> /tmp/customer_distributed_test.txt
Kafka 콘솔 소비자로 데이터를 확인합니다.
kafka-console-consumer.sh --bootstrap-server kafka1:9092 --topic customer_distributed_connect_test --from-beginning
결과:
{"schema":{"type":"string","optional":false},"payload":"kubernetes"}
{"schema":{"type":"string","optional":false},"payload":"netty"}
Kafka 토픽 데이터를 파일로 내보내기 (사용자 정의 Sink)
REST API를 사용하여 사용자 정의 Sink 커넥터를 생성합니다.
요청 URL: http://kafka1:8083/connectors
요청 Body:
{
"name": "customer-distributed-sink",
"config": {
"connector.class": "book_8.CustomFileSinkConnector",
"tasks.max": "1",
"topics": "customer_distributed_connect_test",
"file": "/tmp/customer_distributed_sink.txt"
}
}
현재 커넥터 목록 확인:
curl http://kafka1:8083/connectors -s | jq
결과:
[ "customer-distributed-source", "distributed-console-sink", "distributed-console-source", "customer-distributed-sink" ]
/tmp/customer_distributed_sink.txt 파일 내용을 확인합니다.
kubernetes netty