본문 바로가기
Kafka

커머스 정산 파이프라인에 Kafka + Avro 도입하기 — 5,000건을 한 번에 보내는 법

by Logan Jung 2026. 4. 18.
반응형

커머스 정산 시스템에 Kafka + Avro를 도입하며 설계한 구조와 그 이유를 정리한다.

들어가며

커머스 서비스에서 정산(Settlement)은 결제와는 성격이 다르다. 결제는 수 초 내에 응답이 와야 하는 동기 처리지만, 정산은 보통 하루 단위로 쌓인 대량 데이터를 일괄 처리한다. PG사로부터 받아오는 정산 내역 한 건의 페이로드는 작지만, 5,000건씩 묶이면 수십 MB에 달한다.

이런 대량 배치성 데이터를 REST API로 받아 DB에 곧바로 꽂으면 문제가 많다. 요청 타임아웃, 트랜잭션 길이 폭증, 다른 서비스로의 데이터 전파 어려움, 재처리 복잡성. 이 글에서는 실제 프로젝트에서 정산 파이프라인에 Kafka + Avro + Confluent Schema Registry를 도입한 과정을 공유한다.

다룰 내용은 다음과 같다.

  • 왜 정산에는 Kafka가 어울리는가
  • JSON 대신 Avro를 선택한 이유
  • Producer<T> 포트로 Kafka를 추상화하기
  • 대용량 레코드 전송을 위한 Producer/Consumer 튜닝
  • Schema Registry 연동과 타입 안전한 역직렬화

1. 왜 정산은 Kafka인가

먼저 도메인 흐름부터 보자. 이 프로젝트에서 정산 파이프라인은 다음과 같이 흐른다.

Scheduler / API 호출
  ↓
SettlementsService.send()
  ↓
[1] PG사(토스) 정산 조회 API 호출 (최대 5000건)
  ↓
[2] Avro 레코드로 직렬화
  ↓
[3] Kafka "settlements" 토픽으로 전송
  ↓
[4] KafkaConsumer가 수신
  ↓
[5] 도메인 엔티티 변환 후 bulk insert

PG사 호출과 DB 저장 사이에 Kafka를 끼워넣은 이유는 크게 세 가지다.

① 책임 분리: PG API 호출 실패와 DB 저장 실패를 분리할 수 있다. Consumer가 실패하면 Kafka가 메시지를 보관하므로 재처리 가능.
② 확장성: 정산 데이터를 여러 컨슈머(DB 저장용, 분석용, 알림용)가 독립적으로 소비할 수 있다.
③ 백프레셔 완충: PG 응답을 한 번에 받아도, DB 저장은 컨슈머 처리 속도에 맞춰 진행된다.

실제 코드에서 Producer 쪽 흐름은 이렇게 생겼다.

@Service
@RequiredArgsConstructor
public class SettlementsService implements PaymentSettlementsUseCase, SendSettlementsInfoUseCase {
    private final static String SETTLEMENTS_TOPIC = "settlements";
    private final PaymentAPIs mockTossPayment;
    private final Producer<RPaymentSettlements> producer;

    @SneakyThrows
    @Override
    public boolean send() {
        // 1) PG사에서 정산 내역 조회 (최대 5000건)
        List<ResponsePaymentSettlements> response =
            mockTossPayment.requestPaymentSettlement(createPaymentSettlement());

        // 2) 도메인 → Avro 레코드 변환
        List<PaymentSettlements> datum = response.stream()
            .map(ResponsePaymentSettlements::toEntity)
            .toList();

        RPaymentSettlements record = RPaymentSettlements.newBuilder()
            .setSettlements(datum.stream().map(data -> Settlements.newBuilder()
                .setId(data.getId())
                .setPaymentKey(data.getPaymentKey())
                .setTotalAmount(data.getTotalAmount())
                .setPayOutAmount(data.getPayOutAmount())
                .setCanceledAmount(data.getCanceledAmount())
                .setMethod(data.getMethod().toString())
                .setSoldDate(data.getSoldDate().toString())
                .setPaidOutDate(data.getPaidOutDate().toString())
                .build()
            ).toList())
            .build();

        // 3) Kafka 전송
        producer.send(SETTLEMENTS_TOPIC, record);
        return true;
    }
}

주목할 점은 createPaymentSettlement()어제 기준으로 최근 3일치를 조회한다는 것이다. 정산은 PG사 내부 확정까지 며칠의 지연이 있기 때문에, 실시간이 아니라 T-1 배치로 설계하는 것이 자연스럽다.

