Kafka Connect 완벽 가이드: 커넥터를 활용한 데이터 파이프라인 구축부터 스키마 관리와 장애 대응까지
1장: Kafka Connect로 데이터 파이프라인 구축하기#
1.1 들어가며: Kafka를 데이터 파이프라인의 중심으로 만들기#
카프카는 단순한 메시지 브로커를 넘어 대규모 데이터가 이동하는 데이터 스트리밍 플랫폼으로 활용되고 있습니다.
쇼핑몰에서 주문이 발생했다고 가정해 보겠습니다.
주문 데이터는 데이터베이스에 저장되어야 하고, 동시에 결제 시스템이나 재고 시스템으로 전달되어야 합니다. 분석 시스템에서는 실시간으로 주문량을 집계해야 하며, 장기 보관을 위해 오브젝트 스토리지에 저장할 수도 있습니다.
전통적인 방식이라면 각각의 시스템 사이에 별도의 프로그램을 개발해야 합니다.
DB → 주문 처리 프로그램 → Kafka
Kafka → 재고 처리 프로그램 → 재고 DB
Kafka → 분석 프로그램 → 분석 시스템
Kafka → 저장 프로그램 → Object Storage시스템이 늘어날수록 개발해야 할 연동 프로그램도 함께 증가합니다.
Kafka Connect는 이러한 반복적인 데이터 연동 작업을 표준화합니다.
┌─ PostgreSQL
│
├─ MySQL
│
Database ───────┤
│
└─ Oracle
│
▼
┌──────────────┐
│ Kafka Connect│
└──────────────┘
│
▼
Kafka
│
┌────────────┼────────────┐
▼ ▼ ▼
Elasticsearch S3 Database개발자가 모든 데이터 이동 로직을 직접 작성하는 대신, 이미 만들어진 Connector를 설치하고 설정하는 방식으로 데이터 파이프라인을 구성할 수 있습니다.
Kafka Connect의 핵심은 데이터를 어떻게 처리할 것인가보다 데이터를 어떻게 안정적으로 이동시킬 것인가에 집중한다는 데 있습니다.
1.2 Kafka Connect란 무엇인가#
Kafka Connect는 Apache Kafka와 외부 시스템 사이의 데이터 이동을 표준화하기 위한 프레임워크입니다.
대표적으로 다음과 같은 시스템을 연결할 수 있습니다.
- 관계형 데이터베이스
- NoSQL 데이터베이스
- Elasticsearch
- 오브젝트 스토리지
- 파일 시스템
- 로그 수집 시스템
- 메시징 시스템
- 데이터 웨어하우스
Kafka Connect는 크게 Source Connector와 Sink Connector로 구분합니다.
1.2.1 Source Connector#
Source Connector는 외부 시스템에서 데이터를 가져와 Kafka 토픽으로 전달합니다.
외부 시스템
│
▼
Source Connector
│
▼
Kafka Topic예를 들어 JDBC Source Connector를 사용하면 데이터베이스의 데이터를 Kafka 토픽으로 가져올 수 있습니다.
MySQL
│
▼
JDBC Source Connector
│
▼
Kafka
│
▼
orders-topic1.2.2 Sink Connector#
Sink Connector는 Kafka 토픽의 데이터를 외부 시스템으로 전달합니다.
Kafka Topic
│
▼
Sink Connector
│
▼
외부 시스템예를 들어 Kafka 토픽의 데이터를 Elasticsearch로 전달하면 실시간 검색 시스템을 구축할 수 있습니다.
Kafka
│
▼
Elasticsearch Sink Connector
│
▼
Elasticsearch1.3 Kafka Connect의 핵심 구성 요소#
Kafka Connect를 이해하려면 Connector, Task, Worker, Converter, Transform의 관계를 이해해야 합니다.
1.3.1 Connector#
Connector는 외부 시스템과 Kafka 사이의 연결 방법과 데이터 이동 작업을 정의하는 구성 요소입니다.
예를 들어 다음과 같은 Connector가 존재할 수 있습니다.
JdbcSourceConnector
JdbcSinkConnector
ElasticsearchSinkConnector
S3SinkConnector
MongoDbSourceConnector
MongoDbSinkConnectorConnector 자체가 모든 데이터를 직접 처리하는 것은 아닙니다.
실제 데이터 이동 작업은 Task가 담당합니다.
1.3.2 Task#
Task는 Connector가 실제로 실행시키는 작업 단위입니다.
하나의 Connector에 여러 개의 Task를 구성할 수 있으며, Connector가 처리할 수 있는 작업을 여러 Task로 나누어 병렬 처리할 수 있습니다.
Connector
│
├── Task 1
├── Task 2
├── Task 3
└── Task 4다만 tasks.max를 크게 설정한다고 해서 항상 처리량이 증가하는 것은 아닙니다.
실제 병렬 처리 수준은 Connector 자체의 구현과 외부 시스템의 특성에 영향을 받습니다.
1.3.3 Worker#
Worker는 Connector와 Task를 실제로 실행하는 Kafka Connect 프로세스입니다.
Kafka Connect는 크게 Standalone 모드와 Distributed 모드로 실행할 수 있습니다.
Standalone
Kafka Connect Worker
├─ Connector A
└─ Connector B분산 모드에서는 여러 Worker가 하나의 Connect 클러스터를 구성합니다.
Kafka Connect Cluster
┌──────────────┐
│ Worker 1 │
└──────────────┘
│
┌──────────────┐
│ Worker 2 │
└──────────────┘
│
┌──────────────┐
│ Worker 3 │
└──────────────┘Worker에 장애가 발생하면 Connect 클러스터가 작업을 재분배할 수 있기 때문에 운영 환경에서는 일반적으로 Distributed 모드를 사용합니다.
1.3.4 Converter#
Converter는 Kafka Connect의 데이터와 Kafka 메시지 사이의 직렬화 및 역직렬화를 담당합니다.
대표적인 형식은 다음과 같습니다.
- JSON
- JSON Schema
- Avro
- Protobuf
- String
- Byte Array
예를 들어 Avro Converter를 사용하면 Schema Registry와 함께 데이터 스키마를 관리할 수 있습니다.
Connect Record
│
▼
Avro Converter
│
├── Schema Registry
│
▼
Kafka Message1.3.5 Single Message Transform#
Single Message Transform, 즉 SMT는 Kafka Connect가 처리하는 개별 메시지를 변환하는 기능입니다.
예를 들어 다음과 같은 작업을 수행할 수 있습니다.
- 필드 이름 변경
- 필드 추가
- 필드 삭제
- 데이터 타입 변환
- 토픽 이름 변경
- 메시지 필드 추출
간단한 데이터 변환에는 SMT가 유용하지만, 복잡한 비즈니스 로직을 SMT에 넣는 것은 적절하지 않습니다.
복잡한 처리는 Kafka Streams나 별도의 애플리케이션으로 분리하는 것이 일반적입니다.
1.4 Kafka Connect는 어떻게 데이터를 이동시키는가#
Kafka Connect의 전체 흐름을 살펴보겠습니다.
외부 시스템
│
▼
Connector
│
▼
Task
│
▼
Converter
│
▼
KafkaSink의 경우에는 반대 방향으로 동작합니다.
Kafka
│
▼
Converter
│
▼
Task
│
▼
Connector
│
▼
외부 시스템일반적인 실행 과정은 다음과 같습니다.
1.4.1 Connector 설정#
관리자는 Connector의 이름과 연결 정보, 대상 토픽, 데이터 처리 방식 등을 설정합니다.
1.4.2 Worker가 Connector 실행#
Kafka Connect Worker가 Connector 설정을 받아 실행합니다.
1.4.3 Task 생성#
Connector가 필요한 Task를 생성하고 Worker가 이를 실행합니다.
1.4.4 데이터 이동#
Task가 외부 시스템에서 데이터를 읽거나 Kafka에서 데이터를 읽어 목적지로 전달합니다.
1.4.5 Offset 관리#
Kafka Connect는 작업 진행 위치를 관리합니다.
Source Connector에서는 외부 시스템에서 어디까지 데이터를 읽었는지를 추적할 수 있고, Sink Connector에서는 Kafka의 어느 위치까지 처리했는지를 관리합니다.
따라서 Worker가 재시작되더라도 처음부터 모든 데이터를 무조건 다시 처리하는 것이 아니라 저장된 진행 상태를 기반으로 작업을 이어갈 수 있습니다.
다만 Offset 관리가 곧 중복이 절대 발생하지 않는다는 의미는 아닙니다.
Connector와 외부 시스템의 특성에 따라 재처리나 중복이 발생할 수 있기 때문에 실제 운영에서는 멱등성과 중복 처리 전략까지 함께 설계해야 합니다.
1.5 Kafka Connect REST API로 Connector 관리하기#
Kafka Connect의 Distributed 모드에서는 REST API를 통해 Connector를 관리할 수 있습니다.
예를 들어 Connector를 등록할 때 다음과 같은 형태를 사용합니다.
{
"name": "orders-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"tasks.max": "1",
"connection.url": "jdbc:postgresql://postgres:5432/shop",
"connection.user": "kafka",
"connection.password": "password",
"mode": "incrementing",
"incrementing.column.name": "id",
"topic.prefix": "shop-"
}
}Connector의 상태도 REST API를 통해 확인할 수 있습니다.
이러한 REST 기반 관리 방식은 자동화된 운영 환경에서 특히 유용합니다.
CI/CD 파이프라인을 통해 Connector 설정을 배포하거나 운영 시스템에서 Connector 상태를 모니터링할 수도 있습니다.
1.6 JDBC Connector로 데이터베이스와 Kafka 연결하기#
JDBC Connector는 관계형 데이터베이스와 Kafka를 연결할 때 대표적으로 사용하는 Connector입니다.
JDBC Source Connector는 관계형 데이터베이스에서 데이터를 읽어 Kafka 토픽으로 전달하고, JDBC Sink Connector는 Kafka 데이터를 데이터베이스에 저장합니다.
JDBC Source
MySQL ───────────────────► Kafka
│
│
JDBC Sink ▼
MySQL ◄──────────────────── Kafka1.6.1 JDBC Source Connector#
JDBC Source Connector는 데이터베이스를 주기적으로 조회하여 새로운 데이터를 Kafka로 전달할 수 있습니다.
대표적인 설정은 다음과 같습니다.
connector.class=io.confluent.connect.jdbc.JdbcSourceConnector
connection.url=jdbc:postgresql://localhost:5432/shop
connection.user=kafka
connection.password=password
topic.prefix=shop-데이터를 가져오는 방식은 mode를 통해 결정할 수 있습니다.
대표적으로 다음과 같은 방식이 있습니다.
bulk
incrementing
timestamp
timestamp+incrementingincrementing 방식은 증가하는 숫자형 컬럼을 기준으로 새로운 데이터를 찾는 방식입니다.
예를 들어 다음과 같은 테이블이 있다고 가정합니다.
orders
id customer amount
1 Kim 10000
2 Lee 25000
3 Park 18000다음과 같이 설정할 수 있습니다.
mode=incrementing
incrementing.column.name=id그러면 Connector는 id를 기준으로 새로운 데이터를 추적할 수 있습니다.
JDBC Source Connector는 테이블별 증분 처리와 다양한 조회 방식을 지원하며, Connector가 마지막으로 처리한 위치를 추적해 장애 이후 작업을 이어갈 수 있습니다.
1.6.2 JDBC Sink Connector#
JDBC Sink Connector는 Kafka 토픽의 데이터를 관계형 데이터베이스에 저장합니다.
Kafka Topic
│
▼
JDBC Sink Connector
│
▼
PostgreSQL대표적인 설정은 다음과 같습니다.
connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
connection.url=jdbc:postgresql://localhost:5432/shop
connection.user=kafka
connection.password=password
topics=orders데이터 저장 방식은 insert.mode 등을 통해 제어할 수 있습니다.
대표적인 방식은 다음과 같습니다.
insert
upsert
update기본 키 처리 방식은 pk.mode와 pk.fields를 이용해 설정할 수 있습니다.
insert.mode=upsert
pk.mode=record_value
pk.fields=id이러한 설정을 사용하면 단순한 데이터 적재뿐만 아니라 기존 레코드 갱신을 포함한 파이프라인을 구성할 수 있습니다.
1.6.3 JDBC Driver도 반드시 확인해야 한다#
JDBC Connector를 설치했다고 모든 데이터베이스에 바로 연결할 수 있는 것은 아닙니다.
사용하려는 데이터베이스에 맞는 JDBC Driver가 Connect Worker에서 사용할 수 있어야 합니다.
특히 여러 Worker로 구성된 Distributed 환경에서는 필요한 Connector와 JDBC Driver를 각 Worker에서 사용할 수 있도록 구성해야 합니다.
Worker 1
├─ JDBC Connector
└─ PostgreSQL Driver
Worker 2
├─ JDBC Connector
└─ PostgreSQL Driver
Worker 3
├─ JDBC Connector
└─ PostgreSQL Driver이 부분을 놓치면 특정 Worker에서 Task가 실행될 때 Connector를 찾지 못하거나 데이터베이스 연결에 실패할 수 있습니다.
1.7 FileStream Connector와 로그 데이터 수집#
Kafka Connect에는 파일 데이터를 Kafka로 보내거나 Kafka 데이터를 파일로 저장하는 FileStream Source와 Sink Connector가 있습니다.
예를 들어 다음과 같이 구성할 수 있습니다.
log.txt
│
▼
FileStream Source
│
▼
KafkaSource 설정의 기본적인 형태는 다음과 같습니다.
name=local-file-source
connector.class=FileStreamSource
tasks.max=1
file=/tmp/test.txt
topic=connect-testSink는 반대 방향으로 동작합니다.
Kafka
│
▼
FileStream Sink
│
▼
output.txtname=local-file-sink
connector.class=FileStreamSink
tasks.max=1
file=/tmp/test.sink.txt
topics=connect-test다만 여기서 중요한 변화가 있습니다.
FileStream Connector는 학습과 데모를 위한 용도로 보는 것이 적절하며 운영 환경에서 일반적인 파일 수집 솔루션으로 사용하는 것은 권장되지 않습니다.
현재 Confluent 문서에서도 FileStream Connector를 학습 및 데모 목적으로만 사용하도록 안내하고 있으며, 운영 환경에서 파일을 읽는 경우에는 Spool Dir 계열 Connector 등의 별도 솔루션을 고려하도록 안내합니다. 또한 FileStream Connector 아티팩트는 별도의 플러그인 경로 설정이 필요할 수 있습니다.
따라서 실무에서는 다음과 같이 판단하는 것이 좋습니다.
학습 / 테스트
│
▼
FileStream Connector
운영 환경
│
├─ Spool Dir 계열 Connector
├─ Filebeat
├─ Fluent Bit
└─ Logstash1.8 Logstash와 Kafka Connect는 어떻게 다른가#
Logstash는 Kafka Connect의 Connector가 아닙니다.
둘은 모두 데이터 파이프라인 구축에 활용할 수 있지만 접근 방식이 다릅니다.
Logstash는 Elastic Stack에서 사용하는 데이터 수집 및 변환 파이프라인 도구이며 다음과 같은 구조를 사용합니다.
Input
│
▼
Filter
│
▼
Output예를 들어 로그 파일을 Kafka로 보내는 파이프라인을 구성할 수 있습니다.
로그 파일
│
▼
Logstash Input
│
▼
Filter
│
▼
Kafka Output
│
▼
KafkaLogstash 설정은 대략 다음과 같은 형태입니다.
input {
file {
path => "/var/log/app.log"
}
}
filter {
grok {
match => {
"message" => "%{COMBINEDAPACHELOG}"
}
}
}
output {
kafka {
bootstrap_servers => "kafka:9092"
topic_id => "application-logs"
}
}Logstash는 복잡한 로그 파싱과 필터링이 필요한 환경에서 유용합니다.
반면 Kafka Connect는 이미 만들어진 Connector를 활용해 외부 시스템과 Kafka 사이의 데이터 이동을 표준화하는 데 더 초점이 맞춰져 있습니다.
따라서 둘은 경쟁 관계라기보다 함께 사용할 수도 있습니다.
Application Log
│
▼
Logstash
│
▼
Kafka
│
├────► Elasticsearch
│
└────► S31.9 다양한 Kafka Connector 활용하기#
Kafka Connect 생태계에는 데이터베이스뿐만 아니라 다양한 시스템을 연결하는 Connector가 존재합니다.
1.9.1 Elasticsearch Connector#
Kafka의 데이터를 Elasticsearch로 전달할 수 있습니다.
Application
│
▼
Kafka
│
▼
Elasticsearch실시간 검색이나 로그 분석 시스템을 구축할 때 활용할 수 있습니다.
1.9.2 Amazon S3 Connector#
Kafka 데이터를 Amazon S3와 같은 오브젝트 스토리지에 저장하는 파이프라인을 구성할 수 있습니다.
Kafka
│
▼
S3 Sink Connector
│
▼
Amazon S3장기간 데이터를 보관하거나 데이터 레이크를 구성하는 경우 활용할 수 있습니다.
1.9.3 MongoDB Connector#
MongoDB와 Kafka 사이의 데이터 이동에도 Connector를 사용할 수 있습니다.
MongoDB
│
▼
Kafka
│
▼
다른 시스템MongoDB의 변경 데이터를 Kafka 기반 이벤트 파이프라인으로 연결하는 구조도 구성할 수 있습니다.
1.9.4 기타 Connector#
환경에 따라 다음과 같은 시스템과도 연동할 수 있습니다.
- Cassandra
- Redis
- Elasticsearch
- OpenSearch
- S3
- Azure Blob Storage
- Google Cloud Storage
- Snowflake
- BigQuery
- MongoDB
- JDBC 기반 관계형 데이터베이스
Connector의 지원 범위와 유지보수 상태는 제품 및 배포 환경에 따라 달라지므로 실제 운영 전에는 현재 지원 버전과 호환성을 확인해야 합니다.
1.10 Schema Registry로 데이터 구조 관리하기#
Kafka를 실제 시스템에 도입하면 또 하나의 문제가 등장합니다.
바로 데이터 구조가 변경되는 문제입니다.
예를 들어 처음에는 주문 메시지가 다음과 같았다고 가정하겠습니다.
{
"id": 1001,
"customer": "Kim",
"amount": 25000
}그런데 시간이 지나면서 currency 필드를 추가합니다.
{
"id": 1001,
"customer": "Kim",
"amount": 25000,
"currency": "KRW"
}Producer와 Consumer가 서로 다른 버전의 데이터 구조를 사용하고 있다면 문제가 발생할 수 있습니다.
Schema Registry는 이러한 문제를 해결하기 위해 메시지의 스키마를 중앙에서 관리하고 버전 및 호환성을 관리하는 데 사용됩니다.
Producer
│
▼
Serializer
│
├────────► Schema Registry
│
▼
Kafka
│
▼
Deserializer
│
▼
Consumer대표적인 스키마 형식은 다음과 같습니다.
- Avro
- JSON Schema
- Protobuf
1.10.1 Schema Registry가 필요한 이유#
스키마를 관리하면 다음과 같은 장점이 있습니다.
- 데이터 구조의 명확성
- 스키마 버전 관리
- Producer와 Consumer 간 호환성 관리
- 잘못된 데이터 구조의 전달 방지
- 데이터 변경에 대한 통제
예를 들어 다음과 같은 변경을 생각해 볼 수 있습니다.
Version 1
id
name
amount
↓
Version 2
id
name
amount
currency새로운 필드가 추가되더라도 기존 Consumer가 계속 데이터를 읽을 수 있는지 검증할 수 있습니다.
1.10.2 Kafka Connect와 Schema Registry#
Kafka Connect에서는 Converter와 Schema Registry를 함께 구성할 수 있습니다.
예를 들어 Avro Converter를 사용할 경우 다음과 같은 형태가 됩니다.
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://schema-registry:8081JDBC Source Connector와 Avro Converter를 함께 사용하면 데이터베이스의 구조를 기반으로 Kafka 메시지의 스키마를 생성하고 Schema Registry와 연계할 수 있습니다.
1.11 데이터 손실과 중복을 줄이는 운영 전략#
데이터 파이프라인을 구축할 때는 단순히 데이터를 이동시키는 것만으로 충분하지 않습니다.
다음과 같은 상황을 고려해야 합니다.
Producer 장애
Broker 장애
Network 장애
Worker 장애
Database 장애
Connector 장애
Schema 변경
잘못된 데이터1.11.1 Producer의 acks 설정#
Producer에서는 acks 설정을 통해 메시지 전달 확인 수준을 결정할 수 있습니다.
대표적으로 다음과 같습니다.
acks=0
acks=1
acks=allacks=all 또는 acks=-1은 가능한 복제본들이 메시지를 확인하도록 하는 방식으로 데이터 내구성을 높이는 데 활용할 수 있습니다.
그러나 acks 하나만으로 데이터 손실 문제가 완전히 해결되는 것은 아닙니다.
복제 설정, min.insync.replicas, Producer 재시도 및 멱등성 등 여러 설정을 함께 고려해야 합니다.
1.11.2 Topic Replication Factor#
Kafka Topic의 파티션은 여러 Broker에 복제할 수 있습니다.
Partition 0
├─ Broker 1
├─ Broker 2
└─ Broker 3Broker 하나에 문제가 발생하더라도 다른 복제본을 활용할 수 있기 때문에 단일 Broker에 데이터를 저장하는 것보다 장애 대응 능력이 높아집니다.
운영 환경에서는 단순히 복제 개수만 정하는 것이 아니라 다음 항목을 함께 고려해야 합니다.
Replication Factor
+
min.insync.replicas
+
acks=all
+
Producer Retry1.11.3 Dead Letter Queue#
Kafka Connect에서는 처리할 수 없는 레코드에 대한 오류 처리 전략을 구성할 수 있습니다.
대표적으로 다음과 같은 설정을 사용할 수 있습니다.
errors.tolerance=all
errors.log.enable=true
errors.deadletterqueue.topic.name=orders-dlq
errors.deadletterqueue.topic.replication.factor=3
errors.deadletterqueue.context.headers.enable=trueDLQ를 사용하면 정상적으로 처리할 수 없는 레코드를 별도의 Kafka Topic으로 보내고 나중에 원인을 분석할 수 있습니다.
Kafka Topic
│
▼
Sink Connector
│
├── 정상 ──────────► Database
│
└── 오류 ──────────► DLQ
│
▼
오류 분석단, DLQ는 데이터 손실을 자동으로 해결하는 기능이 아닙니다.
오류 레코드를 별도의 Topic으로 격리하는 기능에 가깝습니다.
따라서 운영 환경에서는 DLQ를 구성한 뒤 다음과 같은 처리 절차도 필요합니다.
오류 발생
│
▼
DLQ 저장
│
▼
원인 분석
│
▼
데이터 수정
│
▼
재처리1.12 Exactly Once를 무조건 기대하면 안 되는 이유#
Kafka를 공부하다 보면 Exactly Once라는 용어를 자주 접하게 됩니다.
하지만 Kafka Connect의 모든 Connector가 모든 외부 시스템에서 동일한 의미의 Exactly Once를 보장한다고 생각해서는 안 됩니다.
데이터 파이프라인은 다음과 같이 여러 시스템을 거쳐갑니다.
Database
↓
Kafka Connect
↓
Kafka
↓
Kafka Connect
↓
DatabaseKafka 내부에서 제공하는 보장과 외부 시스템에 데이터를 실제로 기록하는 과정의 보장은 서로 다를 수 있습니다.
따라서 실무에서는 다음을 함께 고려해야 합니다.
- At-least-once 처리
- 중복 데이터 처리
- Idempotent Sink
- Upsert
- Primary Key
- Offset 관리
- 재처리 전략
- DLQ
- 데이터 정합성 검증
특히 데이터베이스에 데이터를 저장할 때는 upsert나 적절한 Primary Key 설계를 통해 동일한 이벤트가 다시 들어와도 결과가 망가지지 않도록 설계하는 것이 중요합니다.
1.13 Kafka Connect 운영 환경 설계#
개발 환경에서는 다음과 같이 하나의 Worker만 사용해도 충분할 수 있습니다.
Kafka
│
▼
Kafka Connect
│
├─ JDBC Source
└─ JDBC Sink하지만 운영 환경에서는 여러 Worker로 Connect 클러스터를 구성할 수 있습니다.
Kafka Cluster
│
│
┌──────────┴──────────┐
│ │
Connect Worker 1 Connect Worker 2
│ │
┌────┴────┐ ┌────┴────┐
│ │ │ │
Task Task Task Task이때 모든 Worker가 동일한 Connector 플러그인과 필요한 라이브러리를 사용할 수 있도록 관리해야 합니다.
특히 JDBC Connector처럼 외부 Driver가 필요한 Connector는 Worker마다 필요한 Driver를 설치해야 합니다.
운영 환경에서는 다음 항목도 함께 관리하는 것이 좋습니다.
- Connector 상태 모니터링
- Task 실패 감지
- Worker 상태 모니터링
- Kafka Lag 모니터링
- DLQ 모니터링
- Connector 설정 버전 관리
- 비밀번호 및 인증정보 보호
- TLS 및 인증 설정
- Connector Plugin 버전 관리
- 데이터 재처리 절차
1.14 실전 데이터 파이프라인 구성하기#
이제 지금까지 배운 내용을 하나의 예제로 연결해 보겠습니다.
온라인 쇼핑몰의 주문 데이터가 PostgreSQL에 저장되고 있다고 가정하겠습니다.
PostgreSQL
│
│ JDBC Source
▼
Kafka
│
├──────────────► Elasticsearch
│
├──────────────► S3
│
└──────────────► 분석 시스템1.14.1 데이터베이스에서 Kafka로 가져오기#
JDBC Source Connector를 사용합니다.
PostgreSQL
│
▼
JDBC Source Connector
│
▼
orders-topic1.14.2 Kafka에서 검색 시스템으로 보내기#
Elasticsearch Sink Connector를 사용할 수 있습니다.
orders-topic
│
▼
Elasticsearch Sink
│
▼
Elasticsearch1.14.3 Kafka에서 장기 보관 시스템으로 보내기#
S3 Sink Connector 등을 활용하면 Kafka 데이터를 오브젝트 스토리지에 저장할 수 있습니다.
orders-topic
│
▼
S3 Sink
│
▼
Amazon S3이렇게 하면 Kafka를 데이터 흐름의 중심에 두고 다양한 시스템으로 데이터를 분배하는 구조를 만들 수 있습니다.
1.15 Kafka Connect를 사용할 때 흔히 하는 실수#
1.15.1 Connector와 Task를 같은 개념으로 생각하기#
Connector는 작업을 정의하고 Task는 실제 실행 단위입니다.
Connector
↓
Task
↓
실제 데이터 처리둘을 동일한 개념으로 이해하면 tasks.max 설정이나 병렬 처리 구조를 이해하기 어렵습니다.
1.15.2 tasks.max만 늘리면 빨라진다고 생각하기#
Task 수를 늘린다고 무조건 처리량이 증가하지는 않습니다.
Connector 구현 방식과 데이터 소스의 병렬 처리 능력에 따라 실제 성능이 결정됩니다.
1.15.3 FileStream Connector를 운영 환경에 그대로 사용하기#
FileStream Connector는 Kafka Connect의 구조를 이해하기에는 좋은 예제지만 운영 환경의 대규모 로그 수집 솔루션으로 사용하는 것은 적절하지 않습니다. 공식 문서에서도 학습 및 데모 용도로 안내하고 있습니다.
1.15.4 Schema 없이 대규모 시스템을 운영하기#
JSON만으로도 간단한 시스템을 만들 수 있지만 여러 팀과 서비스가 장기간 데이터를 공유한다면 데이터 구조 관리가 중요해집니다.
Schema Registry와 Avro, JSON Schema, Protobuf 등의 사용을 검토할 수 있습니다.
1.15.5 DLQ만 만들고 끝내기#
DLQ에 데이터를 넣는 것보다 중요한 것은 DLQ를 어떻게 처리할 것인가입니다.
오류
↓
DLQ
↓
모니터링
↓
원인 분석
↓
수정
↓
재처리이 과정이 없다면 DLQ는 단순한 데이터 저장소가 되어버립니다.
1.16 Kafka Connect를 도입해야 하는 경우#
Kafka Connect가 모든 데이터 연동 문제의 정답은 아닙니다.
다음과 같은 경우 Kafka Connect가 특히 적합합니다.
- 이미 존재하는 Connector가 있는 경우
- 단순한 데이터 이동이 필요한 경우
- DB와 Kafka를 연결해야 하는 경우
- Kafka 데이터를 검색 시스템에 전달해야 하는 경우
- Kafka 데이터를 Object Storage에 저장해야 하는 경우
- 여러 외부 시스템을 Kafka 중심으로 연결해야 하는 경우
- 데이터 이동 파이프라인을 표준화해야 하는 경우
반대로 복잡한 비즈니스 로직이 핵심이라면 별도의 애플리케이션이나 Kafka Streams 등을 검토하는 것이 적절합니다.
단순 데이터 이동
│
▼
Kafka Connect
복잡한 스트림 처리
│
▼
Kafka Streams
복잡한 비즈니스 로직
│
▼
별도 애플리케이션1.17 Kafka Connect 전체 구조 정리#
지금까지의 내용을 하나의 그림으로 정리하면 다음과 같습니다.
┌──────────────────┐
│ External DB │
└────────┬─────────┘
│
JDBC Source
│
▼
┌──────────────┐ ┌──────────────┐
│ Elasticsearch│◄─────│ │
└──────────────┘ │ Kafka │
│ │
┌──────────────┐ └──────┬───────┘
│ S3 │◄─────────────┤
└──────────────┘ Sink │
│
┌───────┴───────┐
│ │
Database AnalyticsKafka Connect는 이 구조에서 데이터 이동을 담당합니다.
Connector
↓
Task
↓
Worker
↓
Converter
↓
Kafka
↓
외부 시스템그리고 Schema Registry를 추가하면 데이터 구조까지 관리할 수 있습니다.
Producer / Connect
│
▼
Schema Registry
│
▼
Kafka
│
▼
Consumer / Connect결국 Kafka Connect의 핵심은 Kafka를 다양한 시스템을 연결하는 데이터 허브로 활용하는 것입니다.
1.18 요약: Kafka Connect 핵심 정리#
이번 장에서는 Kafka Connect를 이용해 Kafka와 외부 시스템을 연결하는 방법을 살펴보았습니다.
핵심 내용을 정리하면 다음과 같습니다.
- Kafka Connect는 Kafka와 외부 시스템 사이의 데이터 이동을 표준화하는 프레임워크입니다.
- Source Connector는 외부 시스템의 데이터를 Kafka로 가져옵니다.
- Sink Connector는 Kafka 데이터를 외부 시스템으로 전달합니다.
- Connector는 데이터 이동 작업을 정의하고 Task는 실제 작업을 수행합니다.
- Worker는 Connector와 Task를 실행하는 Kafka Connect 프로세스입니다.
- Converter는 데이터의 직렬화와 역직렬화를 담당합니다.
- SMT는 개별 메시지를 간단하게 변환할 때 사용할 수 있습니다.
- JDBC Connector를 사용하면 관계형 데이터베이스와 Kafka를 연결할 수 있습니다.
- FileStream Connector는 학습과 테스트에는 유용하지만 운영 환경의 파일 수집 용도로는 적절하지 않을 수 있습니다.
- Logstash는 Kafka Connect Connector가 아니라 별도의 데이터 수집 및 변환 도구입니다.
- Schema Registry를 사용하면 데이터 스키마와 호환성을 관리할 수 있습니다.
- DLQ는 오류 데이터를 별도 Topic으로 격리하여 분석과 재처리를 가능하게 합니다.
- 데이터 안정성을 확보하려면
acks, 복제,min.insync.replicas, Offset, 멱등성, 재처리 전략 등을 종합적으로 고려해야 합니다. - Connector와 JDBC Driver의 버전 및 설치 상태는 모든 Connect Worker에서 일관되게 관리해야 합니다.
Kafka Connect를 이해하면 Kafka를 단순한 메시지 브로커가 아니라 데이터베이스, 검색 시스템, 스토리지, 분석 시스템을 연결하는 데이터 파이프라인의 중심으로 활용할 수 있습니다.
다음 단계에서는 Kafka Connect를 통해 데이터를 이동시키는 것을 넘어, Kafka에 들어온 데이터를 실시간으로 필터링하고 집계하며 조합하는 Kafka Streams를 살펴볼 수 있습니다.
1.19 자체 점검 문제#
1.19.1 기본 개념#
- Kafka Connect의 역할을 설명하십시오.
- Source Connector와 Sink Connector의 차이를 설명하십시오.
- Connector와 Task의 차이를 설명하십시오.
- Worker의 역할을 설명하십시오.
- Converter와 SMT의 역할을 각각 설명하십시오.
1.19.2 데이터베이스 연동#
- JDBC Source Connector를 이용하여 PostgreSQL의 데이터를 Kafka로 가져오는 방법을 설명하십시오.
- JDBC Sink Connector를 이용하여 Kafka 데이터를 데이터베이스에 저장하는 방법을 설명하십시오.
incrementing,timestamp,timestamp+incrementing방식의 차이를 설명하십시오.- Distributed 모드에서 JDBC Driver를 각 Worker에 설치해야 하는 이유를 설명하십시오.
1.19.3 데이터 파이프라인 설계#
- FileStream Connector가 운영 환경에서 적합하지 않을 수 있는 이유를 설명하십시오.
- Logstash와 Kafka Connect의 차이를 설명하십시오.
- Kafka 데이터를 Elasticsearch로 전달하는 파이프라인을 설계하십시오.
- Kafka 데이터를 Object Storage에 장기 보관하는 파이프라인을 설계하십시오.
1.19.4 데이터 품질과 장애 대응#
- Schema Registry가 필요한 이유를 설명하십시오.
- Schema Registry에서 스키마 호환성이 중요한 이유를 설명하십시오.
acks=all의 의미를 설명하십시오.- Topic Replication Factor와
min.insync.replicas의 관계를 설명하십시오. - Dead Letter Queue가 필요한 상황을 설명하십시오.
- DLQ에 저장된 데이터를 어떻게 재처리할 것인지 설명하십시오.
- Kafka Connect에서 중복 데이터가 발생할 수 있는 상황과 이를 줄이기 위한 방법을 설명하십시오.
1.19.5 실전 문제#
- PostgreSQL → Kafka → Elasticsearch 구조의 주문 데이터 파이프라인을 설계하십시오.
- Kafka → S3 데이터 파이프라인을 구축한다면 어떤 데이터를 장기 보관할 것인지 설명하십시오.
- Connector Worker 한 대가 장애를 일으켰을 때 Distributed 모드에서 어떤 일이 발생할 수 있는지 설명하십시오.
- Connector가 특정 메시지를 처리하지 못했을 때 DLQ를 이용한 장애 대응 절차를 설계하십시오.
- 자신의 프로젝트에서 Kafka Connect를 도입한다면 어떤 Source Connector와 Sink Connector를 사용할지 선정하고 전체 데이터 흐름을 설계하십시오.