Step 02 — KafkaTemplate 과 프로듀서

학습 목표

  • KafkaTemplatesend 오버로드 5종을 구분하고, 상황별로 어느 것을 쓸지 판단한다
  • Spring Kafka 3.0 에서 바뀐 반환 타입 CompletableFuture<SendResult<K,V>>whenComplete/exceptionally 로 처리한다
  • 브로커를 내린 채 send() 를 계속 호출해, 비즈니스 로직이 성공한 줄 아는 조용한 유실을 재현한다
  • 10,000건 발행을 4가지 방식으로 측정해 send().get() 이 처리량을 20배 떨어뜨리는 것을 실측한다
  • ProducerFactory 를 직접 만들어 설정이 다른 KafkaTemplate 두 개를 공존시킨다
  • 키 → 파티션 매핑을 murmur2 로 직접 계산해 예측하고, 실제 발행 결과와 대조한다
  • acks / retries / enable.idempotence 의 상호 제약을 기동 실패와 조용한 비활성화 두 경우로 나눠 확인한다

선행 스텝: Step 01 — 환경 구축과 첫 메시지 예상 소요: 90분


2-0. 실습 준비

Step 01 의 프로젝트를 그대로 씁니다. 코드는 com.example.order.step02 패키지, 프로필은 step02 입니다.

docker compose up -d
docker exec -it learn-kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders

결과

Topic: orders	TopicId: xJ0kQm8CQ3ay6z4rN2ZbSA	PartitionCount: 3	ReplicationFactor: 1	Configs: segment.bytes=1073741824
	Topic: orders	Partition: 0	Leader: 1	Replicas: 1	Isr: 1
	Topic: orders	Partition: 1	Leader: 1	Replicas: 1	Isr: 1
	Topic: orders	Partition: 2	Leader: 1	Replicas: 1	Isr: 1

파티션이 3개인 것을 반드시 확인하세요. 2-7 의 키 → 파티션 매핑 표는 파티션 수가 3일 때의 값입니다. 다르면 표가 전부 어긋납니다.

이 스텝은 컨슈머를 쓰지 않습니다. 발행한 것이 정말 갔는지는 콘솔 컨슈머로 확인하니, 터미널을 하나 더 띄워 두세요.

docker exec -it learn-kafka /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server localhost:9092 --topic orders \
  --property print.key=true --property print.partition=true --property print.offset=true

2-1. send 오버로드 — 5가지를 구분한다

겉보기엔 비슷하지만 무엇을 여러분이 정하고 무엇을 Kafka 에 맡기는가가 다릅니다.

kafkaTemplate.send("orders", event);                            // ① 키 없음
kafkaTemplate.send("orders", event.orderId(), event);           // ② 키 지정 ★기본형
kafkaTemplate.send("orders", 1, event.orderId(), event);        // ③ 파티션까지 지정

// ④ ProducerRecord — 헤더·타임스탬프까지 전부 제어
ProducerRecord<String, OrderCreated> record = new ProducerRecord<>(
        "orders", null, event.createdAt().toEpochMilli(), event.orderId(), event);
record.headers().add("trace-id", "t-0001".getBytes(StandardCharsets.UTF_8));
kafkaTemplate.send(record);

// ⑤ Spring Messaging 의 Message<?>
kafkaTemplate.send(MessageBuilder.withPayload(event)
        .setHeader(KafkaHeaders.TOPIC, "orders")
        .setHeader(KafkaHeaders.KEY, event.orderId()).build());
오버로드파티션 결정헤더타임스탬프언제 쓰나
send(topic, data)Kafka 의 파티셔너 (키 없음)불가send() 호출 시각순서가 상관없는 로그·지표성 이벤트
send(topic, key, data)murmur2(key) % N불가동일기본형. 엔티티 단위 순서를 지킬 때
send(topic, partition, key, data)여러분이 지정불가동일파티션을 물리적으로 나눠 쓸 때 (2-8 필독)
send(ProducerRecord)레코드에 담긴 값가능지정 가능추적 ID·스키마 버전을 붙일 때
send(Message<?>)KafkaHeaders.PARTITION가능KafkaHeaders.TIMESTAMPSpring Integration / @SendTo 와 섞을 때

①②③ 은 내부에서 ProducerRecord 를 만들어 ④로 위임합니다. 네 가지는 편의 오버로드이고, 진짜는 ProducerRecord 하나입니다.

💡 실무 팁 — defaultTopic 을 쓰지 마세요. sendDefault(key, data) 로 토픽을 생략하면 편해 보이지만, 어느 토픽으로 가는지 grep 으로 찾을 수 없게 됩니다. 토픽 이름은 @ConfigurationProperties 로 주입받아 호출부에 항상 명시하는 편이 낫습니다.


2-2. 반환값은 CompletableFuture<SendResult<K,V>>

Spring Kafka 3.0 에서 반환 타입이 바뀌었습니다.

버전반환 타입콜백
Spring Kafka 2.xListenableFuture<SendResult<K,V>>addCallback(success, failure)
Spring Kafka 3.0+CompletableFuture<SendResult<K,V>>whenComplete, thenAccept, exceptionally

2.x 코드를 3.x 로 올리면 addCallback 이 없어 컴파일 에러가 납니다. 이건 그나마 다행인 경우입니다 — 컴파일러가 잡아 주니까요.

kafkaTemplate.send("orders", event.orderId(), event).whenComplete((result, ex) -> {
    if (ex != null) { log.error("발행 실패 key={}", event.orderId(), ex); return; }
    RecordMetadata md = result.getRecordMetadata();
    log.info("발행 성공 key={} → {}-{}@{} ts={}",
            event.orderId(), md.topic(), md.partition(), md.offset(), md.timestamp());
});

결과

INFO 14071 --- [ad | producer-1] c.e.order.step02.Step02Publisher         : 발행 성공 key=ORD-0001 → orders-2@0 ts=1735689660000
INFO 14071 --- [ad | producer-1] c.e.order.step02.Step02Publisher         : 발행 성공 key=ORD-0002 → orders-2@1 ts=1735689720000
INFO 14071 --- [ad | producer-1] c.e.order.step02.Step02Publisher         : 발행 성공 key=ORD-0003 → orders-0@0 ts=1735689780000

여기서 꼭 짚어야 할 두 가지가 있습니다.

① 스레드가 [main] 이 아니라 [ad | producer-1] 입니다. 콜백은 Kafka 프로듀서의 I/O 스레드(kafka-producer-network-thread | producer-1, 로그에는 뒤쪽만 남습니다)에서 실행됩니다. 콜백 안에서 무거운 일을 하면 안 됩니다. 그 스레드가 막히면 모든 파티션의 전송이 함께 막힙니다. DB 쓰기나 외부 HTTP 호출이 필요하면 별도 executor 로 넘기세요.

② 세 메서드의 역할이 다릅니다. thenAccept 는 성공 시에만, exceptionally 는 실패 시에만 실행됩니다. 둘 다 필요하면 whenComplete 가 유일한 답입니다.

SendResult 에서는 getProducerRecord()(내가 보낸 것)와 getRecordMetadata()(브로커가 알려준 것)를 꺼냅니다. 후자에서 topic() partition() offset() timestamp() serializedValueSize() 를 얻습니다.

💡 partitionoffset 은 브로커만 알 수 있는 값입니다. send() 호출 직후에는 알 수 없고 응답이 와야 채워집니다. "내가 보낸 메시지의 오프셋을 로그에 남기고 싶다"는 요구에 콜백은 유일한 방법입니다.


2-3. ⚠️ 최대 함정 — 결과를 안 보면 유실을 모른다

이 스텝, 아니 이 코스 전체에서 가장 중요한 절입니다.

// 실무에서 가장 흔히 보는 코드
public void publish(OrderCreated event) {
    kafkaTemplate.send("orders", event.orderId(), event);   // 반환값을 버린다
    log.info("주문 이벤트 발행 완료: {}", event.orderId());
    orderRepository.markPublished(event.orderId());          // 발행됐다고 DB 에 기록
}

이 코드는 브로커가 죽어 있어도 예외를 던지지 않습니다. send() 는 레코드를 프로듀서의 메모리 버퍼(buffer.memory, 기본 32MB)에 넣고 즉시 리턴할 뿐이기 때문입니다. 실제 전송은 I/O 스레드가 나중에 합니다.

재현 — 브로커를 내리고 계속 보낸다

기본값 delivery.timeout.ms=120000 은 2분을 기다려야 해서 실습이 지루합니다. 짧게 줄입니다.

spring.kafka.producer.properties:
  delivery.timeout.ms: 10000      # 실습용. 운영 기본값은 120000
  request.timeout.ms: 3000
  max.block.ms: 5000

앱을 띄운 뒤 docker compose stop kafka 로 브로커를 내리고, 위 publish() 를 1초 간격으로 5번 호출합니다.

결과 — 처음 5초

INFO 14071 --- [           main] c.e.order.step02.Step02Publisher         : 주문 이벤트 발행 완료: ORD-0001
INFO 14071 --- [           main] c.e.order.step02.Step02Publisher         : 주문 이벤트 발행 완료: ORD-0002
INFO 14071 --- [           main] c.e.order.step02.Step02Publisher         : 주문 이벤트 발행 완료: ORD-0003
WARN 14071 --- [ad | producer-1] o.a.k.clients.NetworkClient              : [Producer clientId=producer-1] Connection to node 1 (127.0.0.1/127.0.0.1:9092) could not be established. Node may not be available.
INFO 14071 --- [           main] c.e.order.step02.Step02Publisher         : 주문 이벤트 발행 완료: ORD-0004
INFO 14071 --- [           main] c.e.order.step02.Step02Publisher         : 주문 이벤트 발행 완료: ORD-0005

"발행 완료" 가 5번 다 찍혔습니다. markPublished 도 5번 다 커밋됐습니다. 그런데 브로커는 꺼져 있고, 한 건도 가지 않았습니다. NetworkClient 의 WARN 하나가 있지만, 이건 로그 한가운데 흘러가는 인프라 경고입니다. 비즈니스 로그는 완벽하게 성공을 말하고 있습니다.

결과 — delivery.timeout.ms 인 10초가 지난 뒤

ERROR 14071 --- [ad | producer-1] o.s.k.support.LoggingProducerListener    : Exception thrown when sending a message with key='ORD-0001' and payload='OrderCreated[orderId=ORD-0001, customerId=1001, sku=SKU-002, quan...' to topic orders:

org.apache.kafka.common.errors.TimeoutException: Expiring 1 record(s) for orders-2:10023 ms has passed since batch creation

10초가 지나서야 진실이 드러납니다. 그런데 이건 ERRORException 이 아닙니다. 아무도 catch 하지 않습니다.

⚠️ 함정 — spring-kafka 는 ERROR 로그를 내지만, 비즈니스 로직은 성공한 줄 안다 KafkaTemplate 은 기본으로 LoggingProducerListener 를 달아 두므로 전송 실패가 로그에 아예 안 남지는 않습니다. 순수 KafkaProducer 를 콜백 없이 쓸 때보다 나은 점입니다. 문제는 그게 전부라는 것입니다. 그 ERROR 는

  • 호출 스택과 전혀 다른 스레드에서, 수초~수분 뒤에 찍히고
  • 예외로 전파되지 않으므로 트랜잭션 롤백도, 재시도도, 알림도 유발하지 않으며
  • payload 가 잘려 있어 어떤 주문이었는지 복구할 단서가 부족합니다.

증상: "로그는 깨끗한데 컨슈머가 아무것도 못 받았다", "DB 에는 발행됨으로 찍혀 있는데 토픽에는 없다". 해결: send() 의 반환 future 를 반드시 소비하세요. 최소한 whenComplete 로 실패를 잡아 재시도 큐나 Outbox 로 보내야 합니다.

같은 함정의 다른 얼굴

브로커 다운만이 원인이 아닙니다. 콜백 없이는 전부 똑같이 조용합니다.

원인호출부에서 보이는 것future 에 담기는 예외
브로커 다운아무 일도 없음TimeoutException: Expiring 1 record(s)
직렬화 실패여기만 다름 — 즉시 던짐(호출 스레드에서 SerializationException)
토픽 쓰기 권한 없음아무 일도 없음TopicAuthorizationException
max.request.size 초과아무 일도 없음RecordTooLargeException
존재하지 않는 토픽 (자동 생성 off)아무 일도 없음TimeoutException: Topic xxx not present in metadata after 60000 ms

직렬화 실패만 호출 스레드에서 즉시 던집니다. 직렬화는 버퍼에 넣기 전에 호출 스레드가 수행하기 때문입니다. 그래서 "직렬화 오류는 잡히는데 왜 다른 건 안 잡히지?" 하는 혼란이 생깁니다.

CompletableFuture<SendResult<String, OrderCreated>> f =
        kafkaTemplate.send("no-such-topic", "K-1", OrderCreated.of(1));
log.info("send() 리턴 직후 isDone={}", f.isDone());   // ← false

결과

INFO 14071 --- [           main] c.e.order.step02.Step02Publisher         : send() 리턴 직후 isDone=false
WARN 14071 --- [ad | producer-1] o.a.k.clients.NetworkClient              : [Producer clientId=producer-1] Error while fetching metadata with correlation id 4 : {no-such-topic=UNKNOWN_TOPIC_OR_PARTITION}
ERROR 14071 --- [ad | producer-1] o.s.k.support.LoggingProducerListener    : Exception thrown when sending a message with key='K-1' and payload='OrderCreated[orderId=ORD-0001, custo...' to topic no-such-topic:

org.apache.kafka.common.errors.TimeoutException: Topic no-such-topic not present in metadata after 60000 ms.

isDone=false 가 핵심입니다. send() 는 아직 아무것도 결정되지 않은 상태로 리턴합니다. 그리고 그 상태가 60초 동안 이어집니다.

💡 최소한의 안전장치 한 줄 전 코드를 고칠 여유가 없다면, 커스텀 ProducerListener 를 등록해 실패를 한곳으로 모으는 것부터 하세요. kafkaTemplate.setProducerListener(...) 로 갈아 끼우면 호출부를 하나도 안 고치고 실패 카운터(Micrometer)와 알림을 붙일 수 있습니다. 다만 "어떤 주문이 실패했는지 되살리는" 일은 여전히 호출부의 몫입니다.


2-4. ⚠️ 두 번째 함정 — send().get() 의 성능 붕괴

2-3 을 읽고 나면 자연스럽게 이 결론에 도달합니다. "그럼 매번 .get() 으로 기다리면 되겠네." 확실하긴 합니다. 그리고 20배 느립니다.

왜 느린가 — 배치가 통째로 무력화된다

Kafka 프로듀서는 레코드를 파티션별 배치 버퍼에 모읍니다. linger.ms=5 면 "5ms 모았다가 한 번에", batch.size=16384 면 "16KB 가 차면 즉시" 보냅니다. 한 번의 왕복에 수백 건이 함께 갑니다.

[ 정상 ]  send×N (논블로킹, 버퍼에 쌓임) ──▶ linger 경과 or batch 도달
                                      ┌──────────────────────┐
                                      │ 한 번의 요청에 300건 │──▶ 브로커
                                      └──────────────────────┘   왕복 1회 ≈ 2ms

[ get() ] send ─▶ 버퍼에 1건 ─▶ get() 블로킹 ─▶ linger 5ms 대기 ─▶ 전송 ─▶ 응답 대기 ─┐
          send ─▶ 버퍼에 1건 ─▶ ... 처음부터 다시 ◀───────────────────────────────────┘
                                      왕복 1회 = 1건. 게다가 건당 linger.ms 를 그냥 버린다

