Kafka Connect API를 활용한 데이터 파이프라인 구축

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로 가져오는 커넥터를 개발합니다. SourceConnectorSourceTask 클래스를 구현해야 합니다.

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 taskClass() {
        return CustomFileSourceTask.class;
    }

    @Override
    public List> taskConfigs(int maxTasks) {
        List> configs = new ArrayList<>();
        Map config = new HashMap<>();
        if (filename != null) {
            config.put(FILE_CONFIG, filename);
        }
        config.put(TOPIC_CONFIG, topic);
        configs.add(config);
        return configs;
    }

    @Override
    public void stop() {
    }

    @Override
    public ConfigDef config() {
        return CONFIG_DEF;
    }
}

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 토픽의 데이터를 외부 시스템으로 내보내는 커넥터를 개발합니다. SinkConnectorSinkTask 클래스를 구현해야 합니다.

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 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

태그: kafka kafka connect data pipeline source connector sink connector

8월 24일 08:48에 게시됨