목차
kafka topic message 처리하기
아래 Java/Gradle 코드는 Kafka 2.2 시기의 실습 기록이다. 현재 환경에서 그대로 빌드되는 예제로 취급하지 않는다. 특히 자동 오프셋 커밋과 비동기 전송 결과 미확인은 메시지 손실 또는 중복을 만들 수 있다. 이 글 끝에서 처리 흐름과 수정 기준을 설명한다.
- Producer와 Consumer를 통해 topic의 메시지를 처리해본다.
요구사항
임의의 숫자 2개를 설정하여 2개의 숫자 범위 값에 들어오면 valid topic으로 존재하지 않으면 invalid topic으로 분류한다.
초기 숫자가 10 20이면
10~20 사이의 수는 valid 그렇지 않으면 invalid topic으로 이동
-
원본메시지(raw-message)는 kafka topic에서 각각 이벤트를 읽는다.
-
이벤트를 검사하고 올바른 메시지는 valid-message topic에 쓰고, 잘못된 메시지는 invalid-message 토픽에 기록한다.
[구성도]
- raw-messages topic에서 이벤트를 읽고, 메시지를 검사하고, 조건 범위 내에 숫자가 존재하지 않으면 invalid-message topic으로, 정상적인 이벤트는 valid-message topic으로 routing 한다.

구현
-
모델 정의
-
zooKeeper 실행
-
kafka 서버 실행(prot:9092)
-
3 개의 topic(source-topic, valid-topic, invalid-topic)을 생성
-
kafka에서 raw-messages(source-topic) topic을 읽고 검사 기능을 처리하는 Producer, Consumer와 작업 구현
0. 모델 정의 - JSON Data 정의
{ "name" : "hello", "number" : "30" }
{ "name" : "ik", "number" : "45" }
{ "name" : "honggil", "number" : "50" }
{ "name" : "baby", "number" : "1" }
{ "name" : "kate", "number" : "80" }
1. zookeeper 실행
$ ./bin/zookeeper-server-start.sh config/zookeeper.properties
2. kafka 서버 실행
$ ./bin/kafka-server-start.sh config/server.properties
3. 3개의 topic 생성
## source-topic, valid-topic, invalid-topic 3개의 topic 생성
$ ./kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic source-topic
$ ./kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic valid-topic
#
$ ./kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic invalid-topic
## 잘만들어졌는지 확인
$ ./kafka-topics.sh --list --bootstrap-server localhost:9092
[3개의 topic 생성 및 확인]

4. kafka에서 raw-messages(source-topic) topic을 읽고 검사 기능을 처리하는 Producer, Consumer와 작업 구현
0. 모델 정의
- JSON Data 정의