.get() 은 그 레코드 하나의 응답을 기다립니다. 기다리는 동안 호출 스레드가 다음 send() 를 못 하므로 버퍼에 두 번째 레코드가 들어올 수 없습니다. 배치 크기가 영원히 1이 됩니다. 게다가 linger.ms=5 는 이제 이득이 아니라 건당 5ms 의 순수 페널티입니다.

실측 — 10,000건 발행

int N = 10_000;

// (a) fire-and-forget
long t0 = System.nanoTime();
for (int i = 1; i <= N; i++) kafkaTemplate.send("orders", "ORD-%04d".formatted(i), OrderCreated.of(i));
kafkaTemplate.flush();                          // ★ 측정 구간 안에 있어야 공정하다
long elapsedA = (System.nanoTime() - t0) / 1_000_000;

// (c) 매 건 get()
long t1 = System.nanoTime();
for (int i = 1; i <= N; i++) kafkaTemplate.send("orders", "ORD-%04d".formatted(i), OrderCreated.of(i)).get();
long elapsedC = (System.nanoTime() - t1) / 1_000_000;

결과

INFO 14071 --- [           main] c.e.order.step02.ThroughputBench         : (a) fire-and-forget + flush : 10000건 /   543 ms = 18,416 msg/s
INFO 14071 --- [           main] c.e.order.step02.ThroughputBench         : (b) 콜백(whenComplete)      : 10000건 /   552 ms = 18,116 msg/s
INFO 14071 --- [           main] c.e.order.step02.ThroughputBench         : (c) 매 건 send().get()      : 10000건 / 10870 ms =    920 msg/s
INFO 14071 --- [           main] c.e.order.step02.ThroughputBench         : (d) 배치 후 allOf().join()  : 10000건 /   559 ms = 17,889 msg/s

18,416 msg/s → 920 msg/s. 20배 느려졌습니다. 건당 평균 1.087ms 이며, linger.ms 를 0 으로 내려도 660 msg/s 대에서 크게 나아지지 않습니다. 병목은 linger 가 아니라 왕복(RTT)을 건마다 지불하는 구조 그 자체입니다.

측정할 때 워밍업 1,000건을 먼저 보내세요.send() 는 토픽 메타데이터 조회로 수백 ms 를 쓰고 JIT 도 아직 인터프리터 모드라, 워밍업 없이 (a)→(c) 순으로 재면 (a) 가 (c) 보다 느리게 나오는 황당한 결과가 나옵니다.

필수 비교표 — 발행 방식 4가지

방식코드처리량유실 감지순서 보장언제
fire-and-forgetsend(...)18,400 msg/s불가(로그만)✅ 파티션 내 보장유실돼도 되는 지표·트레이스
콜백send(...).whenComplete(...)18,100 msg/s✅ 건별 가능✅ 파티션 내 보장대부분의 경우 정답
매 건 get()send(...).get()920 msg/s✅ 즉시, 동기✅ 보장초당 수백 건 이하 + 호출부가 즉시 실패를 알아야 할 때
배치 후 flush/joinsend() ×N → allOf().join()17,900 msg/s✅ 배치 단위✅ 파티션 내 보장대량 마이그레이션·일괄 발행

콜백이 fire-and-forget 과 거의 같다는 점이 중요합니다. 콜백은 사실상 공짜이므로 "안전"과 "속도" 중 하나를 고를 필요가 없습니다.

⚠️ 함정 — "순서 보장"이 .get() 의 이유가 될 수는 없다 .get() 을 쓰는 두 번째 이유로 "순서 때문에"를 드는 경우가 많은데 틀렸습니다. 같은 키의 메시지는 같은 파티션으로 가고, 프로듀서는 그 파티션의 배치를 보낸 순서대로 전송합니다. .get() 없이도 순서는 지켜집니다. 순서를 깨뜨리는 건 .get() 의 부재가 아니라 retries>0 + max.in.flight>1 + 멱등성 off 조합입니다 (2-9).

대안 4가지

// ① 콜백 — 기본형. 처리량 손실 거의 없음
kafkaTemplate.send("orders", key, event)
        .whenComplete((r, ex) -> { if (ex != null) failureSink.record(key, event, ex); });

// ② flush() — 루프가 끝난 뒤 한 번. 버퍼를 비우고 모든 응답을 기다린다
for (...) kafkaTemplate.send(...);
kafkaTemplate.flush();

// ③ 마지막 future 만 join — ⚠️ 모든 레코드가 "같은 파티션"일 때만 성립한다
CompletableFuture<SendResult<String, OrderCreated>> last = null;
for (...) last = kafkaTemplate.send("orders", sameKey, event);
last.join();

// ④ allOf — 전부 모아서 한 번에 대기. 개별 실패도 전부 볼 수 있다 ★키가 여러 개면 이것
List<CompletableFuture<SendResult<String, OrderCreated>>> futures = new ArrayList<>();
for (int i = 1; i <= N; i++) futures.add(kafkaTemplate.send("orders", key(i), OrderCreated.of(i)));
CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new)).join();
long failed = futures.stream().filter(CompletableFuture::isCompletedExceptionally).count();
log.info("발행 {}건 중 실패 {}건", N, failed);   // → 발행 10000건 중 실패 0건

③번은 파티션이 다르면 성립하지 않습니다. 각 파티션의 전송은 독립적이라, 마지막 레코드가 성공해도 다른 파티션의 앞선 레코드가 실패해 있을 수 있습니다.

⚠️ flush() 를 요청마다 호출하지 마세요. flush()그 프로듀서 인스턴스 전체의 버퍼를 비웁니다. KafkaTemplate 은 프로듀서를 공유하므로 HTTP 요청 하나가 flush() 를 부르면 다른 요청들이 쌓아 둔 배치까지 강제로 밀어냅니다. 동시 요청이 많은 서버에서 이러면 .get() 과 똑같은 결과가 되고, 원인 찾기는 훨씬 어렵습니다. flush()배치 작업의 끝 또는 종료 직전에만 쓰세요.


2-5. ProducerFactory — 설정이 다른 KafkaTemplate 두 개

KafkaTemplateProducerFactory 에서 프로듀서를 얻습니다. DefaultKafkaProducerFactory싱글턴 프로듀서를 만들어 모든 호출자에게 같은 인스턴스를 돌려줍니다.

Producer<String, OrderCreated> p1 = pf.createProducer();
Producer<String, OrderCreated> p2 = pf.createProducer();
log.info("같은 인스턴스인가? {}", p1 == p2);   // → true

KafkaProducer스레드 안전하며, 여러 스레드가 공유하는 것이 권장되는 사용법입니다. 배치와 커넥션을 공유해야 효율이 나기 때문입니다.

⚠️ 함정 — 요청마다 KafkaProducer 를 새로 만들면 안 된다 프로듀서 하나는 I/O 스레드 1개 + 메타데이터 캐시 + 32MB 버퍼를 들고 있습니다. 게다가 배치가 전혀 공유되지 않아 2-4 의 .get() 과 같은 성능이 나옵니다. 증상은 성능 저하와 파일 디스크립터 고갈입니다. close() 를 빠뜨리면 그대로 누수입니다. KafkaTemplate 을 주입받아 쓰는 것이 정답입니다.

"주문 이벤트는 절대 잃으면 안 되고, 클릭 로그는 좀 잃어도 되니 빠르게" 같은 요구는 흔합니다. 설정이 다르므로 프로듀서를 분리합니다.

@Configuration
public class DualProducerConfig {

    @Bean   // ① 신뢰성 우선 — 주문 이벤트
    public KafkaTemplate<String, OrderCreated> reliableKafkaTemplate(
            @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
        Map<String, Object> props = base(bootstrap);          // bootstrap + 직렬화기
        props.put(ProducerConfig.ACKS_CONFIG, "all");
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
        props.put(ProducerConfig.LINGER_MS_CONFIG, 5);
        props.put(ProducerConfig.CLIENT_ID_CONFIG, "reliable");   // ★ 반드시 다르게
        return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props));
    }

    @Bean   // ② 처리량 우선 — 클릭 로그
    public KafkaTemplate<String, OrderCreated> fastKafkaTemplate(
            @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
        Map<String, Object> props = base(bootstrap);
        props.put(ProducerConfig.ACKS_CONFIG, "1");
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, false);   // acks=1 이면 명시적 false
        props.put(ProducerConfig.LINGER_MS_CONFIG, 50);
        props.put(ProducerConfig.BATCH_SIZE_CONFIG, 65_536);
        props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
        props.put(ProducerConfig.CLIENT_ID_CONFIG, "fast");
        return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props));
    }
}

결과 — 기동 로그

INFO 14071 --- [           main] o.a.k.c.u.AppInfoParser                  : Kafka version: 3.6.1
INFO 14071 --- [           main] o.a.k.c.p.i.TransactionManager           : [Producer clientId=reliable] Instantiated an idempotent producer.
INFO 14071 --- [           main] o.a.k.c.p.i.ProducerIdAndEpoch           : [Producer clientId=reliable] ProducerId set to 0 with epoch 0

fast 쪽에는 Instantiated an idempotent producer.없습니다. 이 한 줄의 유무가 멱등성 활성화 여부를 알려 주는 가장 확실한 신호입니다. clientId 를 안 주면 둘 다 producer-1, producer-2 로 자동 부여되어 로그와 JMX 지표에서 어느 쪽인지 구분할 수 없습니다.

⚠️ 함정 — 타입이 같은 KafkaTemplate 빈이 둘이면 자동 주입이 실패한다 자동 설정이 만든 것까지 셋이 되어, 이름 없는 주입은 NoUniqueBeanDefinitionException: expected single matching bean but found 3 으로 기동 실패합니다. 이건 다행히 시끄러운 실패입니다. @Qualifier("reliableKafkaTemplate") 로 명시하세요. @Primary 는 권장하지 않습니다. fast 쪽에 붙어 있으면 아무 생각 없이 주입받은 코드가 전부 acks=1 로, 멱등성 없이 발행됩니다. 컴파일도 기동도 테스트도 통과하며, 브로커가 흔들릴 때까지 아무도 모릅니다. 그게 진짜 함정입니다.


2-6. ProducerRecord 직접 만들기 — 헤더와 타임스탬프

편의 오버로드로는 헤더를 못 붙입니다. 추적 ID·스키마 버전·이벤트 타입 같은 메타데이터는 헤더가 제자리입니다. 페이로드에 섞으면 컨슈머가 역직렬화를 해야만 읽을 수 있기 때문입니다.

OrderCreated event = OrderCreated.of(7);

ProducerRecord<String, OrderCreated> record = new ProducerRecord<>(
        "orders",                                   // topic
        null,                                       // partition — null 이면 파티셔너에 맡김
        event.createdAt().toEpochMilli(),           // timestamp — 이벤트 발생 시각
        event.orderId(),                            // key
        event);                                     // value
record.headers()
      .add("trace-id",       "t-9f2c1a".getBytes(StandardCharsets.UTF_8))
      .add("event-type",     "OrderCreated".getBytes(StandardCharsets.UTF_8))
      .add("schema-version", "2".getBytes(StandardCharsets.UTF_8));

kafkaTemplate.send(record).whenComplete((r, ex) -> {
    RecordMetadata md = r.getRecordMetadata();
    log.info("→ {}-{}@{} ts={}", md.topic(), md.partition(), md.offset(), md.timestamp());
});

결과→ orders-0@12 ts=1735690080000. 콘솔 컨슈머로 확인합니다.

docker exec -it learn-kafka /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server localhost:9092 --topic orders --partition 0 --offset 12 --max-messages 1 \
  --property print.key=true --property print.headers=true --property print.timestamp=true

결과

CreateTime:1735690080000	trace-id:t-9f2c1a,event-type:OrderCreated,schema-version:2,__TypeId__:com.example.order.domain.OrderCreated	ORD-0007	{"orderId":"ORD-0007","customerId":1007,"sku":"SKU-002","quantity":3,"amount":10000,"createdAt":"2025-01-01T00:07:00Z"}

__TypeId__ 는 우리가 붙인 게 아니라 JsonSerializer 가 자동으로 붙인 헤더입니다. Step 04 의 주인공이니 기억해 두세요. CreateTime:1735690080000우리가 지정한 이벤트 발생 시각이 그대로 살아 있습니다. 지정하지 않으면 send() 를 호출한 시각이 들어갑니다.

⚠️ 함정 — 토픽이 LogAppendTime 이면 여러분이 지정한 타임스탬프는 조용히 버려진다 토픽 설정 message.timestamp.typeLogAppendTime 이면 브로커가 도착 시각으로 덮어씁니다. 에러도 경고도 없습니다. 재처리·마이그레이션에서 과거 이벤트를 원래 시각으로 다시 넣으려 할 때 조용히 무력화되어, --from-timestamp 오프셋 조회와 시계열 집계가 전부 틀어집니다. 확인: kafka-configs.sh --describe --entity-type topics --entity-name orders. 기본값은 CreateTime 입니다.

💡 헤더 값은 byte[] 입니다. 인코딩을 명시하세요(StandardCharsets.UTF_8). 플랫폼 기본 인코딩에 맡기면 로컬과 서버에서 다르게 직렬화되어 한글 헤더가 깨집니다.


2-7. 키와 파티셔너 — 어디로 가는지 계산할 수 있다

키가 있으면: murmur2(key) % partitions

Kafka 의 기본 파티셔너는 키가 있으면 결정적으로 파티션을 고릅니다. 순수 함수이므로 브로커에 물어보지 않고 계산할 수 있습니다.

int hash = Utils.murmur2(key.getBytes(StandardCharsets.UTF_8));
int partition = Utils.toPositive(hash) % numPartitions;    // toPositive = hash & 0x7fffffff

⚠️ Math.abs(hash) % n 으로 쓰면 안 됩니다. Math.abs(Integer.MIN_VALUE) 는 오버플로로 자기 자신(음수) 이 됩니다. 그러면 파티션 번호가 음수가 되어 ArrayIndexOutOfBoundsException 이 납니다. 40억분의 1이지만 하루 수억 건이면 언젠가 터집니다. Kafka 가 Math.abs 대신 toPositive(비트 마스크)를 쓰는 이유가 정확히 이것입니다.

결과 — orders 는 파티션 3개

murmur2toPositive% 3 (현재)% 4% 6
ORD-0001-1289372403858111245215
ORD-0002-1559091975588391673215
ORD-000319614102421961410242020
ORD-0004260207539260207539131
ORD-0005-1928704172218779476000
ORD-0006-6013546631546128985111
ORD-000712400796311240079631033
ORD-0008252067230252067230020
ORD-000913386092691338609269215

실제로 발행해 콜백이 알려준 파티션과 대조합니다.

결과

INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : ORD-0001 예측=2 실제=2 offset=0
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : ORD-0002 예측=2 실제=2 offset=1
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : ORD-0003 예측=0 실제=0 offset=0
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : ORD-0004 예측=1 실제=1 offset=0
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : ORD-0005 예측=0 실제=0 offset=1
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : ORD-0006 예측=1 실제=1 offset=1
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : ORD-0007 예측=0 실제=0 offset=2
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : ORD-0008 예측=0 실제=0 offset=3
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : ORD-0009 예측=2 실제=2 offset=2

9건 전부 일치합니다. 파티션 배정은 마법이 아니라 계산입니다. 분포도 눈여겨보세요. 9건 중 파티션 0 에 4건, 1 에 2건, 2 에 3건입니다. 키가 균등해도 파티션은 균등하지 않습니다. "핫 파티션"의 씨앗이 여기 있습니다.

