목차
kafka에 Storm 연결하기
-
KafkaSpout : 카프카 데이터 스트림으로 사용한 다음 추가 처리를 위해 데이터를 볼트에 전달한다. 그리고 주키퍼의 정보, 카프카 브로커, 연결할 토픽 등을 SpoutConfig를 통해 설정한다.
-
스파우트는 카프카 컨슈머 역할을 하므로, 목적지를 가리키는 레코드 오프셋 관리가 필요하기 때문에 주키퍼를 사용하며, SpoutConfig의 두 개의 파라미터는 주키퍼의 루트 디렉토리 경로와 특정 스파우트의 ID를 표시한다.
## offset
zkRootDir/consumerID/0
zkRootDir/consumerID/1
zkRootDir/consumerID/2
- SchemeAsMultiScheme : 카프카에서 ByteBuffer가 스톰 튜플로 변환이 사용됐는지 나타낸다.
이 글의 SpoutConfig와 ZooKeeper 오프셋 경로는 당시 사용한 Storm Kafka 연동 방식을 설명한다. Storm과 Kafka 연동 API는 버전에 따라 다르므로, 새 환경에서 이 설정을 그대로 복사하기 전에 사용 중인 Storm 커넥터의 문서를 확인해야 한다.
flowchart LR 생산자 --> 토픽[Kafka 토픽] 토픽 --> 스파우트[KafkaSpout: 소비·튜플 생성] 스파우트 --> 볼트[Storm Bolt: 처리] 볼트 --> 출력[저장 또는 후속 처리]
스파우트는 Kafka 레코드를 Storm 튜플로 바꾸고 볼트에 전달한다. 여기서 확인할 항목은 어느 토픽과 파티션을 읽는지, 오프셋을 언제 확정하는지, 처리 실패 시 재시도하는지다. 재시도 경로가 있으면 같은 레코드가 다시 처리될 수 있으므로 볼트의 쓰기는 중복에 견딜 수 있게 설계한다.
마지막 메모의 SchemeAsMultiScheme은 바이트 데이터를 튜플로 변환하는 스킴을 연결하는 역할이다. 문자열이나 JSON 같은 실제 레코드 형식에 맞는 디코딩 규칙을 적용해야 필드가 예상대로 전달된다.
처리 중 실패하면
예를 들어 스파우트가 주문 이벤트를 읽고 볼트가 집계 결과를 저장한다고 하자. 집계는 저장됐지만 처리 완료를 기록하기 전에 프로세스가 종료되면, 재시작 후 같은 이벤트를 다시 읽을 수 있다. 그러면 단순히 총합 += 금액을 한 번 더 실행하는 구현은 중복 계산한다.
Kafka 레코드 읽기 → Storm 튜플 생성 → 볼트 처리 → 결과 저장 → 처리 완료 기록
↑ 이 사이에 실패하면 재처리 가능
이 경우 이벤트 ID를 결과 저장소의 고유 키로 사용하거나 이미 처리한 ID를 확인하는 등, 재처리해도 결과가 한 번 처리한 것과 같도록 설계할 수 있다. 어떤 지점에서 오프셋이나 완료 신호가 확정되는지는 사용한 커넥터와 Storm 설정에 달려 있다. 배포 전에 장애를 의도적으로 발생시켜 재시도, 중복, 순서가 실제로 어떻게 나타나는지 확인하는 편이 좋다.
현재 커넥터에서 바뀐 연결 정보
원문의 SpoutConfig와 ZooKeeper 경로는 예전 storm-kafka 방식의 학습 자료다. 현재 공식 Storm Kafka 연동 문서는 storm-kafka-client의 KafkaSpoutConfig를 설명한다. 브로커 연결에는 bootstrap.servers에 해당하는 주소를 넣고, 소비할 토픽을 지정한다. 최신 커넥터의 소비 오프셋은 Kafka 소비자 API와 처리 보장 설정을 기준으로 다루며, 예전의 zkRootDir/consumerID/0 경로를 새 배포에 만들 필요가 없다.
TopologyBuilder topology = new TopologyBuilder();
topology.setSpout(
"kafka-spout",
new KafkaSpout<>(KafkaSpoutConfig.builder("broker:9092", "orders").build()),
1
);
topology.setBolt("validate", new OrderValidationBolt(), 2)
.shuffleGrouping("kafka-spout");
이 코드는 Spout→Bolt 연결 형태를 보여 주는 발췌로, OrderValidationBolt 구현과 imports는 프로젝트에서 제공해야 한다. KafkaSpout의 기본 레코드 변환기는 토픽, 파티션, 오프셋, 키, 값을 튜플 필드로 내보낸다. Bolt가 특정 필드를 받는다고 가정했다면 필드 이름과 역직렬화기 설정을 맞춰야 한다. 같은 Kafka 토픽의 파티션 수와 Spout의 병렬도를 함께 확인한다. Spout task를 더 늘려도 읽을 파티션이 부족하면 병렬 처리량이 비례해 늘지 않는다.
flowchart LR K[Kafka orders 토픽] --> S[KafkaSpout: 레코드와 오프셋] S --> B[검증 Bolt] B --> D[DB·다음 스트림] B -->|ack/fail| S
오프셋과 업무 결과 사이의 간격
Spout가 Kafka 레코드를 읽어 Bolt에 전달하고, Bolt가 DB에 결과를 쓴 뒤 ACK를 보낸다고 하자. DB 쓰기는 성공했지만 ACK 전에 worker가 중지되면 레코드가 재전달돼 중복이 생길 수 있다. “Storm에서 ack했다”는 것과 “Kafka 오프셋이 어떤 시점에 커밋됐다”도 커넥터의 처리 보장 설정에 따라 다르다. 외부 저장소에 주문 ID 등의 고유 키를 사용해 중복 처리를 막는다. 반대로 Bolt가 예외를 잡고 성공 ACK를 해 버리면 실패 레코드를 다시 볼 기회를 잃는다.
| 증상 | 확인할 지점 |
|---|---|
| Spout가 데이터를 못 읽음 | 브로커 광고 주소, 토픽·파티션, 인증, 시작 오프셋 |
| Bolt에 값이 비어 옴 | RecordTranslator 출력 필드와 Bolt의 필드 참조 |
| 재시작 뒤 중복 처리 | 실패 시 재생, ACK 시점과 DB 멱등성 |
| 일부 파티션만 읽음 | Spout task 수, 파티션 할당과 커넥터 로그 |
SchemeAsMultiScheme처럼 바이트를 튜플로 바꾸던 예전 API는 사용한 Storm 버전과 함께 읽어야 한다. 새 코드에서는 레코드의 key/value deserializer와 RecordTranslator가 같은 역할의 일부를 담당한다. 단순 Kafka→Kafka 변환이라면 Storm 토폴로지를 운영할 필요가 있는지 Kafka Streams·Connect와도 비교해 선택한다.