private PaymentSettlement createPaymentSettlement() {
    DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd");
    String startDate = LocalDateTime.now(ZoneId.of("Asia/Seoul")).minusDays(3).format(formatter);
    String endDate = LocalDateTime.now(ZoneId.of("Asia/Seoul")).minusDays(1).format(formatter);
    return PaymentSettlement.builder()
        .startDate(startDate).endDate(endDate).page(1).size(5000)
        .build();
}

2. JSON이 아니라 Avro인 이유

Kafka 메시지 포맷으로 JSON을 쓸 수도 있다. 그런데도 Avro를 선택한 이유는 커머스 도메인 특성과 맞닿아 있다.

비교 항목 JSON Avro
직렬화 크기 필드명 포함, 텍스트 기반 바이너리, 필드명 제외 (스키마 별도)
스키마 관리 코드/문서에만 존재 Schema Registry에 중앙 관리
타입 안전성 런타임에 파싱 실패 빌드 타임 코드 생성으로 컴파일 오류
스키마 진화 수동 관리, 깨지기 쉬움 호환성 규칙(BACKWARD, FORWARD 등) 자동 검증

정산 데이터는 5,000건이 한 번에 전송되기 때문에 직렬화 크기가 성능에 큰 영향을 준다. JSON은 필드명을 매 레코드마다 반복 포함시키지만, Avro는 스키마를 분리해 바이너리로 저장하므로 수십 %의 페이로드 절감 효과가 있다.

더 중요한 건 스키마 진화다. 정산 레코드에 새 필드(예: 환불 수수료, 부가세 분리)가 추가되는 일은 비일비재하다. Schema Registry가 호환성 규칙을 강제하면, Producer가 깨뜨릴 수 있는 변경을 사전에 막을 수 있다.

Avro 스키마 정의 (settlements.avsc)

스키마는 JSON 형식의 .avsc 파일로 선언한다.

{
  "namespace": "sw.sustainable.springlabs.fpay.infrastructure.out.mq.record",
  "type": "record",
  "name": "RPaymentSettlements",
  "fields": [
    {
      "name": "settlements",
      "type": {
        "type": "array",
        "items": {
          "name": "Settlements",
          "type": "record",
          "fields": [
            { "name": "id", "type": "int" },
            { "name": "paymentKey", "type": "string" },
            { "name": "method", "type": "string" },
            { "name": "totalAmount", "type": "int" },
            { "name": "canceledAmount", "type": "int" },
            { "name": "payOutAmount", "type": "int" },
            { "name": "soldDate", "type": "string" },
            { "name": "paidOutDate", "type": "string" }
          ]
        }
      }
    }
  ]
}

Gradle에서 Avro 클래스 자동 생성

.avsc로부터 RPaymentSettlements, Settlements 자바 클래스가 빌드 시점에 자동 생성된다.

plugins {
    // AVRO 코드 생성 플러그인
    id("com.github.davidmc24.gradle.plugin.avro") version "1.9.1"
}

dependencies {
    implementation("org.apache.avro:avro:1.11.3")
    implementation("io.confluent:kafka-avro-serializer:7.0.1")
}

tasks.withType<com.github.davidmc24.gradle.plugin.avro.GenerateAvroJavaTask> {
    fieldVisibility = "PRIVATE"
    setCreateSetters("false")
    setSource("src/main/avro")
    setOutputDir(file("build/generated-sources"))
}

생성된 클래스는 Builder 패턴을 제공하므로 Producer 코드에서 RPaymentSettlements.newBuilder().setSettlements(...).build()처럼 타입 안전하게 작성할 수 있다. 스키마가 바뀌면 빌드가 먼저 실패하기 때문에 런타임 장애를 사전에 차단한다.

참고로 이 프로젝트에서는 compileJavagenerateTestAvroJava에 의존하도록 설정되어 있어, 테스트 컴파일 전에 Avro 코드 생성이 항상 선행된다.

tasks.compileJava {
    dependsOn(tasks.generateTestAvroJava)
}

3. Producer 포트 추상화

헥사고날 아키텍처를 따르므로, 도메인이 Kafka에 직접 의존하지 않는다. Producer<T>라는 제네릭 포트로 추상화한다.

// application/port/out/mq/Producer.java
public interface Producer<T> {
    boolean send(String topic, T record);
}

실제 구현은 infrastructure 계층에 위치한 Kafka 어댑터다.