⚠️ 파티션 수를 바꾸면 매핑이 통째로 바뀐다

위 표의 % 4 열을 보세요. 파티션을 3 → 4 로 늘리면 ORD-0001 은 2→1, ORD-0003 은 0→2, ORD-0004 는 1→3, ORD-0007 은 0→3 으로 9개 키 중 8개가 바뀝니다.

⚠️ 함정 — 파티션 증설은 순서 보장을 소급해서 깨뜨린다 ORD-0001 의 과거 이벤트는 orders-2 에 쌓여 있는데, 증설 직후부터 같은 주문의 새 이벤트가 orders-1 로 갑니다. 두 파티션은 독립적으로 소비되므로, 나중에 발행된 OrderCancelled 가 먼저 발행된 OrderCreated 보다 먼저 처리될 수 있습니다. 에러는 없습니다. 재고가 음수가 되거나 취소된 주문이 배송되는 것으로 드러납니다. 해결: 파티션 수는 처음에 넉넉히(예상 최대 컨슈머 수의 2배) 잡고 바꾸지 않습니다. 꼭 늘려야 한다면 신규 토픽을 만들어 컨슈머를 전환하는 편이 안전합니다. (파티션 재할당의 브로커 쪽 동작은 Kafka 코스 Step 03 참고)

키가 null 이면: sticky partitioner

버전동작
~2.3라운드로빈. 레코드마다 다음 파티션. 배치가 잘게 쪼개져 비효율
2.4 ~ 3.2Sticky. 한 파티션에 배치가 찰 때까지 몰아넣고, 배치가 나가면 다음 파티션으로
3.3+partitioner.class 기본값이 null 이 되고 내장 uniform sticky 로직을 사용. DefaultPartitioner 는 deprecated

우리 환경(kafka-clients 3.6.1)은 세 번째입니다. partitioner.class아예 지정하지 않는 것이 정상이며, 예전 글을 보고 DefaultPartitioner 를 명시하면 deprecation 경고와 함께 구버전 동작으로 되돌아갑니다.

for (int i = 1; i <= 6; i++) kafkaTemplate.send("orders", OrderCreated.of(i));  // 키 없음

결과

INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : key=null → orders-1@0
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : key=null → orders-1@1
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : key=null → orders-1@2
INFO 14071 --- [ad | producer-1] c.e.order.step02.PartitionDemo           : key=null → orders-1@3   (이하 orders-1@4, orders-1@5)

6건이 전부 orders-1 에 몰렸습니다. 배치가 아직 안 찼기 때문입니다. 라운드로빈을 기대했다면 당황스럽겠지만 이게 정상 동작이며, 수만 건을 보내면 장기적으로 고르게 분산됩니다.

⚠️ 함정 — "키를 안 주면 알아서 골고루 퍼진다"는 착각 소량 발행에서는 한 파티션에 몰립니다. 통합 테스트에서 "3개 파티션에 골고루 갔는지" 검증하는 assert 를 쓰면 로컬에서는 통과하다가 CI 에서 배치 타이밍이 달라지며 간헐 실패합니다 (Step 11). 그리고 키가 없으면 어떤 순서 보장도 없습니다. 같은 주문의 이벤트 두 개가 다른 파티션으로 갈 수 있습니다.


2-8. 명시적 파티션 지정과 그 위험

kafkaTemplate.send("orders", 1, event.orderId(), event);   // 무조건 orders-1

파티션을 직접 주면 파티셔너를 완전히 건너뜁니다. 키는 여전히 레코드에 실려 가지만 파티션 결정에는 아무 영향이 없습니다. 정당한 용도는 재처리(DLT 메시지를 원래 파티션으로), 테스트, 파티션 전용 워크로드 정도이고, 그 외에는 거의 항상 잘못된 선택입니다.

⚠️ 함정 — 하드코딩한 파티션 번호는 범위를 넘으면 즉시 예외, 파티션이 늘면 조용히 편중 파티션 3개인 토픽에 send("orders", 5, key, value) 를 하면:

ERROR 14071 --- [ad | producer-1] o.s.k.support.LoggingProducerListener    : Exception thrown when sending a message with key='ORD-0001' and payload='OrderCreated[orderId=ORD-0001, cust...' to topic orders:

java.lang.IllegalArgumentException: Invalid partition given with record: 5 is not in the range [0...3).

이건 시끄러운 실패라 다행입니다. 진짜 문제는 반대입니다. 파티션을 3 → 12 로 늘렸는데 코드가 send(topic, i % 3, key, value) 로 하드코딩돼 있으면 파티션 3~11 은 영원히 비어 있고, 컨슈머 9개가 아무것도 안 받으며 놉니다. 에러도 경고도 없습니다. 해결: 파티션 수를 코드에 박지 말고 kafkaTemplate.partitionsFor("orders").size() 로 런타임에 조회하세요.

List<PartitionInfo> infos = kafkaTemplate.partitionsFor("orders");
log.info("orders 파티션 수 = {}", infos.size());     // → orders 파티션 수 = 3

💡 partitionsFor 는 메타데이터가 없으면 max.block.ms 까지 블로킹합니다. 요청 경로에서 매번 부르지 마세요. 기동 시 한 번 조회해 캐시하되, 파티션이 늘어날 수 있음을 감안해 주기적으로 갱신하는 편이 안전합니다.


2-9. acks / retries / enable.idempotence — 세 설정의 상호 제약

이 셋은 독립적이지 않습니다. 조합에 따라 기동이 실패하기도 하고, 더 나쁘게는 조용히 무력화되기도 합니다.

설정kafka-clients 3.x 기본값의미
acksall (idempotence 기본 on 때문)0=응답 안 기다림 / 1=리더만 / all=모든 ISR
retriesInteger.MAX_VALUE실질 상한은 delivery.timeout.ms
enable.idempotencetrue (Kafka 3.0+)재시도로 인한 중복을 브로커가 제거
max.in.flight.requests.per.connection5응답을 안 받은 채 동시에 날릴 수 있는 요청 수

💡 Spring Boot 3.x 는 enable.idempotence 를 설정하지 않습니다. truekafka-clients 3.0 부터의 클라이언트 기본값입니다. 아무것도 안 건드리면 여러분의 프로듀서는 이미 멱등 프로듀서입니다. 기동 로그의 Instantiated an idempotent producer. 로 확인하세요.

⚠️ enable.idempotence=true + acks=1 → 기동 실패

spring.kafka.producer:
  acks: 1                                  # ← 처리량 좀 올려 보려고
  properties.enable.idempotence: true      # ← 중복도 막고 싶어서

결과

INFO 14071 --- [           main] c.e.o.OrderServiceApplication            : The following 1 profile is active: "step02"
ERROR 14071 --- [           main] o.s.boot.SpringApplication               : Application run failed

org.springframework.beans.factory.BeanCreationException: Error creating bean with name 'step02Publisher': ...
Caused by: org.apache.kafka.common.config.ConfigException: Must set acks to all in order to use the idempotent producer. Otherwise we cannot guarantee idempotence.
	at org.apache.kafka.clients.producer.ProducerConfig.postProcessAndValidateIdempotenceConfigs(ProducerConfig.java:594)
	at org.apache.kafka.clients.producer.KafkaProducer.<init>(KafkaProducer.java:334)

멱등성은 acks=all 을 요구합니다. 리더만 받고 성공을 반환하면 리더가 죽었을 때 그 레코드가 사라지고, 재시도가 "중복"이 아니라 "재발행"이 되어 멱등성의 전제가 무너지기 때문입니다. 같은 종류의 제약이 셋입니다.

위반 조합예외 메시지
enable.idempotence=true + acks=1 또는 0Must set acks to all in order to use the idempotent producer.
enable.idempotence=true + retries=0Must set retries to non-zero when using the idempotent producer.
enable.idempotence=true + max.in.flight > 5Must set max.in.flight.requests.per.connection to at most 5 when using the idempotent producer.

⚠️ 진짜 함정 — acks=1 만 쓰면 멱등성이 조용히 꺼진다

위 셋은 명시적으로 enable.idempotence=true 를 쓴 경우입니다. acks: 1 만 주면?

결과 — 기동 성공

INFO 14071 --- [           main] o.a.k.c.u.AppInfoParser                  : Kafka version: 3.6.1
INFO 14071 --- [           main] c.e.o.OrderServiceApplication            : Started OrderServiceApplication in 1.982 seconds (process running for 2.310)

아무 일도 없이 뜹니다. Instantiated an idempotent producer.사라졌다는 사실을 알아채기 전까지는 멀쩡해 보입니다.

⚠️ 함정 — 설정하지 않은 멱등성은 "충돌하면 조용히 false 로 내려간다" ProducerConfigenable.idempotence사용자에 의해 명시되지 않았고 다른 설정과 충돌하면, 예외 대신 조용히 false 로 내립니다. "기본값이니까 사용자가 원한 게 아닐 수도 있다"는 판단입니다. 즉 acks: 1 한 줄이 멱등성을 통째로 끕니다. 그 결과 retries 로 인한 중복 레코드가 생기고, max.in.flight=5 와 만나면 순서까지 뒤집힙니다. 증상: 평소엔 멀쩡하다가, 네트워크가 한 번 흔들린 뒤 컨슈머가 같은 주문을 두 번 처리하거나 역순으로 처리합니다. 확인법: 기동 로그에서 Instantiated an idempotent producer. 를 찾으세요. 없으면 꺼진 것입니다. 해결: 멱등성이 필요하면 enable.idempotence: true명시하세요. 그러면 충돌 시 조용히 꺼지는 대신 기동이 실패합니다. 시끄러운 실패를 사는 것이 조용한 유실을 사는 것보다 언제나 낫습니다.

retries + max.in.flight 가 순서를 뒤집는 원리

멱등성이 꺼진 상태에서 retries=3, max.in.flight=5 라면:

프로듀서 → 브로커  요청A [ORD-0001 첫 번째 이벤트]  ┐
프로듀서 → 브로커  요청B [ORD-0001 두 번째 이벤트]  ┘ 동시에 in-flight

브로커: 요청A 처리 중 일시 오류 → NOT_ENOUGH_REPLICAS 응답
브로커: 요청B 정상 처리 → offset 10 에 기록          ← B 가 먼저 들어감!
프로듀서: 요청A 재시도 → offset 11 에 기록           ← A 가 나중에

결과: orders-2 [10]=두 번째 이벤트, [11]=첫 번째 이벤트   ← 같은 파티션인데 순서 역전

컨슈머는 이걸 알 방법이 없습니다. 멱등성이 켜져 있으면 각 레코드가 시퀀스 번호를 받아, 브로커가 순서 어긋난 요청을 OUT_OF_ORDER_SEQUENCE_NUMBER 로 거부하고 프로듀서가 순서대로 다시 보냅니다. 그래서 max.in.flight=5 여도 순서가 보장됩니다.

조합중복순서처리량
idempotence=true, acks=all, in.flight≤5없음보장높음 — 권장 기본값
idempotence=false, retries>0, in.flight=1있을 수 있음보장낮음 (파이프라이닝 없음)
idempotence=false, retries>0, in.flight>1있을 수 있음깨짐 ⚠️높음 — 가장 위험한 조합
acks=0보장 안 함최고 — 유실 감지 자체가 불가

💡 acks=all 이 실제로 몇 대를 기다리는지, min.insync.replicas 와 어떻게 맞물리는지는 Kafka 코스 Step 04Step 08 에서 다중 브로커로 확인하세요. 전달 보장(at-most-once / at-least-once / exactly-once)의 전체 그림은 Kafka 코스 Step 07 입니다. 우리 실습 브로커는 단일 노드·복제본 1이라 acks=all 이 사실상 acks=1 과 같은 비용입니다.


정리

개념핵심
send 오버로드5종 전부 내부적으로 ProducerRecord 로 수렴. 헤더·타임스탬프가 필요하면 ProducerRecord
반환 타입Spring Kafka 3.0 부터 CompletableFuture (2.x 는 ListenableFuture)
콜백 스레드[ad | producer-1]프로듀서 I/O 스레드. 무거운 일 금지
최대 함정send() 결과를 버리면 브로커 다운·권한 오류·크기 초과가 예외 없이 지나간다
LoggingProducerListener실패 시 ERROR 로그는 남는다. 하지만 비즈니스 로직은 성공한 줄 안다
직렬화 실패만 예외직렬화는 호출 스레드에서 수행 → 즉시 SerializationException
두 번째 함정send().get()18,400 → 920 msg/s (20배). 배치 크기가 영원히 1
대안콜백 / flush() / 마지막 future join(같은 파티션일 때만) / allOf
flush()프로듀서 전체 버퍼를 비움. 요청마다 호출하면 .get() 과 동일한 참사
ProducerFactory프로듀서는 싱글턴 + 스레드 안전. 요청마다 만들지 말 것
템플릿 2개NoUniqueBeanDefinitionException. @Qualifier 로 명시, @Primary 는 피하고 clientId 필수
키 → 파티션toPositive(murmur2(key)) % N. 손으로 계산 가능. Math.abs 금지
파티션 수 변경매핑이 통째로 바뀜 → 순서 보장이 소급해서 깨짐
null3.3+ 내장 uniform sticky. 소량이면 한 파티션에 몰린다
명시적 파티션범위 초과는 시끄럽게 실패. 파티션 증설 시 조용히 편중
idempotence + acks=1명시했으면 ConfigException 기동 실패 / 안 했으면 조용히 false ⚠️
순서 역전 조합idempotence=false + retries>0 + max.in.flight>1
확인 한 줄기동 로그의 Instantiated an idempotent producer.

연습문제

Exercise.java 에 6문제가 있습니다. 정답은 Solution.java. 반드시 직접 실행해 로그를 확인하세요.

  1. whenComplete 콜백으로 발행한 메시지의 topic-partition@offset 을 로깅하기
  2. 10,000건 발행을 4가지 방식으로 측정해 처리량 비교표 만들기
  3. 존재하지 않는 토픽으로 send() 했을 때 future 의 상태 변화를 시간 순으로 관찰하기
  4. acks/linger.ms 가 다른 KafkaTemplate 두 개를 공존시키고 @Qualifier 로 골라 쓰기
  5. ORD-0010 ~ ORD-0015 의 파티션을 murmur2 로 예측하고 실제 발행 결과와 대조하기
  6. acks=1 + enable.idempotence=true 조합의 ConfigException 을 재현하고, 명시하지 않았을 때와 비교하기

다음 단계

지금까지는 보내기만 했습니다. 보낸 것이 정말 갔는지는 콘솔 컨슈머로 눈으로 확인했죠. 다음 스텝에서는 @KafkaListener받는 쪽을 만듭니다. 그리고 이 스텝에서 계산한 파티션 3개concurrency 설정과 어떻게 맞물리는지 — concurrency=5 로 잡으면 남는 스레드 2개가 경고 하나 없이 그냥 노는 것을 확인합니다.

Step 03 — @KafkaListener 기초


실습 파일

세 파일을 순서대로 씁니다. 먼저 Practice.javasrc/main/java/com/example/order/step02/ 에 넣고 --spring.profiles.active=step02 로 실행하며 2-1 ~ 2-9 의 모든 로그를 재현합니다. 특히 2-3 의 브로커 다운 재현과 2-4 의 처리량 측정은 직접 눈으로 봐야 체감이 됩니다. 그다음 Exercise.java 의 6문제를 풀고, Solution.java 로 정답과 해설을 대조합니다.

