WebFlux로 고성능 애플리케이션을 개발해야 하는 상황이라 기술검토가 필요했습니다. WebFlux BlockHound 로 블로킹 코드를 찾아내는 게 목표였습니다.
2023년 말에 사내 위키에 정리해뒀던 내용인데, 지금 봐도 쓸모가 있을 것 같아 옮겨 둡니다. 중간에 원인을 못 찾고 넘어간 부분도 그대로 남겨뒀습니다.
WebFlux BlockHound
WebFlux는 비동기 넌블로킹 방식의 리액티브 프로그래밍을 지원하는데, 리액티브 프로그래밍은 하나의 블록킹 코드만 있어도 리액티브 플로우가 깨지고 성능이 급격히 나빠집니다.
때문에 블록킹 코드를 검출하는 도구인 BlockHound를 적용하였습니다.
블록하운드는 런타임 시 블로킹 코드가 실행되면 BlockingOperationError를 발생시킵니다. 컴파일 시 찾아주는 게 아니기 때문에 테스트는 필수입니다.
스레드 관점에서 보면 이렇습니다. 톰캣은 요청마다 스레드를 하나씩 주기 때문에 한 요청이 멈춰도 다른 요청은 다른 스레드에서 굴러가는데, WebFlux는 적은 수의 이벤트 루프가 여러 요청을 번갈아 처리하기 때문에 그 스레드가 멈추면 얹혀 있던 요청이 전부 같이 멈춥니다.