[Consumer.java]
import org.apache.kafka.clients.consumer.ConsumerRecords;
import java.util.Properties;
public interface Consumer {
public static Properties createConfig(String servers, String groupId) {
Properties props = new Properties();
props.put("bootstrap.servers", servers);
props.put("group.id", groupId);
props.put("enable.auto.commit", "true");
props.put("auto.commit.interval.ms", "1000");
props.put("auto.offset.reset", "earliest");
props.put("session.timeout.ms", "30000");
props.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
return props;
}
public ConsumerRecords<String, String> consume();
}
[Reader.java]
- consumer의 구현체
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import java.time.Duration;
import java.util.Collections;
public class Reader implements Consumer {
private final KafkaConsumer<String, String> consumer; //1. kafka consumer load
private final String topic;
public Reader(String servers, String groupId, String topic) {
this.consumer = new KafkaConsumer<String, String>(Consumer.createConfig(servers, groupId));
this.topic = topic;
}
@Override
public ConsumerRecords<String, String> consume() {
this.consumer.subscribe(Collections.singleton(this.topic)); //2. topic publish
return consumer.poll(Duration.ofMillis(100)); //3. set timeout
}
}
[Producer.java]
import java.util.Properties;
public interface Producer {
public void produce(String message);
public static Properties createConfig(String servers) {
Properties props = new Properties();
props.put("bootstrap.servers", servers);
props.put("acks", "all");
props.put("retries", 0);
props.put("batch.size", 1000);
props.put("linger.ms", 1);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
return props;
}
}
mport java.util.Properties;
public interface Producer {
public void produce(String message);
public static Properties createConfig(String servers) {
Properties props = new Properties();
props.put("bootstrap.servers", servers);
props.put("acks", "all");
props.put("retries", 0);
props.put("batch.size", 1000);
props.put("linger.ms", 1);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
return props;
}
}
[Validator.java]
- producer 구현체
import com.google.gson.Gson;
import com.study.kafka.model.CheckNumber;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
public class Validator implements Producer {
private final KafkaProducer<String, String> producer;
private final String validTopic;
private final String invalidTopic;
private final int startNum;
private final int endNum;
public Validator(String servers, String validTopic, String invalidTopic, int startNum, int endNum) {
this.producer = new KafkaProducer<String, String>(Producer.createConfig(servers));
this.validTopic = validTopic;
this.invalidTopic = invalidTopic;
this.startNum = startNum;
this.endNum = endNum;
}
@Override
public void produce(String message) {
ProducerRecord<String, String> pr = null;
Gson gson = new Gson();
CheckNumber chkNum = null;
System.out.println("data : " + message);
try {
chkNum = gson.fromJson(message, CheckNumber.class);
System.out.println("fromJson : " + chkNum.getName() + ", nunber=" + chkNum.getNumber());
String topic = validate(Integer.valueOf(chkNum.getNumber()), this.startNum, this.endNum) ? this.validTopic : this.invalidTopic;
String resultMsg = this.startNum + "~" + this.endNum + " number is " + topic;
chkNum.setResultMsg(resultMsg);
pr = new ProducerRecord<String, String>(topic, gson.toJson(chkNum));
producer.send(pr);
} catch (Exception e) {
System.out.println("Exception...");
pr = new ProducerRecord<String, String>(this.invalidTopic, message);
producer.send(pr);
return;
}
}
private boolean validate(int src, int startNum, int endNum) {
return (src >= startNum && src <= endNum);
}
}
[CheckNumber]
import lombok.Data;
@Data
public class CheckNumber {
private String name;
private String number;
private String resultMsg;
}
[Main.java]
import com.study.kafka.consumer.Reader;
import com.study.kafka.producer.Validator;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import java.util.Scanner;
public class Main {
//private static final Logger log = LoggerFactory.getLogger(Reader.class);
public static void main(String[] args) throws Exception {
String servers = args[0];
String groupId = args[1];
String sourceTopic = args[2];
String validTopic = args[3];
String invalidTopic = args[4];
Scanner scan = new Scanner(System.in);
System.out.print("input start number ->");
int startNum = scan.nextInt();
System.out.print("input end number ->");
int endNum = scan.nextInt();
System.out.println();
Reader reader = new Reader(servers, groupId, sourceTopic);
Validator validator = new Validator(servers, validTopic, invalidTopic, startNum, endNum);
while (true) {
ConsumerRecords<String, String> consumerRecords = reader.consume();
for (ConsumerRecord<String, String> record : consumerRecords) {
if (record.value() != null && !record.value().equals(""))
validator.produce(record.value());
}
}
}
}
[build.gradle]
apply plugin: 'java'
apply plugin:'application'
apply plugin: 'idea'
group 'study'
version '1.0-SNAPSHOT'
sourceCompatibility = 1.8
repositories {
mavenCentral()
}
dependencies {
annotationProcessor group: 'org.projectlombok', name: 'lombok', version: '1.18.4'
compileOnly group: 'org.projectlombok', name: 'lombok', version: '1.18.4'
compile group: 'org.apache.kafka', name: 'kafka-clients', version: '2.2.1'
compile group: 'com.google.code.gson', name: 'gson', version: '2.8.5'
testCompile group: 'junit', name: 'junit', version: '4.12'
}
mainClassName = 'com.study.kafka.Main'
jar {
manifest {
attributes 'Title': 'My Application', 'Main-Class': mainClassName
}
archiveFileName = 'CheckNumber.jar'
dependsOn configurations.runtime
from {
configurations.compile.collect {it.isDirectory()? it: zipTree(it)}
}{
exclude "META-INF/*.SF"
exclude "META-INF/*.DSA"
exclude "META-INF/*.RSA"
}
}
sourceSets.main.resources {
srcDirs = ['src/main/java']
include '**/*.xml'
}
결과
- source-topic producer console 실행
$ ./bin/kafka-console-producer.sh --broker-list localhost:9092 --topic source-topic
- invalid-topic consumer console 실행
$ ./bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic invalid-topic
- valid-topic consumer console 실행
$ ./bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic valid-topic
- CheckNumber.jar 실행
$ java -jar CheckNumber.jar localhost:9092 EventValidNum source-topic valid-topic invalid-topic