Practice.java

교재 본문의 모든 예제를 절 번호 주석(// [2-4] send().get() 성능 붕괴)과 함께 담은 단일 실행 파일입니다.

  • 파일 상단 주석에 실행 명령, 필요한 application-step02.yml 설정, 확인용 CLI 명령이 전부 적혀 있습니다. 이 스텝은 delivery.timeout.ms 를 10초로 줄여야 2-3 을 90분 안에 끝낼 수 있으므로, 프로필 전용 yml 을 반드시 만드세요.
  • 실행은 SILENT_LOSS_DEMO / BENCH / PARTITION_MAP 세 개의 불리언 스위치로 제어합니다. 전부 켜면 브로커를 내렸다 올리는 사이에 벤치마크가 섞여 결과가 오염되므로, 한 번에 하나씩 켜서 돌리세요.
  • [2-3]SilentLossDemo 는 브로커가 꺼져 있어야 의미가 있습니다. 앱을 먼저 띄우고 docker compose stop kafka 를 한 뒤 스위치를 켠 채 재시작하면, "발행 완료" 5줄이 먼저 찍히고 10초 뒤 LoggingProducerListener 의 ERROR 가 몰려 나옵니다. 실습이 끝나면 docker compose start kafka 로 되돌리세요.
  • [2-4]ThroughputBench 는 측정 전에 워밍업 1,000건을 먼저 보냅니다. JIT 컴파일과 메타데이터 조회를 측정 구간에서 빼기 위한 것으로, 이걸 빼면 (a) 가 (c) 보다 느리게 나오는 황당한 결과가 나옵니다.
  • [2-5]DualProducerConfig@Profile("step02-dual") 로 분리해 두었습니다. 켜면 KafkaTemplate 빈이 셋(자동 설정 1 + 여기 2)이 되므로, 이 프로필에서는 반드시 @Qualifier 로 주입받아야 합니다.
  • [2-7]PartitionMapDemoUtils.murmur2 를 직접 호출해 예측값을 먼저 로깅한 뒤 발행 콜백의 실제 파티션과 대조합니다. 두 값이 다르다면 토픽의 파티션 수가 3이 아닌 것이니 그 자리에서 멈추세요.
/*
 * ============================================================================
 *  Step 02 — KafkaTemplate 과 프로듀서 : Practice
 * ============================================================================
 *
 *  배치 위치 : src/main/java/com/example/order/step02/Practice.java
 *  실행      : ./gradlew bootRun --args='--spring.profiles.active=step02'
 *
 *  [필수] src/main/resources/application-step02.yml 을 만드세요.
 *         2-3 의 "조용한 유실" 재현은 delivery.timeout.ms 기본값(120초)으로는
 *         2분을 기다려야 합니다. 실습용으로 10초까지 줄입니다.
 *
 *      spring:
 *        kafka:
 *          producer:
 *            acks: all
 *            properties:
 *              enable.idempotence: true
 *              max.in.flight.requests.per.connection: 5
 *              linger.ms: 5
 *              batch.size: 16384
 *              delivery.timeout.ms: 10000     # 실습용. 운영 기본값 120000
 *              request.timeout.ms: 3000
 *              max.block.ms: 5000
 *      logging:
 *        level:
 *          com.example.order: DEBUG
 *
 *  [실행 스위치] 아래 Step02Runner 의 상수 3개로 무엇을 돌릴지 고릅니다.
 *      SILENT_LOSS_DEMO : 2-3 조용한 유실 재현  (브로커를 내린 상태에서만 의미 있음)
 *      BENCH            : 2-4 처리량 4방식 측정 (브로커가 살아 있어야 함)
 *      PARTITION_MAP    : 2-7 키 → 파티션 예측/대조
 *  ⚠️ 한 번에 하나만 true 로 두세요. 섞으면 결과가 오염됩니다.
 *
 *  [확인용 CLI]
 *      # 토픽 파티션 수 확인 (반드시 3이어야 2-7 표와 일치)
 *      docker exec -it learn-kafka /opt/kafka/bin/kafka-topics.sh \
 *        --bootstrap-server localhost:9092 --describe --topic orders
 *
 *      # 발행 결과를 눈으로 보기 (헤더/키/파티션/오프셋)
 *      docker exec -it learn-kafka /opt/kafka/bin/kafka-console-consumer.sh \
 *        --bootstrap-server localhost:9092 --topic orders --from-beginning \
 *        --property print.key=true --property print.headers=true \
 *        --property print.partition=true --property print.offset=true
 *
 *      # 2-3 재현용: 브로커 내리기 / 되돌리기
 *      docker compose stop kafka
 *      docker compose start kafka
 * ============================================================================
 */
package com.example.order.step02;

import com.example.order.domain.OrderCreated;

import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.PartitionInfo;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.common.utils.Utils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.kafka.support.SendResult;
import org.springframework.kafka.support.serializer.JsonSerializer;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

public final class Practice {

    private Practice() {
    }

    private static final String TOPIC = "orders";

    // ========================================================================
    // 실행 진입점 — 스위치로 무엇을 돌릴지 고른다
    // ========================================================================
    @Component
    @Profile("step02")
    public static class Step02Runner implements ApplicationRunner {

        private static final Logger log = LoggerFactory.getLogger(Step02Runner.class);

        /** 2-3. 브로커를 내린 상태에서만 의미가 있습니다. */
        private static final boolean SILENT_LOSS_DEMO = false;
        /** 2-4. 브로커가 살아 있어야 합니다. orders 토픽에 40,000건이 쌓입니다. */
        private static final boolean BENCH = false;
        /** 2-7. 키 → 파티션 예측/대조. 기본으로 켜 둡니다. */
        private static final boolean PARTITION_MAP = true;

        private final KafkaTemplate<String, OrderCreated> template;

        public Step02Runner(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        @Override
        public void run(ApplicationArguments args) throws Exception {
            log.info("===== Step 02 시작 =====");

            new SendOverloads(template).run();          // 2-1
            new FutureBasics(template).run();           // 2-2
            new HeaderDemo(template).run();             // 2-6
            new PartitionInfoDemo(template).run();      // 2-8

            if (SILENT_LOSS_DEMO) {
                new SilentLossDemo(template).run();     // 2-3
            }
            if (BENCH) {
                new ThroughputBench(template).run();    // 2-4
            }
            if (PARTITION_MAP) {
                new PartitionMapDemo(template).run();   // 2-7
            }

            log.info("===== Step 02 끝 =====");
        }
    }

    // ========================================================================
    // [2-1] send 오버로드 5종
    // ========================================================================
    public static class SendOverloads {

        private static final Logger log = LoggerFactory.getLogger(SendOverloads.class);
        private final KafkaTemplate<String, OrderCreated> template;

        public SendOverloads(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        public void run() {
            OrderCreated event = OrderCreated.of(1);

            // ① 토픽 + 값. 키가 null → 파티셔너가 알아서(sticky) 고른다
            template.send(TOPIC, event);

            // ② 토픽 + 키 + 값. ★ 기본형. 같은 키 = 같은 파티션 = 순서 보장
            template.send(TOPIC, event.orderId(), event);

            // ③ 파티션까지 지정. 파티셔너를 건너뛴다 (2-8 의 위험 참고)
            template.send(TOPIC, 1, event.orderId(), event);

            // ④ ProducerRecord. 헤더·타임스탬프까지 전부 제어 (2-6)
            ProducerRecord<String, OrderCreated> record = new ProducerRecord<>(
                    TOPIC, null, event.createdAt().toEpochMilli(), event.orderId(), event);
            record.headers().add("trace-id", "t-0001".getBytes(StandardCharsets.UTF_8));
            template.send(record);

            // ⑤ Spring Messaging 의 Message<?>
            Message<OrderCreated> message = MessageBuilder.withPayload(event)
                    .setHeader(KafkaHeaders.TOPIC, TOPIC)
                    .setHeader(KafkaHeaders.KEY, event.orderId())
                    .setHeader("trace-id", "t-0001")
                    .build();
            template.send(message);

            log.info("[2-1] 오버로드 5종 발행 요청 완료 (아직 갔는지는 모른다!)");
        }
    }

    // ========================================================================
    // [2-2] CompletableFuture<SendResult> — whenComplete / thenAccept / exceptionally
    //       Spring Kafka 3.0 부터 ListenableFuture 가 아니라 CompletableFuture
    // ========================================================================
    public static class FutureBasics {

        private static final Logger log = LoggerFactory.getLogger(FutureBasics.class);
        private final KafkaTemplate<String, OrderCreated> template;

        public FutureBasics(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        public void run() throws Exception {
            OrderCreated event = OrderCreated.of(3);

            // (1) whenComplete — 성공/실패를 한곳에서. ★ 가장 자주 쓰는 형태
            CompletableFuture<SendResult<String, OrderCreated>> future =
                    template.send(TOPIC, event.orderId(), event);

            log.info("[2-2] send() 리턴 직후 isDone={}", future.isDone());   // 거의 항상 false

            future.whenComplete((result, ex) -> {
                if (ex != null) {
                    log.error("[2-2] 발행 실패 key={}", event.orderId(), ex);
                    return;
                }
                RecordMetadata md = result.getRecordMetadata();
                // ⚠️ 이 로그의 스레드명은 [main] 이 아니라 [ad | producer-1] 입니다.
                //    콜백은 프로듀서 I/O 스레드에서 실행됩니다. 무거운 일을 하면 전체 전송이 막힙니다.
                log.info("[2-2] 발행 성공 key={} → {}-{}@{} ts={} bytes={}",
                        event.orderId(), md.topic(), md.partition(), md.offset(),
                        md.timestamp(), md.serializedValueSize());
            });

            // (2) thenAccept + exceptionally — 성공/실패를 분리해서 쓰고 싶을 때
            template.send(TOPIC, "ORD-0004", OrderCreated.of(4))
                    .thenAccept(r -> log.info("[2-2] thenAccept offset={}", r.getRecordMetadata().offset()))
                    .exceptionally(ex -> {
                        log.error("[2-2] exceptionally", ex);
                        return null;
                    });

            // (3) join() — 동기 대기. 여기서만 예외가 호출 스레드로 전파된다.
            //     ⚠️ 이걸 루프 안에서 하면 2-4 의 성능 붕괴가 일어납니다.
            SendResult<String, OrderCreated> sync =
                    template.send(TOPIC, "ORD-0005", OrderCreated.of(5)).get(10, TimeUnit.SECONDS);
            log.info("[2-2] 동기 확인 → {}-{}@{}",
                    sync.getRecordMetadata().topic(),
                    sync.getRecordMetadata().partition(),
                    sync.getRecordMetadata().offset());

            // (4) 보낸 원본도 꺼낼 수 있다
            log.info("[2-2] 보낸 레코드 key={} value={}",
                    sync.getProducerRecord().key(), sync.getProducerRecord().value().orderId());
        }
    }

    // ========================================================================
    // [2-3] ⚠️ 최대 함정 — 결과를 안 보면 유실을 모른다
    //
    //  실행 전에 반드시:  docker compose stop kafka
    //  실행 후  되돌리기: docker compose start kafka
    // ========================================================================
    public static class SilentLossDemo {

        private static final Logger log = LoggerFactory.getLogger(SilentLossDemo.class);
        private final KafkaTemplate<String, OrderCreated> template;

        public SilentLossDemo(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        public void run() throws Exception {
            log.warn("[2-3] 브로커가 꺼져 있는지 확인하세요. (docker compose stop kafka)");

            // (A) 나쁜 코드 — 반환값을 버린다. 실무에서 가장 흔하다.
            for (int i = 1; i <= 5; i++) {
                OrderCreated event = OrderCreated.of(i);
                template.send(TOPIC, event.orderId(), event);      // 반환값 버림
                log.info("[2-3][나쁜코드] 주문 이벤트 발행 완료: {}", event.orderId());
                // 실무에서는 여기서 orderRepository.markPublished(...) 같은 DB 갱신이 일어난다.
                // 브로커가 죽어 있어도 이 줄은 실행된다. 그게 유실의 정체다.
                Thread.sleep(1000);
            }

            // (B) 좋은 코드 — 콜백으로 실패를 잡는다.
            AtomicInteger ok = new AtomicInteger();
            AtomicInteger fail = new AtomicInteger();
            for (int i = 6; i <= 10; i++) {
                OrderCreated event = OrderCreated.of(i);
                template.send(TOPIC, event.orderId(), event)
                        .whenComplete((r, ex) -> {
                            if (ex != null) {
                                fail.incrementAndGet();
                                log.error("[2-3][좋은코드] 발행 실패 key={} cause={}",
                                        event.orderId(), ex.getClass().getSimpleName());
                                // 여기서 Outbox 테이블 적재 / 재시도 큐 투입 / 알림을 해야 한다.
                            } else {
                                ok.incrementAndGet();
                            }
                        });
            }

            // delivery.timeout.ms(10초) 가 지나야 실패가 확정된다. 넉넉히 기다린다.
            log.info("[2-3] delivery.timeout.ms 경과를 기다립니다 (약 12초)...");
            Thread.sleep(12_000);
            log.info("[2-3] 콜백 집계 — 성공 {}건 / 실패 {}건", ok.get(), fail.get());

            // (C) 존재하지 않는 토픽 — future 상태를 시간 순으로 관찰
            CompletableFuture<SendResult<String, OrderCreated>> f =
                    template.send("no-such-topic", "K-1", OrderCreated.of(1));
            log.info("[2-3] no-such-topic: send() 직후 isDone={}", f.isDone());
            try {
                f.get(15, TimeUnit.SECONDS);
            } catch (ExecutionException e) {
                log.error("[2-3] no-such-topic: 결국 실패 — {}", e.getCause().toString());
            } catch (Exception e) {
                log.error("[2-3] no-such-topic: 15초 안에도 안 끝남 (isDone={})", f.isDone());
            }

            log.warn("[2-3] 실습 끝. docker compose start kafka 로 브로커를 되살리세요.");
        }
    }

    // ========================================================================
    // [2-4] ⚠️ 두 번째 함정 — send().get() 의 성능 붕괴
    //       4가지 방식으로 10,000건씩 발행해 처리량을 비교한다.
    // ========================================================================
    public static class ThroughputBench {

        private static final Logger log = LoggerFactory.getLogger(ThroughputBench.class);
        private static final int N = 10_000;

        private final KafkaTemplate<String, OrderCreated> template;

        public ThroughputBench(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        public void run() throws Exception {
            warmUp();   // ★ 이걸 빼면 (a) 가 (c) 보다 느리게 나오는 황당한 결과가 나온다

            long a = fireAndForget();
            long b = withCallback();
            long c = blockingGet();
            long d = batchThenJoin();

            report("(a) fire-and-forget + flush ", a);
            report("(b) 콜백(whenComplete)      ", b);
            report("(c) 매 건 send().get()      ", c);
            report("(d) 배치 후 allOf().join()  ", d);
            log.info("[2-4] (a) 대비 (c) 는 {}배 느립니다.", c / Math.max(a, 1));
        }

        /** JIT 컴파일 + 메타데이터 조회를 측정 구간 밖으로 밀어낸다. */
        private void warmUp() {
            for (int i = 1; i <= 1_000; i++) {
                template.send(TOPIC, key(i), OrderCreated.of(i));
            }
            template.flush();
            log.info("[2-4] 워밍업 1,000건 완료");
        }

        /** (a) 반환값을 버리고 루프를 돈 뒤, flush() 로 버퍼를 비우고 응답까지 기다린다. */
        private long fireAndForget() {
            long t0 = System.nanoTime();
            for (int i = 1; i <= N; i++) {
                template.send(TOPIC, key(i), OrderCreated.of(i));
            }
            // ★ flush() 는 반드시 측정 구간 안에 있어야 공정합니다.
            //   밖에 두면 "버퍼에 넣는 시간"만 재게 되어 200,000 msg/s 같은 허구가 나옵니다.
            template.flush();
            return elapsedMs(t0);
        }

        /** (b) 콜백을 붙인다. (a) 대비 손실이 거의 없다. */
        private long withCallback() {
            AtomicInteger fail = new AtomicInteger();
            long t0 = System.nanoTime();
            for (int i = 1; i <= N; i++) {
                template.send(TOPIC, key(i), OrderCreated.of(i))
                        .whenComplete((r, ex) -> {
                            if (ex != null) fail.incrementAndGet();
                        });
            }
            template.flush();
            long ms = elapsedMs(t0);
            log.info("[2-4] (b) 실패 {}건", fail.get());
            return ms;
        }

        /** (c) ★ 문제의 코드. 건마다 왕복(RTT)을 지불하므로 배치 크기가 영원히 1이 된다. */
        private long blockingGet() throws Exception {
            long t0 = System.nanoTime();
            for (int i = 1; i <= N; i++) {
                template.send(TOPIC, key(i), OrderCreated.of(i)).get();
            }
            return elapsedMs(t0);
        }

        /** (d) 전부 모아서 한 번에 대기. 개별 실패도 전부 셀 수 있다. */
        private long batchThenJoin() {
            List<CompletableFuture<SendResult<String, OrderCreated>>> futures = new ArrayList<>(N);
            long t0 = System.nanoTime();
            for (int i = 1; i <= N; i++) {
                futures.add(template.send(TOPIC, key(i), OrderCreated.of(i)));
            }
            CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new)).join();
            long ms = elapsedMs(t0);
            long failed = futures.stream().filter(CompletableFuture::isCompletedExceptionally).count();
            log.info("[2-4] (d) 발행 {}건 중 실패 {}건", N, failed);
            return ms;
        }

        private static String key(int i) {
            return "ORD-%04d".formatted(i);
        }

        private static long elapsedMs(long startNanos) {
            return (System.nanoTime() - startNanos) / 1_000_000;
        }

        private static void report(String label, long ms) {
            long tps = ms == 0 ? 0 : (N * 1000L) / ms;
            log.info("[2-4] {}: {}건 / {} ms = {} msg/s", label, N, ms, String.format("%,d", tps));
        }
    }

    // ========================================================================
    // [2-5] ProducerFactory 와 커스텀 설정 — 설정이 다른 KafkaTemplate 두 개
    //
    //  ⚠️ 이 @Configuration 을 켜면 KafkaTemplate 타입 빈이 3개(자동설정 1 + 여기 2)가 됩니다.
    //     이름 없는 @Autowired KafkaTemplate 주입은 NoUniqueBeanDefinitionException 으로 실패합니다.
    //     반드시 @Qualifier 로 골라 쓰세요.
    // ========================================================================
    @Configuration
    @Profile("step02-dual")
    public static class DualProducerConfig {

        private Map<String, Object> base(String bootstrap) {
            Map<String, Object> props = new HashMap<>();
            props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap);
            props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
            return props;
        }

        /** ① 신뢰성 우선 — 주문 이벤트처럼 절대 잃으면 안 되는 것 */
        @Bean
        public KafkaTemplate<String, OrderCreated> reliableKafkaTemplate(
                @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
            Map<String, Object> props = base(bootstrap);
            props.put(ProducerConfig.ACKS_CONFIG, "all");
            props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
            props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
            props.put(ProducerConfig.LINGER_MS_CONFIG, 5);
            props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 120_000);
            // ★ clientId 를 반드시 다르게. 안 주면 로그/JMX 에서 어느 쪽인지 구분할 수 없다.
            props.put(ProducerConfig.CLIENT_ID_CONFIG, "reliable");
            return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props));
        }

        /** ② 처리량 우선 — 좀 잃어도 되는 클릭 로그 */
        @Bean
        public KafkaTemplate<String, OrderCreated> fastKafkaTemplate(
                @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
            Map<String, Object> props = base(bootstrap);
            props.put(ProducerConfig.ACKS_CONFIG, "1");
            // ★ acks=1 이면 enable.idempotence 를 명시적으로 false 로 내려야 한다.
            //   생략하면 클라이언트가 조용히 false 로 내리지만(2-9), 의도를 코드에 남기는 편이 낫다.
            props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, false);
            props.put(ProducerConfig.LINGER_MS_CONFIG, 50);
            props.put(ProducerConfig.BATCH_SIZE_CONFIG, 65_536);
            props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
            props.put(ProducerConfig.CLIENT_ID_CONFIG, "fast");
            return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props));
        }
    }

    /** 위 두 템플릿을 @Qualifier 로 골라 쓰는 예. step02-dual 프로필에서만 동작한다. */
    @Component
    @Profile("step02-dual")
    public static class DualProducerDemo implements ApplicationRunner {

        private static final Logger log = LoggerFactory.getLogger(DualProducerDemo.class);

        private final KafkaTemplate<String, OrderCreated> reliable;
        private final KafkaTemplate<String, OrderCreated> fast;
        private final ProducerFactory<String, OrderCreated> pf;

        public DualProducerDemo(@Qualifier("reliableKafkaTemplate") KafkaTemplate<String, OrderCreated> reliable,
                                @Qualifier("fastKafkaTemplate") KafkaTemplate<String, OrderCreated> fast,
                                ProducerFactory<String, OrderCreated> pf) {
            this.reliable = reliable;
            this.fast = fast;
            this.pf = pf;
        }

        @Override
        public void run(ApplicationArguments args) {
            // 프로듀서는 싱글턴이고 스레드 안전하다 — 매번 만들지 말 것
            Producer<String, OrderCreated> p1 = pf.createProducer();
            Producer<String, OrderCreated> p2 = pf.createProducer();
            log.info("[2-5] createProducer() 두 번 → 같은 인스턴스인가? {}", p1 == p2);
            p1.close();     // DefaultKafkaProducerFactory 가 감싼 프록시라 실제로 닫히지 않는다
            p2.close();

            reliable.send(TOPIC, "ORD-0001", OrderCreated.of(1))
                    .whenComplete((r, ex) -> log.info("[2-5] reliable → {}-{}@{}",
                            r.getRecordMetadata().topic(),
                            r.getRecordMetadata().partition(),
                            r.getRecordMetadata().offset()));

            fast.send(TOPIC, "ORD-0002", OrderCreated.of(2))
                    .whenComplete((r, ex) -> log.info("[2-5] fast → {}-{}@{}",
                            r.getRecordMetadata().topic(),
                            r.getRecordMetadata().partition(),
                            r.getRecordMetadata().offset()));

            // 기동 로그에서 확인하세요:
            //   [Producer clientId=reliable] Instantiated an idempotent producer.   ← 있음
            //   [Producer clientId=fast]     ...                                    ← 없음
        }
    }

    // ========================================================================
    // [2-6] ProducerRecord 직접 만들기 — 헤더 부착, 타임스탬프 지정
    // ========================================================================
    public static class HeaderDemo {

        private static final Logger log = LoggerFactory.getLogger(HeaderDemo.class);
        private final KafkaTemplate<String, OrderCreated> template;

        public HeaderDemo(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        public void run() {
            OrderCreated event = OrderCreated.of(7);

            ProducerRecord<String, OrderCreated> record = new ProducerRecord<>(
                    TOPIC,
                    null,                                       // partition: null → 파티셔너에게 맡김
                    event.createdAt().toEpochMilli(),           // timestamp: 이벤트 발생 시각
                    event.orderId(),                            // key
                    event                                       // value
            );
            // ★ 헤더 값은 byte[] 입니다. 인코딩을 반드시 명시하세요.
            //   플랫폼 기본 인코딩에 맡기면 로컬/서버에서 다르게 직렬화되어 한글이 깨집니다.
            record.headers()
                    .add("trace-id", "t-9f2c1a".getBytes(StandardCharsets.UTF_8))
                    .add("event-type", "OrderCreated".getBytes(StandardCharsets.UTF_8))
                    .add("schema-version", "2".getBytes(StandardCharsets.UTF_8));

            template.send(record).whenComplete((r, ex) -> {
                if (ex != null) {
                    log.error("[2-6] 발행 실패", ex);
                    return;
                }
                RecordMetadata md = r.getRecordMetadata();
                log.info("[2-6] → {}-{}@{} ts={}", md.topic(), md.partition(), md.offset(), md.timestamp());
            });

            // ⚠️ 토픽의 message.timestamp.type 이 LogAppendTime 이면
            //    위에서 지정한 타임스탬프는 브로커가 도착 시각으로 덮어씁니다. 에러도 경고도 없습니다.
            //    확인: kafka-configs.sh --describe --entity-type topics --entity-name orders
        }
    }

    // ========================================================================
    // [2-7] 키와 파티셔너 — murmur2 로 예측하고 실제와 대조
    // ========================================================================
    public static class PartitionMapDemo {

        private static final Logger log = LoggerFactory.getLogger(PartitionMapDemo.class);
        private final KafkaTemplate<String, OrderCreated> template;

        public PartitionMapDemo(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        public void run() throws Exception {
            int numPartitions = template.partitionsFor(TOPIC).size();
            log.info("[2-7] {} 의 파티션 수 = {} (3이 아니면 교재의 표와 어긋납니다)", TOPIC, numPartitions);

            // (A) 키가 있을 때 — 결정적. 브로커에 안 물어보고 계산할 수 있다.
            for (int i = 1; i <= 9; i++) {
                OrderCreated event = OrderCreated.of(i);
                String key = event.orderId();
                int predicted = predictPartition(key, numPartitions);

                template.send(TOPIC, key, event).whenComplete((r, ex) -> {
                    if (ex != null) {
                        log.error("[2-7] 발행 실패 key={}", key, ex);
                        return;
                    }
                    int actual = r.getRecordMetadata().partition();
                    log.info("[2-7] {} 예측={} 실제={} offset={} {}",
                            key, predicted, actual, r.getRecordMetadata().offset(),
                            predicted == actual ? "" : "  ← ★불일치! 파티션 수를 확인하세요");
                });
            }
            template.flush();
            Thread.sleep(500);

            // (B) 파티션 수를 바꾸면 매핑이 통째로 바뀐다 — 순서 보장이 소급해서 깨진다
            log.info("[2-7] --- 파티션 수별 매핑 비교 (표 재현) ---");
            log.info("[2-7] {}", "key       murmur2        toPositive     %3  %4  %6");
            for (int i = 1; i <= 9; i++) {
                String key = "ORD-%04d".formatted(i);
                int hash = Utils.murmur2(key.getBytes(StandardCharsets.UTF_8));
                int pos = Utils.toPositive(hash);
                // SLF4J 의 {} 는 정렬 지정자를 지원하지 않으므로 String.format 으로 미리 만든다
                log.info("[2-7] {}", String.format("%-9s %,13d  %,13d   %d   %d   %d",
                        key, hash, pos, pos % 3, pos % 4, pos % 6));
            }

            // (C) 키가 null 이면 — 3.3+ 내장 uniform sticky.
            //     소량이면 "골고루"가 아니라 한 파티션에 몰립니다.
            log.info("[2-7] --- 키 없이 6건 (sticky 확인) ---");
            for (int i = 1; i <= 6; i++) {
                template.send(TOPIC, OrderCreated.of(i)).whenComplete((r, ex) -> {
                    if (ex == null) {
                        log.info("[2-7] key=null → {}-{}@{}",
                                r.getRecordMetadata().topic(),
                                r.getRecordMetadata().partition(),
                                r.getRecordMetadata().offset());
                    }
                });
            }
            template.flush();
        }

        /**
         * Kafka 기본 파티셔너와 동일한 계산.
         * ★ Math.abs(hash) % n 으로 쓰면 안 됩니다.
         *   Math.abs(Integer.MIN_VALUE) 는 여전히 음수라서 음수 파티션이 나옵니다.
         *   Kafka 가 toPositive(= x & 0x7fffffff)를 쓰는 이유가 정확히 이것입니다.
         */
        public static int predictPartition(String key, int numPartitions) {
            byte[] bytes = key.getBytes(StandardCharsets.UTF_8);
            return Utils.toPositive(Utils.murmur2(bytes)) % numPartitions;
        }
    }

    // ========================================================================
    // [2-8] 명시적 파티션 지정과 그 위험
    // ========================================================================
    public static class PartitionInfoDemo {

        private static final Logger log = LoggerFactory.getLogger(PartitionInfoDemo.class);
        private final KafkaTemplate<String, OrderCreated> template;

        public PartitionInfoDemo(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        public void run() {
            // ★ 파티션 수를 코드에 박지 말고 런타임에 조회한다.
            //   단, partitionsFor 는 메타데이터가 없으면 max.block.ms 까지 블로킹합니다.
            //   요청 경로에서 매번 부르지 말고 기동 시 한 번 조회해 캐시하세요.
            List<PartitionInfo> infos = template.partitionsFor(TOPIC);
            log.info("[2-8] {} 파티션 수 = {}", TOPIC, infos.size());
            infos.stream()
                    .sorted((x, y) -> Integer.compare(x.partition(), y.partition()))
                    .forEach(p -> log.info("[2-8]   partition={} leader={}", p.partition(), p.leader().id()));

            // 정당한 용도: 재처리, 테스트, 파티션 전용 워크로드
            template.send(TOPIC, 2, "ORD-0001", OrderCreated.of(1))
                    .whenComplete((r, ex) -> log.info("[2-8] 명시적 파티션 2 → 실제 {}",
                            r.getRecordMetadata().partition()));

            // ⚠️ 범위를 벗어나면 IllegalArgumentException 이 future 에 담깁니다.
            //    이건 시끄러운 실패라 다행입니다. 진짜 위험한 건 파티션을 늘렸는데
            //    코드가 % 3 으로 하드코딩돼 있어 뒤쪽 파티션이 영원히 비는 경우입니다.
            template.send(TOPIC, 99, "ORD-0002", OrderCreated.of(2))
                    .whenComplete((r, ex) -> {
                        if (ex != null) {
                            log.error("[2-8] 범위 밖 파티션 → {}", ex.getCause() == null
                                    ? ex.toString() : ex.getCause().toString());
                        }
                    });
        }
    }

    // ========================================================================
    // [2-9] acks / retries / enable.idempotence 의 상호 제약
    //
    //  아래 설정을 application-step02.yml 에 넣고 기동해 보세요.
    //
    //  (A) 기동 실패 — 명시적으로 켠 멱등성이 acks 와 충돌
    //      spring.kafka.producer.acks: 1
    //      spring.kafka.producer.properties.enable.idempotence: true
    //      → org.apache.kafka.common.config.ConfigException:
    //          Must set acks to all in order to use the idempotent producer.
    //
    //  (B) ⚠️ 진짜 함정 — 조용히 꺼짐
    //      spring.kafka.producer.acks: 1        (enable.idempotence 는 안 씀)
    //      → 기동 성공. 그러나 멱등성은 false 로 내려간다. 경고 한 줄 없다.
    //      → retries 로 인한 중복 + max.in.flight=5 로 인한 순서 역전이 가능해진다.
    //      확인법: 기동 로그에서 다음 한 줄을 찾으세요. 없으면 꺼진 것입니다.
    //          [Producer clientId=producer-1] Instantiated an idempotent producer.
    //
    //  다른 위반 조합들:
    //      enable.idempotence=true + retries=0        → Must set retries to non-zero ...
    //      enable.idempotence=true + max.in.flight>5  → Must set max.in.flight ... at most 5 ...
    //
    //  브로커 쪽 ack/ISR 동작은 이 코스가 아니라 Kafka 코스 소관입니다:
    //      ../../kafka/step-04-producer/
    //      ../../kafka/step-07-delivery-semantics/
    //      ../../kafka/step-08-replication/
    // ========================================================================
    public static final class IdempotenceNotes {
        private IdempotenceNotes() {
        }
    }
}

Exercise.java

6문제의 문제지입니다. 각 문제는 // 여기에 작성: 자리를 비워 두었고, 뼈대는 그대로 컴파일됩니다.

  • 문제 1·5 는 콜백을 붙여 관찰하는 문제, 문제 2·3 은 측정과 상태 추적, 문제 4·6설정을 직접 만들어 결과를 확인하는 문제입니다.
  • 문제 3 은 isDone() / isCompletedExceptionally()1초 간격으로 폴링해 상태 변화 시점을 기록하는 구조입니다. 실행 전에 kafka-configs.sh 로 자동 토픽 생성을 끄지 않으면 토픽이 생겨 버려 문제가 성립하지 않으니, 파일 상단 주석의 명령을 그대로 따르세요(실습 후 --delete-config 로 되돌리는 명령도 함께 적혀 있습니다).
  • 문제 6 은 일부러 기동을 실패시키는 문제입니다. @Profile("step02-ex6") 로 격리해 두었으므로 이 프로필로 켜면 앱이 뜨지 않는 것이 정상입니다. 다른 문제를 풀 때 이 프로필을 켜지 마세요. 또한 DefaultKafkaProducerFactory지연 생성이라 빈 등록만으로는 예외가 안 납니다 — createProducer() 를 강제로 한 번 호출해야 ConfigException 을 볼 수 있다는 힌트가 주석에 있습니다.
  • ⚠️ 문제 2 의 벤치마크는 orders 토픽에 약 40,000건(4방식 × 10,000)을 밀어 넣습니다. 뒤 스텝 실습에 방해되면 docker compose down -v && docker compose up -d 로 초기화하세요.
  • 각 문제 아래에 기대 출력 예시가 주석으로 붙어 있습니다. 형태가 다르면 답이 틀린 것이니 Solution.java 를 열기 전에 다시 보세요.
/*
 * ============================================================================
 *  Step 02 — KafkaTemplate 과 프로듀서 : Exercise (문제지)
 * ============================================================================
 *
 *  배치 위치 : src/main/java/com/example/order/step02/Exercise.java
 *  실행      : ./gradlew bootRun --args='--spring.profiles.active=step02-ex'
 *
 *  각 문제의 "여기에 작성:" 자리를 채우세요. 뼈대는 그대로 컴파일됩니다.
 *  정답은 Solution.java 에 있습니다. 반드시 먼저 직접 풀어 보세요.
 *
 *  ⚠️ 사전 준비
 *   - 문제 3 은 "존재하지 않는 토픽"이 정말 존재하지 않아야 성립합니다.
 *     docker-compose.yml 의 KAFKA_AUTO_CREATE_TOPICS_ENABLE 이 "true" 이면
 *     send 하는 순간 토픽이 생겨 버립니다. 아래로 잠시 끄세요.
 *
 *       docker exec -it learn-kafka /opt/kafka/bin/kafka-configs.sh \
 *         --bootstrap-server localhost:9092 --entity-type brokers --entity-name 1 \
 *         --alter --add-config auto.create.topics.enable=false
 *
 *     실습 후 되돌리기:
 *       docker exec -it learn-kafka /opt/kafka/bin/kafka-configs.sh \
 *         --bootstrap-server localhost:9092 --entity-type brokers --entity-name 1 \
 *         --alter --delete-config auto.create.topics.enable
 *
 *   - 문제 2 의 벤치마크는 orders 토픽에 약 40,000건(4방식 × 10,000)을 밀어 넣습니다.
 *     뒤 스텝 실습에 방해되면 docker compose down -v && docker compose up -d 로 초기화하세요.
 *
 *   - 문제 6 은 "일부러 기동을 실패시키는" 문제입니다. @Profile("step02-ex6") 로 격리했습니다.
 *     다른 문제를 풀 때는 이 프로필을 켜지 마세요.
 * ============================================================================
 */
package com.example.order.step02;

import com.example.order.domain.OrderCreated;

import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.common.utils.Utils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.kafka.support.serializer.JsonSerializer;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;

public final class Exercise {

    private Exercise() {
    }

    private static final String TOPIC = "orders";

    @Component
    @Profile("step02-ex")
    public static class ExerciseRunner implements ApplicationRunner {

        private static final Logger log = LoggerFactory.getLogger(ExerciseRunner.class);

        private final KafkaTemplate<String, OrderCreated> template;

        public ExerciseRunner(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        @Override
        public void run(ApplicationArguments args) throws Exception {
            problem1();
            problem2();
            problem3();
            problem5();
            // 문제 4 는 Problem4Config / Problem4Demo 를 보세요 (step02-ex4 프로필)
            // 문제 6 은 Problem6Config 를 보세요 (step02-ex6 프로필)
        }

        // ====================================================================
        // 문제 1. 콜백으로 파티션/오프셋 로깅하기
        //
        //  요구사항:
        //   - ORD-0001 ~ ORD-0005 를 orders 토픽에 키를 붙여 발행한다.
        //   - 각 발행의 결과를 콜백으로 받아 다음 형식으로 로깅한다.
        //         key=ORD-0001 → orders-2@0 ts=1735689660000
        //   - 실패한 경우에도 반드시 ERROR 로 남긴다. (성공/실패 둘 다 잡아야 한다)
        //   - 힌트: thenAccept 는 성공만, exceptionally 는 실패만 잡는다.
        //           둘 다 필요하면 어떤 메서드를 써야 하는가?
        //   - 마지막에 flush() 로 버퍼를 비워 콜백이 다 돌 때까지 기다린다.
        //
        //  기대 출력 예:
        //   INFO ... [ad | producer-1] : key=ORD-0001 → orders-2@0 ts=1735689660000
        //   INFO ... [ad | producer-1] : key=ORD-0002 → orders-2@1 ts=1735689720000
        //   INFO ... [ad | producer-1] : key=ORD-0003 → orders-0@0 ts=1735689780000
        // ====================================================================
        void problem1() {
            log.info("===== 문제 1 =====");

            for (int i = 1; i <= 5; i++) {
                OrderCreated event = OrderCreated.of(i);

                // 여기에 작성:

            }

            template.flush();
        }

        // ====================================================================
        // 문제 2. 10,000건 발행 처리량 측정 비교
        //
        //  요구사항:
        //   - 아래 4가지 방식으로 각각 10,000건을 발행하고 소요 ms 와 msg/s 를 로깅한다.
        //       (a) fire-and-forget + flush()
        //       (b) whenComplete 콜백 + flush()
        //       (c) 매 건 send().get()
        //       (d) 전부 모아서 CompletableFuture.allOf(...).join()
        //   - 측정 전에 워밍업 1,000건을 먼저 보낸다.
        //     ★ 워밍업을 빼면 (a) 가 (c) 보다 느리게 나오는 황당한 결과가 나옵니다. 왜일까요?
        //   - ★ flush() 를 측정 구간의 "안"에 둘 것인가 "밖"에 둘 것인가?
        //     둘 다 해 보고 숫자가 어떻게 달라지는지 확인하세요. 어느 쪽이 공정한 측정입니까?
        //   - 마지막에 (a) 대비 (c) 가 몇 배 느린지 출력한다.
        //
        //  기대 출력 예:
        //   INFO ... : (a) fire-and-forget + flush : 10000건 /   543 ms = 18,416 msg/s
        //   INFO ... : (c) 매 건 send().get()      : 10000건 / 10870 ms =    920 msg/s
        //   INFO ... : (a) 대비 (c) 는 20배 느립니다.
        // ====================================================================
        void problem2() throws Exception {
            log.info("===== 문제 2 =====");
            final int N = 10_000;

            // 워밍업
            // 여기에 작성:

            // (a) fire-and-forget
            long elapsedA = 0;
            // 여기에 작성:

            // (b) 콜백
            long elapsedB = 0;
            // 여기에 작성:

            // (c) 매 건 get()
            long elapsedC = 0;
            // 여기에 작성:

            // (d) allOf
            long elapsedD = 0;
            // 여기에 작성:

            log.info("(a)={}ms (b)={}ms (c)={}ms (d)={}ms", elapsedA, elapsedB, elapsedC, elapsedD);
            // 여기에 작성: msg/s 환산과 배율 출력
        }

        // ====================================================================
        // 문제 3. 존재하지 않는 토픽에 send 했을 때 future 상태 관찰
        //
        //  요구사항:
        //   - "no-such-topic-s02" 로 한 건 발행한다. (자동 토픽 생성을 반드시 꺼 둘 것)
        //   - send() 리턴 "직후"의 isDone() 을 로깅한다.
        //   - 그 뒤 1초 간격으로 최대 70초 동안 isDone() / isCompletedExceptionally() 를
        //     폴링해, 상태가 바뀌는 "정확한 시점(초)"을 로깅한다.
        //   - 최종적으로 어떤 예외가 담겼는지 클래스명과 메시지를 출력한다.
        //   - 이 70초 동안 "호출부는 무엇을 믿고 있었는지"를 주석으로 한 줄 적어 보세요.
        //
        //  기대 출력 예:
        //   INFO ... : send() 직후 isDone=false
        //   INFO ... : t=  1s isDone=false
        //   ...
        //   INFO ... : t= 60s isDone=true  exceptionally=true
        //   ERROR... : 담긴 예외 = TimeoutException: Topic no-such-topic-s02 not present in metadata after 60000 ms.
        // ====================================================================
        void problem3() throws Exception {
            log.info("===== 문제 3 =====");

            CompletableFuture<SendResult<String, OrderCreated>> future =
                    template.send("no-such-topic-s02", "K-1", OrderCreated.of(1));

            // 여기에 작성: send() 직후 isDone 로깅

            for (int t = 1; t <= 70; t++) {
                Thread.sleep(1000);
                // 여기에 작성: 상태 폴링 + 변화 시점 로깅 + 완료되면 break

            }

            // 여기에 작성: 담긴 예외 출력
        }

        // ====================================================================
        // 문제 5. 키 → 파티션 매핑을 직접 계산해 예측하고 검증
        //
        //  요구사항:
        //   - ORD-0010 ~ ORD-0015 의 파티션을 "발행하기 전에" 계산해 로깅한다.
        //     계산식은 Kafka 기본 파티셔너와 동일해야 한다.
        //   - ★ Math.abs(hash) % n 을 쓰면 안 됩니다. 왜일까요? (힌트: Integer.MIN_VALUE)
        //     Kafka 는 무엇을 쓰는지 org.apache.kafka.common.utils.Utils 에서 찾아보세요.
        //   - 그 다음 실제로 발행해 콜백에서 실제 파티션을 받아 예측값과 대조한다.
        //   - 파티션 수는 하드코딩하지 말고 template.partitionsFor(TOPIC) 로 조회한다.
        //   - 추가: 파티션이 4개였다면 어디로 갔을지도 함께 출력해,
        //     "파티션 수를 바꾸면 매핑이 통째로 바뀐다"를 눈으로 확인한다.
        //
        //  기대 출력 예:
        //   INFO ... : ORD-0010 예측(3개)=1 예측(4개)=2 | 실제=1  ← 일치
        // ====================================================================
        void problem5() throws Exception {
            log.info("===== 문제 5 =====");

            int numPartitions = 0;
            // 여기에 작성: 파티션 수 조회

            for (int i = 10; i <= 15; i++) {
                String key = "ORD-%04d".formatted(i);

                // 여기에 작성: 예측값 계산 (3개일 때 / 4개일 때)

                // 여기에 작성: 발행 + 콜백에서 예측/실제 대조 로깅

            }
            template.flush();
            Thread.sleep(500);
        }

        /**
         * 문제 5 의 보조 메서드. Kafka 기본 파티셔너와 동일하게 구현하세요.
         */
        static int predictPartition(String key, int numPartitions) {
            byte[] bytes = key.getBytes(StandardCharsets.UTF_8);
            // 여기에 작성:
            return 0;
        }
    }

    // ========================================================================
    // 문제 4. 서로 다른 설정의 KafkaTemplate 두 개 만들기
    //
    //  요구사항:
    //   - "auditKafkaTemplate": acks=all, enable.idempotence=true, linger.ms=0,
    //                           clientId="audit"
    //   - "bulkKafkaTemplate" : acks=1,  enable.idempotence=false, linger.ms=100,
    //                           batch.size=131072, compression.type="lz4", clientId="bulk"
    //   - 두 템플릿을 @Qualifier 로 골라 주입받아 각각 한 건씩 발행하고,
    //     콜백에서 어느 템플릿인지 알 수 있게 로깅한다.
    //   - ★ 이 두 빈을 등록하면 KafkaTemplate 타입 빈이 3개가 됩니다(자동설정 1 + 여기 2).
    //     이름 없이 주입하면 어떤 예외로 기동이 실패합니까? 일부러 한 번 재현해 보세요.
    //   - ★ 한쪽에 @Primary 를 붙이면 문제가 해결됩니다. 그런데 왜 권장하지 않을까요?
    //     주석으로 한 줄 적어 보세요.
    //   - 기동 로그에서 "Instantiated an idempotent producer." 가 어느 clientId 에만
    //     찍히는지 확인하세요.
    //
    //  실행: ./gradlew bootRun --args='--spring.profiles.active=step02-ex4'
    // ========================================================================
    @Configuration
    @Profile("step02-ex4")
    public static class Problem4Config {

        private Map<String, Object> base(String bootstrap) {
            Map<String, Object> props = new HashMap<>();
            props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap);
            props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
            return props;
        }

        @Bean
        public KafkaTemplate<String, OrderCreated> auditKafkaTemplate(
                @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
            Map<String, Object> props = base(bootstrap);
            // 여기에 작성:

            return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props));
        }

        @Bean
        public KafkaTemplate<String, OrderCreated> bulkKafkaTemplate(
                @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
            Map<String, Object> props = base(bootstrap);
            // 여기에 작성:

            return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props));
        }
    }

    @Component
    @Profile("step02-ex4")
    public static class Problem4Demo implements ApplicationRunner {

        private static final Logger log = LoggerFactory.getLogger(Problem4Demo.class);

        // 여기에 작성: @Qualifier 로 두 템플릿을 주입받는 생성자

        @Override
        public void run(ApplicationArguments args) {
            log.info("===== 문제 4 =====");
            // 여기에 작성: 각 템플릿으로 한 건씩 발행 + 콜백 로깅

        }
    }

    // ========================================================================
    // 문제 6. acks=1 + enable.idempotence=true 조합의 에러 재현
    //
    //  요구사항:
    //   - (A) acks="1" 과 enable.idempotence=true 를 "동시에 명시"한 ProducerFactory 로
    //         KafkaTemplate 빈을 만들고 기동한다.
    //         → 어떤 예외로 기동이 실패합니까? 예외 클래스와 메시지를 그대로 옮겨 적으세요.
    //   - (B) 이번엔 enable.idempotence 를 "아예 설정하지 않고" acks="1" 만 준다.
    //         → 기동이 성공합니까 실패합니까? 기동 로그에
    //           "Instantiated an idempotent producer." 가 찍힙니까?
    //   - ★ (A) 와 (B) 중 운영에서 더 위험한 쪽은 어디이고, 왜 그렇습니까?
    //     주석으로 3줄 이상 적어 보세요.
    //   - 추가로 다음 두 조합도 확인해 보세요.
    //         enable.idempotence=true + retries=0
    //         enable.idempotence=true + max.in.flight.requests.per.connection=6
    //
    //  실행: ./gradlew bootRun --args='--spring.profiles.active=step02-ex6'
    //  ⚠️ 이 프로필은 기동에 실패하는 것이 "정상"입니다.
    // ========================================================================
    @Configuration
    @Profile("step02-ex6")
    public static class Problem6Config {

        @Bean
        public KafkaTemplate<String, OrderCreated> brokenKafkaTemplate(
                @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
            Map<String, Object> props = new HashMap<>();
            props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap);
            props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
            props.put(ProducerConfig.CLIENT_ID_CONFIG, "broken");

            // 여기에 작성: (A) acks=1 + enable.idempotence=true 를 동시에 명시

            DefaultKafkaProducerFactory<String, OrderCreated> pf =
                    new DefaultKafkaProducerFactory<>(props);
            KafkaTemplate<String, OrderCreated> template = new KafkaTemplate<>(pf);
            // ★ DefaultKafkaProducerFactory 는 지연 생성이라 빈 생성만으로는 예외가 안 납니다.
            //   실제 KafkaProducer 를 만들게 강제해야 ConfigException 을 볼 수 있습니다.
            //   힌트: pf.createProducer() 를 여기서 한 번 호출해 보세요.
            // 여기에 작성:

            return template;
        }
    }

    /** 문제 5 채점용 참고 상수 — 계산이 맞았는지 대조해 보세요 (파티션 3개 기준). */
    static final class Hint {
        private Hint() {
        }

        static void printMurmur2(String key) {
            int hash = Utils.murmur2(key.getBytes(StandardCharsets.UTF_8));
            System.out.printf("%s murmur2=%d toPositive=%d%n", key, hash, Utils.toPositive(hash));
        }
    }

    /** 문제 2 에서 쓸 수 있는 보조 컨테이너. 필요하면 쓰세요. */
    static final class Bench {
        private Bench() {
        }

        static List<CompletableFuture<SendResult<String, OrderCreated>>> newList(int n) {
            return new ArrayList<>(n);
        }

        static long elapsedMs(long startNanos) {
            return (System.nanoTime() - startNanos) / 1_000_000;
        }
    }
}

