Загружаем каталог…
Загружаем каталог…
Kafka 전용 코드 리뷰: 주문 event가 수집되기까지 1. 메시지의 길 POST /api/orders → OrderService → MySQL 주문 + outbox (같은 DB 트랜잭션) OutboxScheduler → OutboxPublisher → KafkaOrderEventSender → orders.paid 토픽 orders.paid 토픽 → AnalyticsConsumer → AnalyticsService → MySQL 수집 기록 토픽 은 Kafka의 메시지 이름표가 붙은 저장 공간이다. 파티션 은 한 토픽을 나눈 줄이다. 프로듀서 는 메시지를 보내고 컨슈머 는 읽는다. 이 코드에서는 OrderEvent 한 개가 주문 완료 사건 한 개를 나타낸다. 2. 주문과 Kafka를 한 번에 직접 처리하지 않는 이유 MySQL과 Kafka는 서로 다른 시스템이다. 주문을 DB에 저장한 직후 서버가 꺼지면 Kafka 전송이 빠질 수 있다. 그래서 OrderService 가 주문과 outbox 행을 같은 MySQL 트랜잭션 에 저장한다. 예약 작업이 나중에 outbox를 읽어 Kafka로 보낸다. Kafka 확인을 받은 뒤에만 대기 행을 지운다. 서버가 멈춘 시점 남는 것 다시 시작하면 DB 트랜잭션이 끝나기 전 주문과 outbox가 함께 취소됨 새 요청으로 다시 시도 DB 커밋 후 Kafka 전송 전 주문과 outbox가 남음 예약 작업이 전송 Kafka 확인 후 outbox 삭제 전 Kafka 메시지와 outbox가 모두 남을 수 있음 다시 전송될 수 있어 수신자가 event_id 로 중복 저장을 막음 이 방식은 최소 한 번 전달 이다. 결제 요청을 한 번만 처리하는 Idempotency-Key , Kafka 프로듀서 설정의 enable.idempotence , 수집 표의 event_id 중복 방지는 각각 다른 구간의 중복을 다룬다. 3. 파티션 3개와 consumer 3개를 읽는 법 OrderTopicConfig 는 orders.paid 토픽에 파티션 3개를 요청한다. AnalyticsConsumer 의 concurrency = "3" 은 애플리케이션 인스턴스 한 대에서 최대 세 소비자를 만든다는 뜻이다. 여러 인스턴스가 모두 coffee-analytics 그룹을 사용하면 파티션 3개를 나눠 읽는다. 이 그룹에서 동시에 유효한 파티션 담당자는 최대 3개다. 이 로컬 구성은 Kafka 브로커가 한 대이고 복제본도 한 개라 브로커 장애에 대비한 구성은 아니다. 4. 코드와 줄 바로 위 설명 1. build.gradle // 플러그인은 Gradle에 기능을 더합니다. Java 컴파일과 Spring Boot 실행 기능을 여기서 켭니다. plugins { id 'java' // Spring Boot는 웹 서버를 쉽게 만드는 Java 도구입니다. 버전을 적어 같은 환경으로 빌드합니다. id 'org.springframework.boot' version '3.5.7' id 'io.spring.dependency-management' version '1.1.7' } group = 'com.example' version = '0.0.1-SNAPSHOT' java { // 이 프로젝트가 Java 17 문법을 사용한다는 뜻입니다. sourceCompatibility = JavaVersion.VERSION_17 } repositories { // 필요한 라이브러리를 내려받을 공개 저장소입니다. mavenCentral() } dependencies { // HTTP 요청을 받고 JSON 응답을 보내는 기능입니다. implementation 'org.springframework.boot:spring-boot-starter-web' // Java 코드에서 SQL을 실행하는 기능입니다. 잔액을 바꾸는 SQL을 직접 쓰기 위해 선택했습니다. implementation 'org.springframework.boot:spring-boot-starter-jdbc' // 사용자 ID와 금액이 올바른지 @Min, @Max 등으로 검사하는 기능입니다. implementation 'org.springframework.boot:spring-boot-starter-validation' implementation 'org.springframework.boot:spring-boot-starter-cache' // Redis에 연결합니다. Redis는 빠른 조회용 저장소이고 결제의 최종 기록은 MySQL에 둡니다. implementation 'org.springframework.boot:spring-boot-starter-data-redis' // Kafka로 주문 완료 메시지를 보내고 받는 기능입니다. implementation 'org.springframework.kafka:spring-kafka' // Flyway는 SQL 파일을 번호 순서대로 실행해 DB 테이블을 만듭니다. implementation 'org.flywaydb:flyway-core' implementation 'org.flywaydb:flyway-mysql' // 실행할 때 MySQL에 접속하는 드라이버가 필요합니다. runtimeOnly 'com.mysql:mysql-connector-j' compileOnly 'org.projectlombok:lombok' // Lombok이 컴파일 중 builder() 코드를 자동 생성합니다. 코드에 직접 적지 않아도 호출할 수 있습니다. annotationProcessor 'org.projectlombok:lombok' testImplementation 'org.springframework.boot:spring-boot-starter-test' // 자동 테스트에서는 설치가 간단한 메모리 DB인 H2를 사용합니다. testRuntimeOnly 'com.h2database:h2' } tasks.withType(JavaCompile).configureEach { // 더 최신 JDK로 빌드하더라도 Java 17에 없는 기능을 잘못 쓰지 않도록 막습니다. options.release = 17 } tasks.named('test') { // JUnit은 자동 테스트를 실행하는 도구입니다. 이 설정이 JUnit 5 테스트를 켭니다. useJUnitPlatform() } 2. compose.yaml # Docker Compose가 함께 실행할 프로그램들을 나열합니다. 여기서는 MySQL, Kafka, Redis입니다. services: mysql: # MySQL 8.4를 사용합니다. MySQL에는 잔액, 주문, 발행 대기 메시지를 저장합니다. image: mysql:8.4 environment: MYSQL_DATABASE: coffee MYSQL_USER: coffee # 로컬 연습용 비밀번호입니다. 실제 서비스에서는 코드에 비밀번호를 그대로 두면 안 됩니다. MYSQL_PASSWORD: coffee MYSQL_ROOT_PASSWORD: root ports: # 왼쪽 13306은 내 컴퓨터의 포트, 오른쪽 3306은 컨테이너 안의 포트입니다. - "13306:3306" volumes: # DB 파일을 Docker 볼륨에 보관합니다. 컨테이너를 다시 만들어도 실습 데이터가 남을 수 있습니다. - mysql_data:/var/lib/mysql healthcheck: test: ["CMD-SHELL", "mysqladmin ping -h localhost -u root -proot"] interval: 3s timeout: 3s retries: 10 kafka: # Kafka는 주문 완료 사건을 다른 프로그램에 전달합니다. 이 파일은 브로커 한 대만 실행합니다. image: apache/kafka:4.1.2 ports: - "9092:9092" redis: # Redis는 메뉴 캐시와 최근 주문 횟수 ZSET을 저장합니다. Redis 데이터가 없어져도 MySQL에서 다시 만들 수 있습니다. image: redis:7.4.11-alpine3.21 ports: - "6379:6379" healthcheck: # healthcheck는 Redis가 명령에 응답할 준비가 됐는지 확인합니다. test: ["CMD", "redis-cli", "ping"] interval: 3s timeout: 3s retries: 10 volumes: mysql_data: 3. src/main/resources/application.yml # Spring Boot의 실행 설정입니다. 들여쓰기로 어떤 설정이 어느 기능에 속하는지 구분합니다. spring: datasource: # ${DB_URL:기본값}은 환경변수 DB_URL이 있으면 그 값을 쓰고, 없으면 기본값을 쓴다는 뜻입니다. # 시간을 UTC로 맞춰 서버마다 주문 시각과 7일 계산이 달라지지 않게 합니다. url: ${DB_URL:jdbc:mysql://localhost:13306/coffee?connectionTimeZone=UTC&forceConnectionTimeZoneToSession=true} username: ${DB_USER:coffee} password: ${DB_PASSWORD:coffee} # 앱이 시작할 때 번호가 붙은 DB 변경 SQL을 실행합니다. flyway: enabled: true kafka: # Kafka가 실행 중인 주소입니다. bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092} consumer: key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer # 이 소비자 그룹에 저장된 읽기 위치가 없을 때만 토픽 앞에서부터 읽습니다. # 처음 소비하는 그룹이라 읽은 위치 기록이 없으면 가장 오래된 메시지부터 읽습니다. auto-offset-reset: earliest properties: # JSON의 타입 정보를 Java 객체로 바꿀 때 허용하는 패키지를 좁힙니다. # 주문 이벤트 DTO가 옮겨진 새 패키지에서 Kafka JSON 객체를 읽습니다. # JSON 메시지를 Java 객체로 바꿀 때 허용할 코드 패키지를 제한합니다. spring.json.trusted.packages: com.example.coffee.order.dto producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.springframework.kafka.support.serializer.JsonSerializer # 설정된 복제본들이 기록을 확인해야 전송 성공으로 간주합니다. 이 실습은 복제본이 1개입니다. # Kafka가 메시지를 받았다고 확인할 때까지 기다리도록 합니다. 브로커가 한 대이므로 복제 내구성은 별개입니다. acks: all properties: # 프로듀서의 재전송 때문에 같은 레코드가 중복 기록되는 일을 줄입니다. 결제 요청 멱등성과는 별개입니다. # Kafka 생산자 내부의 재시도 중복을 줄입니다. HTTP 결제 요청 중복 방지와는 다른 기능입니다. enable.idempotence: true data: redis: # Redis 주소도 환경변수로 바꿀 수 있습니다. host: ${REDIS_HOST:localhost} port: ${REDIS_PORT:6379} connect-timeout: 2s timeout: 2s cache: type: redis # 이전 메뉴 객체를 담은 Redis 캐시와 새 DTO 패키지의 캐시를 분리합니다. cache-names: menus-v2 redis: # 캐시 값을 5분 뒤 만료시킵니다. 만료되면 다음 조회에서 MySQL을 다시 읽습니다. time-to-live: 5m orders: # 주문 완료 사건이 지나가는 Kafka 토픽의 이름입니다. # 주문 완료 메시지를 보낼 Kafka 토픽 이름입니다. 토픽은 메시지를 담는 통로입니다. topic: orders.paid # outbox 발행과 Redis 집계를 실행할 예약 작업 스레드를 둘로 설정합니다. spring.task.scheduling.pool.size: 2 4. src/main/resources/db/migration/V1__init.sql -- 메뉴 표를 만듭니다. 한 행이 커피 메뉴 하나입니다. CREATE TABLE menus ( id BIGINT AUTO_INCREMENT PRIMARY KEY, name VARCHAR(100) NOT NULL UNIQUE, -- 가격을 정수 원 단위로 저장합니다. 소수점 반올림 오류를 피하기 쉽습니다. price BIGINT NOT NULL CHECK (price > 0) ); -- 사용자 한 명의 현재 포인트 잔액을 저장합니다. CREATE TABLE point_accounts ( user_id BIGINT PRIMARY KEY, -- CHECK는 잔액이 음수가 되거나 정한 최대값을 넘지 못하게 DB에서도 검사합니다. balance BIGINT NOT NULL DEFAULT 0 CHECK (balance >= 0 AND balance <= 1000000000000) ); -- 결제가 끝난 주문을 저장합니다. 메뉴 가격이 나중에 변해도 당시 결제액을 기억해야 합니다. CREATE TABLE orders ( id BIGINT AUTO_INCREMENT PRIMARY KEY, -- REFERENCES는 해당 사용자 계정이 실제로 존재해야 한다는 외래키 규칙입니다. user_id BIGINT NOT NULL REFERENCES point_accounts(user_id), menu_id BIGINT NOT NULL REFERENCES menus(id), paid_amount BIGINT NOT NULL CHECK (paid_amount > 0), ordered_at TIMESTAMP(6) NOT NULL ); -- 인덱스는 책의 색인처럼 최근 주문을 찾는 작업을 돕습니다. 대신 저장 공간과 쓰기 비용이 듭니다. CREATE INDEX idx_orders_recent ON orders(ordered_at, menu_id); -- Kafka에 아직 보내지 않았거나 확인받지 못한 주문 사건을 잠시 보관합니다. CREATE TABLE order_outbox ( id BIGINT AUTO_INCREMENT PRIMARY KEY, -- 주문 하나에 발행 대기 기록이 중복 생성되지 않게 UNIQUE를 둡니다. order_id BIGINT NOT NULL UNIQUE REFERENCES orders(id), user_id BIGINT NOT NULL, menu_id BIGINT NOT NULL, paid_amount BIGINT NOT NULL, attempts INTEGER NOT NULL DEFAULT 0, next_attempt_at TIMESTAMP(6) NOT NULL ); -- 다음 전송 시간이 된 outbox 행을 빨리 찾기 위한 색인입니다. CREATE INDEX idx_order_outbox_ready ON order_outbox(next_attempt_at, id); -- 처음 실행할 때 사용할 커피 메뉴 네 개를 넣습니다. INSERT INTO menus(name, price) VALUES ('Americano', 4500), ('Cafe Latte', 5000), ('Cappuccino', 5500), ('Cold Brew', 6000); 5. src/main/resources/db/migration/V2__collected_order_events.sql -- Kafka 메시지를 받은 쪽을 흉내 낸 표입니다. 실제 분석 플랫폼의 역할을 실습용 DB로 표현했습니다. CREATE TABLE collected_order_events ( -- 같은 사건이 Kafka에서 두 번 오더라도 같은 event_id는 한 행만 저장합니다. event_id BIGINT PRIMARY KEY, order_id BIGINT NOT NULL, user_id BIGINT NOT NULL, menu_id BIGINT NOT NULL, paid_amount BIGINT NOT NULL, -- 메시지를 받은 시각입니다. 커피가 결제된 시각과는 다릅니다. collected_at TIMESTAMP(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6) ); 6. src/main/resources/db/migration/V3__order_idempotency.sql -- 클라이언트가 붙인 결제 시도 이름을 저장합니다. 이미 있던 주문은 이 값이 없어서 NULL을 허용합니다. ALTER TABLE orders ADD COLUMN request_key VARCHAR(128); -- 같은 사용자와 같은 요청 키로 주문을 두 개 만들지 못하게 DB가 막습니다. 서버가 여러 대여도 이 규칙은 공유됩니다. CREATE UNIQUE INDEX uq_orders_user_request_key ON orders(user_id, request_key); 7. src/main/java/com/example/coffee/config/kafka/OrderTopicConfig.java // 이 파일은 config 기능의 kafka 패키지에 속합니다. package com.example.coffee.config.kafka; import org.apache.kafka.clients.admin.NewTopic; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.TopicBuilder; // Spring 설정을 담은 클래스입니다. @Configuration public class OrderTopicConfig { @Bean NewTopic orderPaidTopic(@Value("${orders.topic}") String topic) { // 로컬 Kafka 브로커가 한 대여서 복제본도 한 개입니다. 브로커 장애를 견디는 구성은 아닙니다. // 메시지 저장 줄을 세 개로 나눠 한 그룹에서 최대 세 소비자가 병렬 처리하게 합니다. // 하나의 Kafka 토픽을 세 갈래로 나눠 소비자 세 개가 동시에 처리할 수 있게 합니다. // 브로커가 한 대인 연습 환경이라 복제본은 하나입니다. 장애에 안전한 운영 구성을 뜻하지는 않습니다. return TopicBuilder.name(topic).partitions(3).replicas(1).build(); } } 8. src/main/java/com/example/coffee/order/dto/OrderRequest.java // 이 파일은 order 기능의 dto 패키지에 속합니다. package com.example.coffee.order.dto; import jakarta.validation.constraints.Min; import lombok.Builder; @Builder // 주문할 사람과 메뉴를 담습니다. 금액은 보내지 않고 서버가 메뉴 가격을 조회합니다. public record OrderRequest(@Min(1) long userId, @Min(1) long menuId) {} 9. src/main/java/com/example/coffee/order/dto/OrderEvent.java // 이 파일은 order 기능의 dto 패키지에 속합니다. package com.example.coffee.order.dto; import lombok.Builder; @Builder // Kafka로 보낼 사건입니다. eventId는 같은 메시지가 다시 왔는지 확인하는 번호입니다. public record OrderEvent(long eventId, long orderId, long userId, long menuId, long paidAmount) {} 10. src/main/java/com/ex
То, что RADAR обнаружил и классифицировал для этой возможности. Это опубликованный источником текст, а не подтверждение, что предложение ещё действует.
Kafka 전용 코드 리뷰: 커피주문 event가 수집 되기까지. Kafka 전용 코드 리뷰: 주문 event가 수집되기까지 1. 메시지의 길 POST /api/orders → OrderService → MySQL 주문 + outbox (같은 DB 트랜잭션) OutboxScheduler → OutboxPublisher → KafkaOrderEventSender → orders.paid 토픽 orders.paid 토픽 → AnalyticsConsumer → AnalyticsService → MySQL 수집 기록 토픽 은 Kafka의 메시지 이름표가 붙은 저장 공간이다. 파티션 은 한 토픽을 나눈 줄이다. 프로듀서 는 메시지를 보내고 컨슈머 는 읽는다. 이 코드에서는 OrderEvent 한 개가 주문 완료…
Открыть источник