WebFlux BlockHound 사용법
의존성 추가.
implementation 'io.projectreactor.tools:blockhound:1.0.8.RELEASE'
JDK 13버전 이상의 경우, 의존성만 추가하면 아래와 같은 에러가 출력됩니다.
> Task :order-producer:OrderProducerApplication.main()
Exception in thread "main" java.lang.IllegalStateException: The instrumentation have failed.
It looks like you're running on JDK 13+.
You need to add '-XX:+AllowRedefinitionToAddDeleteMethods' JVM flag.
See https://github.com/reactor/BlockHound/issues/33 for more info.
공식 페이지에는 임시 해결방안으로 JVM -XX:+AllowRedefinitionToAddDeleteMethods 옵션을 추가해서 사용하라고 안내하고 있습니다.
아래와 같이 build.gradle 파일에 코드를 추가해서도 해결할 수 있다고 안내되어 있지만,, 동작하지 않습니다!
tasks.withType(Test).all {
if (JavaVersion.current().isCompatibleWith(JavaVersion.VERSION_13)) {
jvmArgs += [
"-XX:+AllowRedefinitionToAddDeleteMethods"
]
}
}
저는 인텔리제이 JVM 옵션을 추가하는 것으로 해결했습니다.
@SpringBootApplication
public class OrderProducerApplication {
public static void main(String[] args) {
BlockHound.install(); //추가
SpringApplication.run(OrderProducerApplication.class, args);
}
}
BlockHound.install()을 추가하면 애플리케이션이 시작되는 시점부터 모든 바이트 코드를 분석해서 블로킹 코드가 존재하는지 검사합니다.
테스트
@RestController
@RequiredArgsConstructor
public class OrderController {
private final OrderService orderService;
@GetMapping("/nonblock")
public Mono<ResponseEntity<Result>> nonblock(){
return Mono.just(ResponseEntity.ok(Result.success(orderService.nonblock())));
}
@GetMapping("/block")
public Mono<ResponseEntity<Result>> block(){
return Mono.just(ResponseEntity.ok(Result.success(orderService.block())));
}
}
@Service
@RequiredArgsConstructor
public class OrderService {
public String nonblock() {
System.out.println("nonblock 호출");
return "nonblock!";
}
public String block() {
System.out.println("block 호출");
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
return "block";
}
}
/nonblock 호출 시, 에러 없이 메서드 수행에 성공합니다.
/block 호출 시, BlockingOperationError를 발생시킵니다.
이렇게 블록하운드를 사용하면 리액티브 애플리케이션 개발 시 많은 도움을 받을 수 있습니다.
BlockHound와 Binder-Kafka-Reactive
위에서 간단히 살펴본 블록하운드를 이용하여 주문 서비스를 개발하였습니다.
주문 서비스는 Spring WebFlux와 Kafka를 사용하며 2가지로 분리되어 있습니다.
- 주문이 들어올 때 Request → Domain 변경하여 Kafka에 메시지를 보내는 Producer 서버
- Kafka로 메시지를 받아 Domain → Entity 변경하여 Repository에 저장하는 Consumer 서버
Producer 서버에 BlockHound를 적용하여 테스트해보았습니다.
//spring
implementation 'org.springframework.boot:spring-boot-starter-webflux'
//blockhound
implementation 'io.projectreactor.tools:blockhound:1.0.8.RELEASE'
//kafka
implementation 'org.springframework.cloud:spring-cloud-stream-binder-kafka-reactive'
@Slf4j
@Service
@RequiredArgsConstructor
public class KafkaProducer {
private final KafkaTemplate<String, Order> kafkaTemplate;
@Value("${spring.kafka.topic}")
private String TOPIC;
public void sendMessage(Order order) {
UUID key = UUID.randomUUID();
kafkaTemplate.send(TOPIC, key.toString(), order);
log.info("■■■ send to kafka : {}", order);
}
}
@Slf4j
@Service
@RequiredArgsConstructor
public class KafkaConsumer {
private final ObjectMapper objectMapper;
@KafkaListener(topics = "order", groupId = "order-group")
public void listen(String message) {
Order order = null;
try {
order = objectMapper.readValue(message, Order.class);
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
log.info("■■■ receive from kafka message : {}, Order : {}", message, order);
}
}
kafkaTemplate.send() 1번째 호출
블록하운드가 Producer 쪽에 블록 코드가 존재한다고 알려줍니다.
2023-12-21 14:41:31.269 [INFO ] AppInfoParser - Kafka version: 3.4.1
2023-12-21 14:41:31.281 [INFO ] Metadata - [Producer clientId=producer-1] Resetting the last see...
2023-12-21 14:41:31.309 [ERROR] AbstractErrorWebExceptionHandler - [3b6bebf2-1]
reactor.blockhound.BlockingOperationError: Blocking call! java.lang.Object#wait
500 Server
kafkaTemplate.send() 2번째 호출부터
똑같이 카프카에서 에러가 발생할 줄 알았지만 이번엔 잘 동작합니다. 이후 몇 번을 시도해봐도 블록하운드는 에러를 발생시키지 않습니다.
2023-12-21 14:45:47.373 [INFO ] KafkaProducer - ■■■ send to kafka : Order(...)
2023-12-21 14:45:47.439 [INFO ] KafkaConsumer - ■■■ receive from kafka message : {"status":n...
결론부터 말하자면 명확한 원인과 해결책은 조금 더 찾아봐야 할 것 같습니다.
일단 동일한 증상에 대해서 이슈로 등록된 것은 확인하였고, 카프카 내부 동작도 조금 더 자세히 살펴보아야 할 것 같습니다.
블록하운드가 동작하지 않는 건 아닌가?
블록하운드가 최초 1회만 동작한 건 아닌가? 라는 의문이 들어 확인해보았습니다.
public void sendMessage(Order order) {
UUID key = UUID.randomUUID();
try {
Thread.sleep(1000);
log.info("BlockHound 동작해요");
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
kafkaTemplate.send(TOPIC, key.toString(), order);
log.info("■■■ send to kafka : {}", order);
}
위에서 살펴봤던 카프카 Producer 내부에서 발생했던 에러메시지와는 다르게 Thread.sleep에서 에러가 발생합니다. 다시 호출해봐도 똑같이 Thread.sleep 에러메시지를 보여줍니다.
2023-12-21 14:57:41.320 [ERROR] AbstractErrorWebExceptionHandler - [25dbadeb-1]
reactor.blockhound.BlockingOperationError: Blocking call! java.lang.Thread.sleep
500 Server
이번엔 슬립을 kafkaTemplate.send 이후에 위치시켰습니다.
public void sendMessage(Order order) {
UUID key = UUID.randomUUID();
kafkaTemplate.send(TOPIC, key.toString(), order);
try {
Thread.sleep(1000);
log.info("BlockHound 동작해?");
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
log.info("■■■ send to kafka : {}", order);
}
이번엔 1번째는 카프카 Producer 내부에서 에러를 발생했고, 2번째 이후부터는 Thread.sleep에서 에러가 발생합니다.
2023-12-21 15:00:20.931 [ERROR] AbstractErrorWebExceptionHandler - [915c7305-1]
reactor.blockhound.BlockingOperationError: Blocking call! java.lang.Object#wait
2023-12-21 15:00:27.106 [ERROR] AbstractErrorWebExceptionHandler - [915c7305-2]
reactor.blockhound.BlockingOperationError: Blocking call! java.lang.Thread.sleep
테스트 결과, 카프카 메시지 발송 최초에만 블록처리가 이루어지고, 2번째부터는 블록처리하지 않습니다.
2023년 12월 21일 추가
원인
org.apache.kafka.clients.producer.internals.ProducerMetadata.awaitUpdate(ProducerMetadata.java:119)
~[kafka-clients-3.4.1.jar:?]
/**
* Wait for metadata update until the current version is larger than the last version we know of
*/
public synchronized void awaitUpdate(final int lastVersion, final long timeoutMs)
throws InterruptedException {
long currentTimeMs = time.milliseconds();
long deadlineMs = currentTimeMs + timeoutMs < 0 ? Long.MAX_VALUE : currentTimeMs + timeoutMs;
time.waitObject(this, () -> {
// Throw fatal exceptions, if there are any. Recoverable topic errors will be handled by the caller.
maybeThrowFatalException();
return updateVersion() > lastVersion || isClosed();
}, deadlineMs);
if (isClosed())
throw new KafkaException("Requested metadata update after close");
}
ProducerMetadata.awaitUpdate() 메서드가 synchronized로 작성된 것 확인.
프로듀서는 메시지를 보내기 전에 토픽의 파티션과 리더를 알아야 하는데, 첫 전송 때는 그 메타데이터가 없어서 브로커 응답을 기다립니다. 그때 스레드가 실제로 멈춥니다. 두 번째부터는 캐시에 있으니 기다리지 않습니다.

2023년 12월 22일
BlockHound와 Reactor-Kafka
이번엔 리액터 카프카를 이용하여 테스트해보겠습니다.
//spring
implementation 'org.springframework.boot:spring-boot-starter-webflux'
//blockhound
implementation 'io.projectreactor.tools:blockhound:1.0.8.RELEASE'
//kafka
implementation 'io.projectreactor.kafka:reactor-kafka:1.3.19'
@Configuration
public class ReactiveKafkaConfig {
@Value("${kafka.topic}") String topic;
@Value("${kafka.bootstrap-servers}") String bootstrapServers;
@Bean
public ReactiveKafkaProducerTemplate<String, Order> reactiveKafkaProducerTemplate(
KafkaProperties properties) {
Map<String, Object> props = properties.buildProducerProperties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
return new ReactiveKafkaProducerTemplate<>(SenderOptions.create(props));
}
}
@Slf4j
@Service
@RequiredArgsConstructor
public class KafkaProducer {
private final ReactiveKafkaProducerTemplate<String, Order> reactiveKafkaProducer;
@Value("${kafka.topic}")
private String TOPIC;
public void sendMessage(Order order) {
String key = UUID.randomUUID().toString();
reactiveKafkaProducer.send(TOPIC, key, order)
.doOnSuccess(result -> log.info("■■■ send to kafka key: {}, message: {}", key, order))
.subscribe();
}
}
호출 테스트 결과, 1번째 호출부터 블록하운드가 아무 에러도 발생시키지 않습니다.
WebFlux BlockHound 를 쓸지에 대한 결론
org.springframework.cloud:spring-cloud-stream-binder-kafka-reactive
→ 최초 1회 메타데이터 가져올 때만 블록처리됩니다.
io.projectreactor.kafka:reactor-kafka:1.3.19
→ 논블록으로 처리됩니다! 이거 쓰세요!
참고
- reactor/BlockHound (GitHub)
- OpenJDK 13 requires -XX:+AllowRedefinitionToAddDeleteMethods flag — BlockHound Issue #33
- Questionable blocking behavior in KafkaSender — reactor-kafka Issue #71
- ProducerMetadata.java — apache/kafka