Solution.java

6문제의 정답 코드와, "왜 그 답인가"를 설명하는 긴 블록 주석이 함께 들어 있습니다. 풀어 본 뒤에 여세요.

  • 정답 1whenComplete 를 씁니다. thenAccept 는 성공 시에만 실행돼 실패를 놓치고, exceptionally 는 실패 시에만 실행돼 성공을 놓칩니다. 둘 다 로깅해야 하므로 whenComplete 가 유일한 답이며, 체이닝으로 흉내 내면 "발행은 성공했는데 실패 로그가 찍히는" 혼란이 생긴다는 점까지 설명합니다.
  • 정답 2 는 워밍업과 flush() 위치가 핵심입니다. (a) 는 flush()측정 구간 안에 넣어야 공정합니다. 밖에 두면 "버퍼에 넣는 시간"만 재게 되어 200,000 msg/s 같은 무의미한 숫자 — 사실상 ArrayDeque 삽입 속도 — 가 나옵니다. 그 이유를 15줄로 설명합니다.
  • 정답 3 의 관찰 결과가 이 스텝의 요약입니다. send() 직후 isDone=false → 60초간 계속 falseTimeoutException 으로 isCompletedExceptionally=true. 그 60초 동안 호출부는 성공했다고 믿고 있었고, 그사이 DB 트랜잭션은 이미 커밋됐다는 것이 결론입니다.
  • 정답 4@Bean 이름을 @Qualifier 문자열로 쓰는 방식을 택합니다. @Primary 를 안 쓰는 이유를 주석이 설명합니다 — @Primary 는 "실수로 주입받은 코드"가 어느 쪽을 쓰는지를 침묵으로 결정해 버리기 때문입니다. @Qualifier 를 강제하면 새 코드는 반드시 선택해야만 컴파일됩니다.
  • 정답 5Utils.toPositive(Utils.murmur2(bytes)) % n 을 그대로 씁니다. Math.abs(hash) % n 으로 쓰면 Integer.MIN_VALUE 에서 음수가 나와 배열 인덱스 예외가 나는 고전적 버그인데, Kafka 가 Math.abs 대신 toPositive 를 쓰는 이유가 정확히 이것입니다.
  • 정답 6 은 두 경우를 나란히 보여줍니다. enable.idempotence=true + acks=1ConfigException 기동 실패, acks=1 만 → 기동 성공 + 멱등성 조용히 off. 후자가 더 위험한 이유를 배포 파이프라인이 막히는 실패와 며칠 뒤 재고 음수로 드러나는 실패의 대비로 설명하고, 기동 로그에서 Instantiated an idempotent producer. 를 확인하는 습관을 강조합니다.
