카프카 설정
config
package com.inflearn.kafka.config;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.LongDeserializer;
import org.apache.kafka.common.serialization.LongSerializer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.*;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class KafkaConfig {
@Bean
public KafkaTemplate<String, String> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
@Bean
public ProducerFactory<String, String> producerFactory() {
Map<String, Object> config = new HashMap<>();
config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); // Long 직렬화
return new DefaultKafkaProducerFactory<>(config);
}
@Bean
public ConsumerFactory<String, Long> consumerFactory() {
Map<String, Object> config = new HashMap<>();
config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
config.put(ConsumerConfig.GROUP_ID_CONFIG, "group_1");
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class); // Long 역직렬화
return new DefaultKafkaConsumerFactory<>(config);
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, Long> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, Long> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
return factory;
}
}
프로듀서
package com.inflearn.kafka.repository;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
@Slf4j
@Component
@RequiredArgsConstructor
public class CouponCreateProducer {
private static final String TOPIC_NAME = "create-coupon";
private final KafkaTemplate<String, Long> kafkaTemplate;
public void create(Long userId) {
log.info("메세지 발행 ");
try {
kafkaTemplate.send(TOPIC_NAME, userId);
log.info("메세지 발행 완료");
} catch (Exception e) {
log.error("메시지 발행 실패 - userId: {}", userId, e);
throw new RuntimeException("Kafka 메시지 발행 실패", e);
}
}
}
컨슈머
package com.inflearn.kafka.repository;
import com.inflearn.kafka.domain.Coupon;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.context.event.ApplicationStartedEvent;
import org.springframework.context.event.EventListener;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
@Component
@RequiredArgsConstructor
@Slf4j
public class CouponCreatedConsumer {
private final CouponRepository couponRepository;
@KafkaListener(topics = "create-coupon", groupId = "group_1")
public void listener(Long userId) {
try {
log.info("쿠폰 생성 시작 - userId: {}", userId);
couponRepository.save(new Coupon(userId));
log.info("쿠폰 생성 완료 - userId: {}", userId);
} catch (Exception e) {
log.error("쿠폰 생성 실패 - userId: {}", userId, e);
// 재시도 정책에 따라 예외를 던질지 말지 결정
throw new RuntimeException("쿠폰 생성 실패", e);
}
}
}
docker-compose.yaml
version: '2'
services:
zookeeper:
image: confluentinc/cp-zookeeper:latest
container_name: zookeeper
environment:
ZOOKEEPER_SERVER_ID: 1
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
ZOOKEEPER_INIT_LIMIT: 5
ZOOKEEPER_SYNC_LIMIT: 2
ports:
- "2181:2181"
kafka:
image: confluentinc/cp-kafka:latest
container_name: kafka
depends_on:
- zookeeper
ports:
- "9092:9092"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:29092,PLAINTEXT_HOST://localhost:9092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
redis:
image: redis:latest
container_name: redis
ports:
- "6379:6379"
volumes:
- redis_data:/data
volumes:
redis_data:
토픽 생성
docker exec -it kafka /bin/kafka-topics --bootstrap-server localhost:9092 --create --topic create-coupon --partitions 1 --replication-factor 1
프로듀서 (메세지전송)
docker exec -it kafka kafka-console-producer --bootstrap-server kafka:9092 --topic create-coupon
컨슈머 (메세지 수신)
docker exec -it kafka kafka-console-consumer --bootstrap-server kafka:9092 --topic create-coupon --from-beginning
package com.inflearn.kafka.repository;
import lombok.RequiredArgsConstructor;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Repository;
@Repository
@RequiredArgsConstructor
public class CouponCountRepository {
private final RedisTemplate<String, String> redisTemplate;
public Long increment() {
return redisTemplate.opsForValue().increment("coupon-count");
}
public void reset() {
redisTemplate.opsForValue().set("coupon-count", "0");
}
public Long isCouponReceived(Long userId) {
return redisTemplate.opsForSet().add("user", userId.toString());
}
}
서비스 로직
package com.inflearn.kafka.service;
import com.inflearn.kafka.repository.CouponCreateProducer;
import com.inflearn.kafka.repository.CouponCountRepository;
import com.inflearn.kafka.repository.CouponRepository;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@Service
@Slf4j
@RequiredArgsConstructor
public class ApplyService {
private final CouponRepository couponRepository;
private final CouponCreateProducer couponCreateProducer;
private final CouponCountRepository couponCountRepository;
public void couponIssued(Long userId) {
log.info("쿠폰 발급 서비스 시작");
Long check = couponCountRepository.isCouponReceived(userId);
if (check != 1) {
return;
}
Long increment = couponCountRepository.increment();
log.info("increment ={}", increment);
if (increment > 100) {
return;
}
couponCreateProducer.create(userId);
}
}
이슈 1
(카프카 ui 사용으로 해결하는것이 더 좋을듯 하다.)
메세지 수신값을 콘솔 로그인에서 확인할때 Long값이 아니여서 터미널창에서 보이지 않음
프로듀서에서 String값으로 보낸 후 컨슈머에서 String값으로 메세지 수신으로 코드 수정
그후에 다시 기존설정값 Long으로 변경했지만 컨슈머에서 받지 못함
다른값이 저장된후 다시 변경하려면 토픽을 삭제하고 다시 설정해야함
docker exec -it kafka /bin/kafka-topics --bootstrap-server localhost:9092 --delete --topic create-coupon
docker exec -it kafka /bin/kafka-topics --bootstrap-server localhost:9092 --list
리스트 조회 시 아직 삭제안되있으면 연결되있는 프로듀서와 컨슈머를 종료시켜야한다.
docker exec -it kafka /bin/kafka-topics --bootstrap-server localhost:9092 --create --topic create-coupon --partitions 1 --replication-factor 1
나중에 보니 애초에 컨슈머 콘솔창 띄울때 long설정을 해주면 가능
이슈 2
@Test
void t2() throws InterruptedException {
int threadCount = 100;
ExecutorService executorService = Executors.newFixedThreadPool(32);
CountDownLatch countDownLatch = new CountDownLatch(threadCount);
for (int i = 0; i < threadCount; i++) {
Long userId = (long) i;
executorService.execute(() -> {
try {
applyService.couponIssued(userId);
} finally {
countDownLatch.countDown();
}
});
}
countDownLatch.await();
Thread.sleep(5000);
couponCountRepository.reset();
long count = couponRepository.count();
assertThat(count).isEqualTo(100);
}
sleep을 5초 준 이유
카프카 -> 데이터베이스까지 흘러가는 시간을 고려해서 sleep 설정
밑에 K6사용한 부하테스트 지표를 보면 알수있듯이 throughput = 28.4 이다.
테스트코드는 0.8초 안에 종료

k6 스크립트
import http from 'k6/http';
import { check, sleep } from 'k6';
export let options = {
vus: 30,
duration: '10s',
};
export default function () {
const url = 'http://localhost:8080/api/coupon';
const payload = JSON.stringify({
userId: 1
});
const params = {
headers: {
'Content-Type': 'application/json',
},
};
const res = http.post(url, payload, params);
// 응답 검사
check(res, {
'status is 200': (r) => r.status === 200,
});
sleep(1);
}