세 토픽의 의미와 판정 기준
입력 토픽 source-topic에서 JSON 한 건을 읽고 number를 정수로 해석한다. 설정 범위가 10~20이라면 10과 20도 유효하다. 범위 안이면 valid-topic, 그 밖이거나 JSON·숫자 해석에 실패하면 invalid-topic에 원본과 실패 이유를 기록한다. 원문 예제의 30, 45, 50, 1, 80은 모두 범위 밖이어서 유효 경로가 실제로 동작하는지 확인할 수 없다. 아래 세 건으로 유효·범위 초과·형식 오류를 각각 시험한다.
{"eventId":"a1","name":"min","number":15}
{"eventId":"a2","name":"jin","number":30}
{"eventId":"a3","name":"sol","number":"unknown"}
기대 결과는 첫 건이 valid-topic, 나머지 두 건이 invalid-topic이다. 문자열 "15"도 허용할지, 숫자형 15만 허용할지는 계약에서 정한다. 입력에 eventId를 두면 재처리 중복을 식별할 수 있다. 오류 토픽에는 원본 본문뿐 아니라 오류 유형, 발생 시각, 소스 파티션·오프셋을 함께 보관하면 다시 처리하기 쉽다. 다만 개인 정보가 포함된 본문을 오류 로그나 토픽에 그대로 복제하지 않도록 보존·접근 정책을 정해야 한다.
flowchart LR S[source-topic] --> P[JSON 파싱과 number 검증] P -->|10 이상 20 이하| V[valid-topic] P -->|범위 밖·형식 오류| I[invalid-topic] V --> O[출력 성공 확인] I --> O O --> C[입력 오프셋 커밋]
전송 성공과 오프셋 사이의 실패
원문은 소비자에 enable.auto.commit=true를 설정하고 producer.send(pr)의 결과를 확인하지 않는다. send는 일반적으로 비동기라서 호출이 반환됐다는 것만으로 브로커 저장 성공을 뜻하지 않는다. 입력 오프셋이 먼저 커밋된 다음 출력 전송이 실패하면 입력은 다시 읽지 않으면서 결과 토픽에도 기록되지 않는 손실이 생긴다. 반대로 출력은 성공했지만 오프셋 커밋 전에 프로세스가 죽으면 입력을 다시 읽어 결과 토픽에 중복이 생길 수 있다.
최소한 출력 전송의 callback 또는 Future 결과를 확인한 뒤 해당 입력의 오프셋을 커밋하는 방식으로 손실 경로를 줄인다. 그래도 출력 성공 후 커밋 전에 죽는 중복 가능성은 남는다. 같은 Kafka 클러스터 안에서 입력 오프셋과 출력 기록을 원자적으로 묶어야 한다면 Kafka 트랜잭션을 검토하고, 외부 DB나 HTTP까지 함께 쓰는 경우에는 별도의 중복 처리·일관성 설계가 필요하다. acks=all도 복제 계수 1인 실습 클러스터에서는 단일 브로커의 장애 복원력을 만들지 못한다.
원문의 Consumer와 Producer는 Kafka 클라이언트 클래스와 이름이 겹쳐 독자가 헷갈릴 수 있다. Reader에서 매번 subscribe하기보다 시작 시 한 번 구독하고 루프에서 poll하는 편이 흐름이 분명하다. 종료 시에는 consumer와 producer를 닫아 버퍼 전송과 그룹 정리를 수행한다. Validator가 예외를 모두 잡아 invalid 토픽으로 보내면 잘못된 데이터와 브로커 전송 실패를 구별할 수 없으므로 파싱 오류만 분기하고 시스템 오류는 로그·재시도·중단으로 처리한다.
| 관찰 | 우선 확인 |
|---|---|
| valid 토픽이 비어 있음 | 테스트에 실제 범위 안 숫자가 있는지, 값 타입이 맞는지 |
| 두 토픽 모두 비어 있음 | 입력 그룹 오프셋, bootstrap-server, consumer 로그 |
| 중복 결과 | 출력 성공 후 입력 오프셋 커밋 전 실패 여부 |
| 일부 결과 누락 | 비동기 send 오류 처리와 자동 커밋 설정 |
처음에는 각 토픽을 콘솔 소비자로 읽어 위 세 결과가 기대대로 나오는지 확인한다. 다음으로 처리 프로세스를 출력 전송 전후에 강제로 종료해 재시작할 때 어떤 결과가 생기는지 관찰한다. 이 실험이 성공·실패 경계를 이해하는 데 단순 정상 실행보다 도움이 된다.