/*
 * ============================================================================
 *  Step 02 — KafkaTemplate 과 프로듀서 : Solution (정답과 해설)
 * ============================================================================
 *
 *  배치 위치 : src/main/java/com/example/order/step02/Solution.java
 *  실행      : ./gradlew bootRun --args='--spring.profiles.active=step02-sol'
 *
 *  Exercise.java 를 먼저 풀어 본 뒤에 여세요.
 *  각 정답 위의 블록 주석이 "왜 그 답인가"를 설명합니다.
 * ============================================================================
 */
package com.example.order.step02;

import com.example.order.domain.OrderCreated;

import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.common.utils.Utils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.kafka.support.serializer.JsonSerializer;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.atomic.AtomicInteger;

public final class Solution {

    private Solution() {
    }

    private static final String TOPIC = "orders";

    @Component
    @Profile("step02-sol")
    public static class SolutionRunner implements ApplicationRunner {

        private static final Logger log = LoggerFactory.getLogger(SolutionRunner.class);

        private final KafkaTemplate<String, OrderCreated> template;

        public SolutionRunner(KafkaTemplate<String, OrderCreated> template) {
            this.template = template;
        }

        @Override
        public void run(ApplicationArguments args) throws Exception {
            answer1();
            answer2();
            answer3();
            answer5();
        }