@Component
@RequiredArgsConstructor
@Slf4j
public class KafkaProducer<T> implements Producer<T> {
    private final KafkaTemplate<String, T> kafkaTemplate;

    @Override
    public boolean send(String topic, T record) {
        log.info("sending payload={} to topic={}", record, topic);
        kafkaTemplate.send(topic, record);
        return false;
    }
}

이 구조의 실용적 이점은 테스트할 때 드러난다. 도메인 서비스 테스트에서는 Producer<RPaymentSettlements> Mock을 주입해 Kafka 없이 비즈니스 로직만 검증할 수 있다. 나중에 Kafka를 AWS Kinesis나 SQS로 바꿔도 어댑터만 새로 쓰면 된다.


4. Producer 설정 — Schema Registry와 크기 튜닝

Confluent의 KafkaAvroSerializer를 Value Serializer로 쓰고, Schema Registry URL을 주입한다.

@Configuration
@ConfigurationProperties(prefix = "spring.kafka")
public class KafkaProducerProperties {
    private String[] bootstrapServers;
    private String schemaRegistryUrl;

    @Bean
    public ProducerFactory<String, Object> producerFactory() {
        Map<String, Object> config = new HashMap<>();
        config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, String.join(",", bootstrapServers));
        config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaAvroSerializer.class);
        config.put(KafkaAvroSerializerConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
        config.put(ProducerConfig.MAX_REQUEST_SIZE_CONFIG, 31_457_280);  // 30MB
        return new DefaultKafkaProducerFactory<>(config);
    }

    @Bean
    public KafkaTemplate<String, Object> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

여기서 놓치기 쉬운 핵심이 MAX_REQUEST_SIZE_CONFIG다. Kafka Producer의 기본 max.request.size1MB인데, 5,000건의 정산 레코드는 Avro로 압축해도 그보다 훨씬 크다. 그래서 30MB로 상향 조정했다. 이 설정을 빼먹으면 프로덕션에서 다음과 같은 예외를 만나게 된다.

RecordTooLargeException: The message is 12345678 bytes when serialized
which is larger than the maximum request size you have configured with
the max.request.size configuration.

브로커 쪽에도 같은 의미의 message.max.bytes, 토픽 단위의 max.message.bytes 제한이 있다. Producer만 올려도 브로커가 거절하므로, 인프라 팀과 사전에 맞춰야 한다.


5. Consumer 설정 — Specific Avro Reader의 중요성

Consumer 쪽 설정은 더 재미있다. 여기서 선택 하나가 Consumer 코드의 타입 안전성을 결정한다.

@Configuration
@ConfigurationProperties(prefix = "spring.kafka")
public class KafkaConsumerProperties {
    private String[] bootstrapServers;
    private String groupId;
    private String schemaRegistryUrl;

    @Bean
    public ConsumerFactory<String, Object> consumerFactory() {
        Map<String, Object> config = new HashMap<>();
        config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, String.join(",", bootstrapServers));
        config.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
        config.put(KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
        config.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, true);  // ★
        config.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, 52_428_800);  // 50MB
        return new DefaultKafkaConsumerFactory<>(config);
    }
}

SPECIFIC_AVRO_READER_CONFIG = true가 핵심이다. 이 옵션이 없을 때있을 때의 차이를 보자.

❌ 이 옵션이 없으면 — Consumer는 메시지를 GenericRecord로 받는다.

// Bad: GenericRecord는 키-값 Map처럼 다뤄야 함
@KafkaListener(topics = "settlements")
public void receive(ConsumerRecord<String, GenericRecord> record) {
    GenericRecord payload = record.value();
    Object settlements = payload.get("settlements");  // 타입 캐스팅 지옥
    // 필드명 오타 → 런타임 NPE
}

✅ 옵션을 켜면 — Avro 플러그인이 생성한 구체 클래스로 역직렬화된다.

@KafkaListener(topics = "settlements")
public void receive(ConsumerRecord<String, RPaymentSettlements> consumerRecord) {
    RPaymentSettlements payload = consumerRecord.value();
    List<Settlements> records = payload.getSettlements();  // 타입 안전
    // IDE 자동완성, 컴파일 타임 검증 모두 됨
}

FETCH_MAX_BYTES_CONFIG = 50MB도 꼭 필요하다. Producer가 30MB까지 보낼 수 있으니, Consumer가 한 번에 fetch 가능한 크기는 그보다 커야 한다.