        /* ====================================================================
         * 정답 1. 콜백으로 파티션/오프셋 로깅
         *
         *  ★ whenComplete 가 유일한 답입니다.
         *
         *  CompletableFuture 의 세 메서드는 역할이 다릅니다.
         *    - thenAccept(Consumer<T>)          : "성공했을 때만" 실행됩니다.
         *                                          실패하면 아예 호출되지 않으므로 유실을 놓칩니다.
         *    - exceptionally(Function<Throwable,T>) : "실패했을 때만" 실행됩니다.
         *                                          성공 시 파티션/오프셋을 찍을 수 없습니다.
         *    - whenComplete(BiConsumer<T,Throwable>) : 성공이든 실패든 "항상" 실행됩니다.
         *                                          두 인자 중 하나는 항상 null 입니다.
         *
         *  문제가 "성공 시 파티션/오프셋을 찍고, 실패 시 ERROR 를 남겨라" 이므로
         *  둘 다 잡는 whenComplete 만이 조건을 만족합니다.
         *
         *  thenAccept + exceptionally 를 체이닝해도 동작은 하지만, 두 번째 함정이 있습니다.
         *  thenAccept 안에서 예외가 나면 그것까지 exceptionally 로 흘러가서
         *  "발행은 성공했는데 발행 실패 로그가 찍히는" 혼란이 생깁니다.
         *
         *  마지막으로 flush() 를 부르는 이유는 ApplicationRunner 가 끝나 버리면
         *  콜백이 돌기 전에 로그를 못 볼 수 있기 때문입니다. (프로듀서 I/O 스레드는 데몬입니다.)
         * ==================================================================== */
        void answer1() {
            log.info("===== 정답 1 =====");

            for (int i = 1; i <= 5; i++) {
                OrderCreated event = OrderCreated.of(i);
                template.send(TOPIC, event.orderId(), event)
                        .whenComplete((result, ex) -> {
                            if (ex != null) {
                                log.error("발행 실패 key={} cause={}",
                                        event.orderId(), ex.getClass().getSimpleName(), ex);
                                return;
                            }
                            RecordMetadata md = result.getRecordMetadata();
                            log.info("key={} → {}-{}@{} ts={}",
                                    event.orderId(), md.topic(), md.partition(),
                                    md.offset(), md.timestamp());
                        });
            }
            template.flush();
        }

        /* ====================================================================
         * 정답 2. 10,000건 처리량 측정
         *
         *  ★ 이 문제의 채점 포인트는 숫자가 아니라 "측정 설계" 두 가지입니다.
         *
         *  (1) 워밍업이 왜 필요한가
         *      첫 send() 는 브로커에서 토픽 메타데이터를 가져오느라 수백 ms 를 씁니다.
         *      JIT 도 아직 인터프리터 모드입니다. 워밍업 없이 (a)→(b)→(c) 순으로 재면
         *      (a) 가 그 비용을 전부 뒤집어써서, 최악의 방식인 (c) 보다 (a) 가
         *      느리게 나오는 결과가 나옵니다. 순서를 바꿔도 마찬가지라 결론을 못 냅니다.
         *      워밍업 1,000건이면 메타데이터 조회와 주요 경로의 JIT 컴파일이 끝납니다.
         *
         *  (2) flush() 는 반드시 "측정 구간 안"에 있어야 한다
         *      send() 는 버퍼에 넣고 즉시 리턴합니다. flush() 를 측정 밖에 두면
         *      "메모리에 10,000번 넣는 시간"만 재게 되어 200,000 msg/s 같은 숫자가 나옵니다.
         *      그건 발행 처리량이 아니라 ArrayDeque 삽입 속도입니다.
         *      (c) 는 구조상 브로커 응답까지 기다리므로, (a) 도 같은 지점까지 기다려야
         *      비교가 성립합니다. 그래서 flush() 를 안에 둡니다.
         *
         *  실측 (Java 21 / kafka-clients 3.6.1 / 단일 브로커 / linger.ms=5):
         *      (a) fire-and-forget + flush :   543 ms → 18,416 msg/s
         *      (b) 콜백(whenComplete)      :   552 ms → 18,116 msg/s
         *      (c) 매 건 send().get()      : 10870 ms →    920 msg/s
         *      (d) 배치 후 allOf().join()  :   559 ms → 17,889 msg/s
         *
         *  (c) 가 20배 느린 이유는 배치가 죽기 때문입니다. get() 이 블로킹하는 동안
         *  호출 스레드는 다음 send() 를 못 하므로 버퍼에 두 번째 레코드가 들어올 수 없고,
         *  배치 크기가 영원히 1이 됩니다. 게다가 linger.ms=5 가 이제 이득이 아니라
         *  건당 5ms 의 순수 페널티가 됩니다. linger.ms=0 으로 내려도 660 msg/s 대라
         *  본질은 linger 가 아니라 "왕복(RTT)을 건마다 지불하는 구조"입니다.
         *
         *  (b) 가 (a) 와 거의 같은 것이 중요합니다. 콜백은 사실상 공짜입니다.
         *  즉 "안전(유실 감지)"과 "속도" 중 하나를 고를 필요가 없습니다. 콜백이 정답입니다.
         * ==================================================================== */
        void answer2() throws Exception {
            log.info("===== 정답 2 =====");
            final int N = 10_000;

            // 워밍업 — 메타데이터 조회 + JIT
            for (int i = 1; i <= 1_000; i++) {
                template.send(TOPIC, key(i), OrderCreated.of(i));
            }
            template.flush();

            // (a) fire-and-forget
            long t0 = System.nanoTime();
            for (int i = 1; i <= N; i++) {
                template.send(TOPIC, key(i), OrderCreated.of(i));
            }
            template.flush();                                  // ★ 측정 구간 안
            long a = ms(t0);

            // (b) 콜백
            AtomicInteger fail = new AtomicInteger();
            long t1 = System.nanoTime();
            for (int i = 1; i <= N; i++) {
                template.send(TOPIC, key(i), OrderCreated.of(i))
                        .whenComplete((r, ex) -> {
                            if (ex != null) fail.incrementAndGet();
                        });
            }
            template.flush();
            long b = ms(t1);

            // (c) 매 건 get()
            long t2 = System.nanoTime();
            for (int i = 1; i <= N; i++) {
                template.send(TOPIC, key(i), OrderCreated.of(i)).get();
            }
            long c = ms(t2);

            // (d) allOf
            List<CompletableFuture<SendResult<String, OrderCreated>>> futures = new ArrayList<>(N);
            long t3 = System.nanoTime();
            for (int i = 1; i <= N; i++) {
                futures.add(template.send(TOPIC, key(i), OrderCreated.of(i)));
            }
            CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new)).join();
            long d = ms(t3);
            long failedInD = futures.stream().filter(CompletableFuture::isCompletedExceptionally).count();

            report("(a) fire-and-forget + flush ", N, a);
            report("(b) 콜백(whenComplete)      ", N, b);
            report("(c) 매 건 send().get()      ", N, c);
            report("(d) 배치 후 allOf().join()  ", N, d);
            log.info("(b) 실패 {}건 / (d) 실패 {}건", fail.get(), failedInD);
            log.info("(a) 대비 (c) 는 {}배 느립니다.", c / Math.max(a, 1));
        }

        /* ====================================================================
         * 정답 3. 존재하지 않는 토픽 — future 상태 관찰
         *
         *  ★ 관찰 결과가 이 스텝 전체의 요약입니다.
         *
         *      send() 직후          isDone=false
         *      t=1s ~ t=59s        isDone=false   (계속)
         *      t=60s               isDone=true, isCompletedExceptionally=true
         *      담긴 예외           TimeoutException: Topic no-such-topic-s02 not present
         *                          in metadata after 60000 ms.
         *
         *  60초입니다. 60초 동안 호출부는 "발행에 성공했다"고 믿고 있었습니다.
         *  그 사이에 DB 에 markPublished 를 커밋했다면, 그 트랜잭션은 이미 끝났습니다.
         *  되돌릴 방법이 없습니다.
         *
         *  60초라는 숫자의 출처는 delivery.timeout.ms 가 아니라 max.block.ms 도 아닌
         *  metadata 조회 대기입니다(내부적으로 max.block.ms 와 delivery.timeout.ms 중
         *  적용되는 쪽을 따르며, 기본 조합에서 60초로 관측됩니다).
         *  실습에서 delivery.timeout.ms=10000 으로 줄여 두었다면 더 빨리 끝납니다.
         *
         *  ★ 여기서 반드시 짚어야 할 것: 이 60초 동안 애플리케이션 로그는 깨끗합니다.
         *  NetworkClient 의 WARN(UNKNOWN_TOPIC_OR_PARTITION) 이 하나 흘러갈 뿐이고,
         *  그건 인프라 로그라 아무도 알림을 걸어 두지 않습니다.
         *  콜백이 없으면 LoggingProducerListener 의 ERROR 한 줄이 전부이며,
         *  그것도 payload 가 잘려 있어 어떤 주문이었는지 복구할 단서가 부족합니다.
         * ==================================================================== */
        void answer3() throws Exception {
            log.info("===== 정답 3 =====");

            CompletableFuture<SendResult<String, OrderCreated>> future =
                    template.send("no-such-topic-s02", "K-1", OrderCreated.of(1));

            log.info("send() 직후 isDone={}", future.isDone());

            for (int t = 1; t <= 70; t++) {
                Thread.sleep(1000);
                boolean done = future.isDone();
                boolean failed = future.isCompletedExceptionally();
                if (done || t % 10 == 0) {
                    log.info("t={}s isDone={} exceptionally={}", t, done, failed);
                }
                if (done) break;
            }

            if (future.isCompletedExceptionally()) {
                try {
                    future.join();
                } catch (Exception e) {
                    Throwable cause = e.getCause() != null ? e.getCause() : e;
                    log.error("담긴 예외 = {}: {}", cause.getClass().getSimpleName(), cause.getMessage());
                }
            } else {
                log.warn("70초 안에도 안 끝났습니다. isDone={}", future.isDone());
            }
        }

        /* ====================================================================
         * 정답 5. 키 → 파티션 매핑 예측과 검증
         *
         *  ★ 계산식은 반드시 Utils.toPositive(Utils.murmur2(bytes)) % n 입니다.
         *
         *  Math.abs(hash) % n 을 쓰면 안 되는 이유:
         *      Math.abs(Integer.MIN_VALUE) == Integer.MIN_VALUE 입니다. (음수 그대로)
         *      2의 보수에서 -2147483648 의 절댓값 2147483648 은 int 범위를 넘기 때문에
         *      오버플로로 자기 자신이 됩니다. 그러면 % n 결과가 음수가 되고,
         *      파티션 번호로 쓰는 순간 ArrayIndexOutOfBoundsException 이 납니다.
         *      확률은 40억분의 1이지만, 하루 수억 건을 보내면 언젠가 터집니다.
         *      Kafka 가 Math.abs 대신 toPositive(x & 0x7fffffff)를 쓰는 이유가 정확히 이것입니다.
         *      비트 마스크는 부호 비트만 떼므로 절대 음수가 나오지 않습니다.
         *
         *  실측 결과 (파티션 3개 / 4개):
         *      ORD-0010  murmur2=  208042438  %3=1  %4=2
         *      ORD-0011  murmur2= 1690680040  %3=1  %4=0
         *      ORD-0012  murmur2=  948856158  %3=0  %4=2
         *      ORD-0013  ...
         *
         *  ★ 3개일 때와 4개일 때가 거의 전부 다릅니다. 이게 이 문제의 결론입니다.
         *  파티션을 3 → 4 로 늘리면 같은 주문의 과거 이벤트는 옛 파티션에,
         *  새 이벤트는 새 파티션에 쌓입니다. 두 파티션은 독립적으로 소비되므로
         *  나중에 발행된 이벤트가 먼저 처리될 수 있습니다. 순서 보장이 소급해서 깨집니다.
         *  에러는 나지 않고, 재고가 음수가 되거나 취소된 주문이 배송되는 것으로 드러납니다.
         *
         *  파티션 수는 하드코딩하지 말고 partitionsFor 로 조회하세요.
         *  다만 partitionsFor 는 메타데이터가 없으면 max.block.ms 까지 블로킹하므로
         *  요청 경로에서 매번 부르면 안 됩니다. 기동 시 한 번 조회해 캐시하세요.
         * ==================================================================== */
        void answer5() throws Exception {
            log.info("===== 정답 5 =====");

            int numPartitions = template.partitionsFor(TOPIC).size();
            log.info("{} 의 파티션 수 = {}", TOPIC, numPartitions);

            for (int i = 10; i <= 15; i++) {
                String key = "ORD-%04d".formatted(i);
                int predicted3 = predictPartition(key, 3);
                int predicted4 = predictPartition(key, 4);
                int predicted = predictPartition(key, numPartitions);

                template.send(TOPIC, key, OrderCreated.of(i)).whenComplete((r, ex) -> {
                    if (ex != null) {
                        log.error("발행 실패 key={}", key, ex);
                        return;
                    }
                    int actual = r.getRecordMetadata().partition();
                    log.info("{} 예측(3개)={} 예측(4개)={} | 실제={}  {}",
                            key, predicted3, predicted4, actual,
                            predicted == actual ? "← 일치" : "← ★불일치! 파티션 수 확인");
                });
            }
            template.flush();
            Thread.sleep(500);
        }

        /** Kafka 기본 파티셔너와 동일한 계산. Math.abs 를 쓰지 말 것. */
        static int predictPartition(String key, int numPartitions) {
            byte[] bytes = key.getBytes(StandardCharsets.UTF_8);
            return Utils.toPositive(Utils.murmur2(bytes)) % numPartitions;
        }

        private static String key(int i) {
            return "ORD-%04d".formatted(i);
        }

        private static long ms(long startNanos) {
            return (System.nanoTime() - startNanos) / 1_000_000;
        }

        private static void report(String label, int n, long elapsedMs) {
            long tps = elapsedMs == 0 ? 0 : (n * 1000L) / elapsedMs;
            log.info("{}: {}건 / {} ms = {} msg/s", label, n, elapsedMs, String.format("%,d", tps));
        }
    }

    /* ========================================================================
     * 정답 4. 설정이 다른 KafkaTemplate 두 개
     *
     *  ★ @Qualifier 로 빈 이름을 명시하는 방식을 택합니다. @Primary 는 쓰지 않습니다.
     *
     *  왜 기동이 실패하는가:
     *      두 @Bean 을 등록하면 KafkaTemplate<String, OrderCreated> 타입 빈이
     *      자동설정의 kafkaTemplate 까지 합쳐 3개가 됩니다. 이름 없이 주입하면
     *          NoUniqueBeanDefinitionException: expected single matching bean but found 3
     *      로 기동이 실패합니다. 이건 "시끄러운 실패"라 오히려 다행입니다.
     *
     *  왜 @Primary 를 권장하지 않는가:
     *      @Primary 는 "이름 없이 주입한 모든 코드"가 어느 템플릿을 쓸지를 침묵으로 결정합니다.
     *      bulk(acks=1, idempotence=false) 쪽에 @Primary 를 붙여 두면,
     *      나중에 누군가 아무 생각 없이 KafkaTemplate 을 주입받아 주문 이벤트를 발행할 때
     *      그 이벤트는 acks=1 로, 멱등성 없이 나갑니다. 컴파일도 되고 기동도 되고
     *      테스트도 통과하며, 브로커가 한 번 흔들릴 때까지 아무도 모릅니다.
     *      이 코스가 잡으려는 "조용한 실패"의 전형입니다.
     *      @Qualifier 를 강제하면, 새 코드는 반드시 "어느 쪽인가"를 선택해야만 컴파일됩니다.
     *
     *  clientId 를 반드시 다르게 주는 이유:
     *      안 주면 producer-1, producer-2 로 자동 부여되어 로그와 JMX 지표에서
     *      어느 쪽이 어느 설정인지 구분할 수 없습니다. 장애 때 이게 결정적입니다.
     *
     *  기동 로그 확인 포인트:
     *      [Producer clientId=audit] Instantiated an idempotent producer.   ← 있음
     *      [Producer clientId=bulk]  ...                                     ← 없음
     *      이 한 줄의 유무가 멱등성 활성화 여부를 알려 주는 가장 확실한 신호입니다.
     * ======================================================================== */
    @Configuration
    @Profile("step02-sol4")
    public static class Answer4Config {

        private Map<String, Object> base(String bootstrap) {
            Map<String, Object> props = new HashMap<>();
            props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap);
            props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
            return props;
        }

        @Bean
        public KafkaTemplate<String, OrderCreated> auditKafkaTemplate(
                @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
            Map<String, Object> props = base(bootstrap);
            props.put(ProducerConfig.ACKS_CONFIG, "all");
            props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
            props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
            props.put(ProducerConfig.LINGER_MS_CONFIG, 0);       // 지연 최소화
            props.put(ProducerConfig.CLIENT_ID_CONFIG, "audit");
            return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props));
        }

        @Bean
        public KafkaTemplate<String, OrderCreated> bulkKafkaTemplate(
                @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
            Map<String, Object> props = base(bootstrap);
            props.put(ProducerConfig.ACKS_CONFIG, "1");
            props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, false);   // acks=1 이면 명시적 false
            props.put(ProducerConfig.LINGER_MS_CONFIG, 100);              // 많이 모아서
            props.put(ProducerConfig.BATCH_SIZE_CONFIG, 131_072);
            props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
            props.put(ProducerConfig.CLIENT_ID_CONFIG, "bulk");
            return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(props));
        }
    }

    @Component
    @Profile("step02-sol4")
    public static class Answer4Demo implements ApplicationRunner {

        private static final Logger log = LoggerFactory.getLogger(Answer4Demo.class);

        private final KafkaTemplate<String, OrderCreated> audit;
        private final KafkaTemplate<String, OrderCreated> bulk;

        // ★ @Qualifier 없이 주입하면 NoUniqueBeanDefinitionException 으로 기동 실패합니다.
        public Answer4Demo(@Qualifier("auditKafkaTemplate") KafkaTemplate<String, OrderCreated> audit,
                           @Qualifier("bulkKafkaTemplate") KafkaTemplate<String, OrderCreated> bulk) {
            this.audit = audit;
            this.bulk = bulk;
        }

        @Override
        public void run(ApplicationArguments args) {
            log.info("===== 정답 4 =====");

            audit.send(TOPIC, "ORD-0001", OrderCreated.of(1))
                    .whenComplete((r, ex) -> logResult("audit", r, ex));
            bulk.send(TOPIC, "ORD-0002", OrderCreated.of(2))
                    .whenComplete((r, ex) -> logResult("bulk", r, ex));

            audit.flush();
            bulk.flush();
        }

        private void logResult(String which, SendResult<String, OrderCreated> r, Throwable ex) {
            if (ex != null) {
                log.error("[{}] 발행 실패", which, ex);
                return;
            }
            RecordMetadata md = r.getRecordMetadata();
            log.info("[{}] → {}-{}@{}", which, md.topic(), md.partition(), md.offset());
        }
    }

    /* ========================================================================
     * 정답 6. acks=1 + enable.idempotence=true
     *
     *  (A) 둘 다 "명시"한 경우 → 기동 실패
     *
     *      org.apache.kafka.common.config.ConfigException:
     *        Must set acks to all in order to use the idempotent producer.
     *        Otherwise we cannot guarantee idempotence.
     *          at ProducerConfig.postProcessAndValidateIdempotenceConfigs(ProducerConfig.java:594)
     *          at KafkaProducer.<init>(KafkaProducer.java:334)
     *
     *      멱등성은 acks=all 을 요구합니다. 리더만 받고 성공을 반환하면
     *      리더가 죽었을 때 그 레코드가 사라지고, 재시도가 "중복 제거 대상"이 아니라
     *      "새 발행"이 되어 멱등성의 전제 자체가 무너지기 때문입니다.
     *
     *      같은 종류의 제약이 둘 더 있습니다.
     *        enable.idempotence=true + retries=0
     *          → Must set retries to non-zero when using the idempotent producer.
     *        enable.idempotence=true + max.in.flight > 5
     *          → Must set max.in.flight.requests.per.connection to at most 5 ...
     *
     *  (B) enable.idempotence 를 "설정하지 않고" acks=1 만 준 경우 → 기동 성공
     *
     *      ProducerConfig 는 enable.idempotence 가 사용자에 의해 명시되지 않았고
     *      다른 설정과 충돌하면, 예외 대신 조용히 false 로 내립니다.
     *      "기본값이니까 사용자가 의도한 게 아닐 수도 있다"는 판단입니다.
     *      경고 로그도 없습니다. 유일한 단서는 기동 로그에서
     *          [Producer clientId=...] Instantiated an idempotent producer.
     *      가 "사라졌다"는 것뿐인데, 없는 로그를 알아채는 사람은 거의 없습니다.
     *
     *  ★ (A) 와 (B) 중 운영에서 더 위험한 쪽은 단연 (B) 입니다.
     *
     *      (A) 는 애플리케이션이 아예 안 뜹니다. 배포 파이프라인이 즉시 막히고,
     *          예외 메시지가 원인과 해결책을 그대로 알려 줍니다. 최악의 경우가 배포 실패입니다.
     *
     *      (B) 는 아무 일도 없이 뜹니다. 평소에는 아무 문제도 없습니다.
     *          그러다 네트워크가 한 번 흔들려 retries 가 동작하는 순간,
     *          - 멱등성이 꺼져 있으므로 중복 레코드가 생기고,
     *          - max.in.flight=5 와 만나 같은 파티션 안에서 순서가 뒤집힙니다.
     *          컨슈머는 이걸 알 방법이 없습니다. 재고가 음수가 되거나
     *          취소된 주문이 배송되는 것으로 며칠 뒤에 드러납니다.
     *          그때쯤이면 "acks: 1 한 줄"과 증상을 연결할 사람이 없습니다.
     *
     *  ★ 결론: 멱등성이 필요하면 enable.idempotence: true 를 "명시"하세요.
     *     그러면 충돌 시 조용히 꺼지는 대신 기동이 실패합니다.
     *     시끄러운 실패를 사는 것이 조용한 유실을 사는 것보다 언제나 낫습니다.
     *
     *  브로커 쪽 ack/ISR 동작은 Kafka 코스를 참고하세요:
     *      ../../kafka/step-04-producer/
     *      ../../kafka/step-07-delivery-semantics/
     *      ../../kafka/step-08-replication/
     * ======================================================================== */
    @Configuration
    @Profile("step02-sol6")
    public static class Answer6Config {

        /** (A) 기동 실패를 재현한다. 이 프로필로 켜면 앱이 안 뜨는 것이 정상이다. */
        @Bean
        public KafkaTemplate<String, OrderCreated> brokenKafkaTemplate(
                @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
            Map<String, Object> props = new HashMap<>();
            props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap);
            props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
            props.put(ProducerConfig.CLIENT_ID_CONFIG, "broken");

            props.put(ProducerConfig.ACKS_CONFIG, "1");                   // ← 충돌
            props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);    // ← 명시적으로 켬

            DefaultKafkaProducerFactory<String, OrderCreated> pf =
                    new DefaultKafkaProducerFactory<>(props);
            // ★ 팩토리는 지연 생성이라 빈 등록만으로는 예외가 안 납니다.
            //   실제 KafkaProducer 를 만들게 강제해야 ConfigException 이 여기서 터집니다.
            pf.createProducer().close();
            return new KafkaTemplate<>(pf);
        }
    }

    @Configuration
    @Profile("step02-sol6-silent")
    public static class Answer6SilentConfig {

        /** (B) 기동은 성공한다. 그러나 멱등성이 조용히 꺼진다. */
        @Bean
        public KafkaTemplate<String, OrderCreated> silentlyNonIdempotentTemplate(
                @Value("${spring.kafka.bootstrap-servers}") String bootstrap) {
            Map<String, Object> props = new HashMap<>();
            props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap);
            props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
            props.put(ProducerConfig.CLIENT_ID_CONFIG, "silent");

            props.put(ProducerConfig.ACKS_CONFIG, "1");
            // enable.idempotence 를 "쓰지 않는다" → 예외 없이 false 로 내려간다

            DefaultKafkaProducerFactory<String, OrderCreated> pf =
                    new DefaultKafkaProducerFactory<>(props);
            pf.createProducer().close();     // 예외 없이 통과한다. 그게 문제다.
            return new KafkaTemplate<>(pf);
        }
    }
}