6. Consumer — 받은 데이터를 도메인으로 변환

Consumer의 역할은 Avro 레코드를 도메인 엔티티로 변환한 뒤 bulk insert하는 것이다.

@Component
@RequiredArgsConstructor
@Slf4j
public class KafkaConsumer {
    private final static String SETTLEMENTS_TOPIC = "settlements";
    private final PaymentSettlementsRepository paymentSettlementsRepository;

    @KafkaListener(topics = SETTLEMENTS_TOPIC)
    public void receive(ConsumerRecord<String, RPaymentSettlements> consumerRecord) {
        RPaymentSettlements payload = consumerRecord.value();
        List<Settlements> records = payload.getSettlements();

        List<PaymentSettlements> rows = records.stream()
            .map(record -> PaymentSettlements.builder()
                .paymentKey(record.getPaymentKey())
                .method(PaymentMethod.valueOf(record.getMethod()))
                .paymentStatus(PaymentStatus.valueOf("SETTLEMENTS_REQUESTED"))
                .totalAmount(record.getTotalAmount())
                .payOutAmount(record.getPayOutAmount())
                .canceledAmount(record.getCanceledAmount())
                .soldDate(Date.valueOf(record.getSoldDate()))
                .paidOutDate(Date.valueOf(record.getPaidOutDate()))
                .build())
            .toList();

        paymentSettlementsRepository.bulkInsert(rows);
    }
}

도메인 엔티티로 변환할 때 paymentStatus = SETTLEMENTS_REQUESTED로 지정해 정산 상태의 출발점을 명시한다. 이후 실제 지급이 완료되면 SETTLEMENTS_COMPLETED로 전이된다.

JPA batch-size는 application.yml에서 100으로 설정되어 있어, 5,000건을 한 번에 insert해도 100건씩 묶여 나간다.

spring:
  jdbc:
    template:
      batch-size: 100

7. 환경별 Kafka 설정 분리

이 프로젝트는 aws, dev, dev-docker, test 네 개 프로파일을 운영한다. Kafka 엔드포인트만 환경별로 갈아끼운다.

# dev 프로파일 (로컬 개발)
spring:
  kafka:
    bootstrap-servers: [ "localhost:19094" ]
    group-id: "settlements_1"
    schema-registry-url: "http://localhost:8085"

---
# dev-docker 프로파일 (Docker Compose)
spring:
  kafka:
    bootstrap-servers: [ "kafka-1:19092" ]
    group-id: "settlements_1"
    schema-registry-url: "http://172.26.0.13:8085"

---
# aws 프로파일 (운영)
spring:
  kafka:
    bootstrap-servers: [ "172.31.5.117:19094" ]
    group-id: "settlements_1"
    schema-registry-url: "http://172.31.5.117:8085"

Docker Compose 환경에서는 컨테이너 네트워크의 서비스명(kafka-1)을 쓰고, 운영에서는 AWS 내부 IP를 사용한다. 코드는 건드리지 않고 Profile 활성화만으로 환경이 바뀐다.


마치며

정산 파이프라인에 Kafka + Avro를 도입하며 얻은 실전 교훈을 정리한다.

  1. 도메인 특성을 먼저 봐라: 정산은 대량 배치 + T-1 지연 + 여러 소비자가 필요한 도메인이다. 이 세 가지 조건 중 하나라도 해당되면 Kafka 도입 가치가 있다.
  2. Avro는 커머스처럼 장기 운영되는 도메인에서 빛난다: 스키마 진화가 불가피한 환경에서 Schema Registry의 호환성 강제가 프로덕션 장애를 미리 차단해 준다.
  3. 크기 설정을 미리 계산해라: Producer max.request.size, Consumer fetch.max.bytes, Broker message.max.bytes, Topic max.message.bytes. 이 네 곳이 일관되게 설정되지 않으면 어느 지점에서든 RecordTooLargeException을 만난다.
  4. Specific Avro Reader를 잊지 말자: SPECIFIC_AVRO_READER_CONFIG=true 한 줄이 Consumer 코드의 타입 안전성을 좌우한다.
  5. 포트로 감싸라: Producer<T>처럼 제네릭 포트로 추상화하면 테스트도 편하고, 향후 메시지 브로커 교체도 유연하다.

Kafka는 "일단 도입"하기보다 "왜 필요한가"부터 답하는 게 중요하다. 정산 도메인은 그 답을 자연스럽게 내놓는 케이스였다.


참고

 

반응형