반응형
Java 에서 비동기 처리를 하는 방법 (Future, Complated Future)
실습(Future)
class JavaFutureStudy {
private fun helloFuture(): Future<String> {
val executor = Executors.newSingleThreadExecutor()
return try {
executor.submit<String> {
println("[${Thread.currentThread().name} executor.submit() 실행됨")
"hello future!"
}
} finally {
executor.shutdown()
}
}
@Test
@DisplayName("Future 테스트 - future.get() 함수는 Blocking 되고 isDone, isCancelled 상태를 갖는다.")
fun getFutureTest() {
val future: Future<String> = helloFuture()
Assertions.assertFalse(future.isDone)
Assertions.assertFalse(future.isCancelled)
// future.get()은 블록킹 메소드
val futureResult = future.get()
Assertions.assertEquals("hello future!", futureResult)
Assertions.assertTrue(future.isDone)
Assertions.assertFalse(future.isCancelled)
}
private fun timeoutFuture(): Future<String> {
val executor = Executors.newSingleThreadExecutor()
return try {
executor.submit<String> {
println("[${Thread.currentThread().name} executor.submit() 실행됨")
Thread.sleep(1000)
"Hello, World!"
}
} finally {
executor.shutdown()
}
}
@Test
@DisplayName("Future 테스트 - Timeout 설정 > timeout 시간동안 Thread 는 Blocking 된다")
fun futureTimeoutTest() {
val future: Future<String> = timeoutFuture()
val futureResult = future.get(1500, java.util.concurrent.TimeUnit.MILLISECONDS)
Assertions.assertEquals("Hello, World!", futureResult)
Assertions.assertTrue(future.isDone)
Assertions.assertFalse(future.isCancelled)
}
@Test
@DisplayName("Future 테스트 - Timeout 설정 Exception > Timeout 시간보다 더 걸리면 TimeoutException 발생한다")
fun futureTimeoutExceptionTest() {
val future: Future<String> = timeoutFuture()
Assertions.assertThrows(java.util.concurrent.TimeoutException::class.java) {
future.get(500, java.util.concurrent.TimeUnit.MILLISECONDS)
}
}
@Test
@DisplayName("Future 테스트 - cancel > Future 를 취소할 수 있다, 취소된 future.get()은 CancellationException 발생한다")
fun futureCancelTest() {
val future: Future<String> = timeoutFuture()
val cancel = future.cancel(true)
Assertions.assertTrue(cancel)
Assertions.assertTrue(future.isDone)
Assertions.assertTrue(future.isCancelled)
Assertions.assertThrows(java.util.concurrent.CancellationException::class.java) {
future.get()
}
val cancelRepeat = future.cancel(true)
Assertions.assertFalse(cancelRepeat)
}
private fun exceptionFuture(): Future<String> {
val executor = Executors.newSingleThreadExecutor()
return try {
executor.submit<String> {
throw IllegalArgumentException("Error")
}
} finally {
executor.shutdown()
}
}
@Test
@DisplayName("Future 의 한계 > Exception 이 발생하든 cancel 이 되든 isDone 은 true 이다 > 완료되거나 에러가 발생했는지 구분이 어렵다 > 오류 헨들링이 어렵다")
fun future(){
val future = helloFuture()
future.cancel(true)
Assertions.assertTrue(future.isDone)
val exceptionFuture = exceptionFuture()
Assertions.assertThrows(java.util.concurrent.ExecutionException::class.java) {
exceptionFuture.get()
}
Assertions.assertTrue(exceptionFuture.isDone)
}
}
Future 의 한계
- future.get() 을 호출하면 메인 Thread 가 Blocking 된다.
- Exception 오류 핸들링이 어렵다.
- future 가 cancel 되도 isDone 은 true 이다.
실습(Complated Future)
package org.example.completable.future
import org.example.printlnWithThreadName
import org.junit.jupiter.api.DisplayName
import org.junit.jupiter.api.Test
import java.util.concurrent.CompletableFuture
import java.util.concurrent.CompletionStage
/**
* CompletionStage Interface
* @see CompletionStage
*/
class CompletionFutureStudy{
private fun finishedStage(): CompletionStage<String> {
val future = CompletableFuture.supplyAsync {
printlnWithThreadName("return helloFinishedStage")
"Hello, CompletableFuture!"
}
Thread.sleep(100)
return future
}
@Test
@DisplayName(
"""
thenAccept 는 FUTURE STAGE 가 DONE 상태이면 메인 THREAD 에서 실행 됨. finishedStage 블락되는 작업이 함께 있으면 메인 스레드가 블락될 수 있음.
비동기 작업이 빨리 끝났는데 Block 되는 작업이 있다면 main thread 가 블락킹 될 수 있음 > Blocking 비동기 상태가 된다
"""
)
fun thenAcceptFinishedStageTest(){
printlnWithThreadName("start main")
finishedStage()
.thenAccept{ printlnWithThreadName("thenAccept >> $it") }
.thenAccept{ printlnWithThreadName("thenAccept2 >> $it") }
printlnWithThreadName("after thenAccept")
Thread.sleep(100)
}
@Test
@DisplayName("thenAcceptAsync 는 FUTURE STAGE 가 DONE 상태여도 별도의 thread pool 에서 실행됨")
fun thenAsyncAcceptFinishedStageTest(){
println("[${Thread.currentThread().name}] start main")
finishedStage()
.thenAcceptAsync{ printlnWithThreadName("thenAcceptAsync >> $it") }
.thenAcceptAsync{ printlnWithThreadName("thenAcceptAsync2 >> $it") }
printlnWithThreadName("after thenAccept")
Thread.sleep(100)
}
private fun runningStage(): CompletionStage<String> {
return CompletableFuture.supplyAsync {
Thread.sleep(1000)
printlnWithThreadName("return helloRunningStage")
"Hello, CompletableFuture!"
}
}
@Test
@DisplayName("thenAccept 는 FUTURE STAGE 가 RUNNING 상태이면 별도의 Thread 에서 실행 됨")
fun thenAcceptRunningStageTest(){
printlnWithThreadName("start main")
runningStage()
.thenAccept{ printlnWithThreadName("thenAccept >> $it") }
.thenAccept{ printlnWithThreadName("thenAccept2 >> $it") }
printlnWithThreadName("after thenAccept")
Thread.sleep(2000)
}
@Test
@DisplayName("thenAcceptAsync 는 FUTURE STAGE 가 RUNNING 상태이면 별도의 Thread 에서 실행 됨")
fun thenAcceptAsyncRunningStageTest(){
printlnWithThreadName("start main")
runningStage()
.thenAcceptAsync{ printlnWithThreadName("thenAcceptAsync >> $it") }
.thenAcceptAsync{ printlnWithThreadName("thenAcceptAsync2 >> $it") }
printlnWithThreadName("after thenAccept")
Thread.sleep(2000)
}
@Test
@DisplayName("completableFuture 는 완료 상태로 변경할 수 있고 이미 완료 된 경우 FALSE 를 반환 한다.")
fun completeTest() {
val future = CompletableFuture<String>()
assert(!future.isDone)
var triggered = future.complete("completed")
assert(future.isDone)
assert(triggered)
assert(future.get() == "completed")
triggered = future.complete("completed2")
assert(future.isDone)
assert(!triggered)
assert(future.get() == "completed")
}
@Test
@DisplayName("completableFuture 는 isCompletedExceptionally 로 오류가 났는지 확인할 수 있다.")
fun isCompletedExceptionallyTest(){
val futureWithException = CompletableFuture.supplyAsync{
return@supplyAsync 1/0
}
Thread.sleep(100)
assert(futureWithException.isDone)
assert(futureWithException.isCompletedExceptionally)
}
private fun waitAndReturn(millis: Int, value: Int): CompletableFuture<Int> {
return CompletableFuture.supplyAsync<Int> {
try {
printlnWithThreadName("waitAndReturn: {$millis}ms")
Thread.sleep(millis.toLong())
} catch (e: InterruptedException) {
throw RuntimeException(e)
}
value
}
}
@Test
@DisplayName("CompletableFuture.allOf() 를 사용하면 각 Future 가 완료될 때까지 기다린다.")
fun allOfTest(){
printlnWithThreadName("start")
val firstFuture = waitAndReturn(100, 1)
val secondFuture = waitAndReturn(500, 2)
val thirdFuture = waitAndReturn(3000, 3)
CompletableFuture.allOf(firstFuture, secondFuture, thirdFuture).thenAcceptAsync{
printlnWithThreadName("after allOf")
printlnWithThreadName(firstFuture.get())
printlnWithThreadName(secondFuture.get())
printlnWithThreadName(thirdFuture.get())
}.join()
println("end")
}
@Test
@DisplayName("anyOf 의 경우 가장 빨리 끝난 Future 의 결과를 가져온다.")
fun anyOfTest(){
printlnWithThreadName("start")
val firstFuture = waitAndReturn(500, 1)
val secondFuture = waitAndReturn(100, 2)
val thirdFuture = waitAndReturn(3000, 3)
CompletableFuture.anyOf(firstFuture, secondFuture, thirdFuture).thenAcceptAsync{
printlnWithThreadName("after anyOf")
printlnWithThreadName("first value $it")
}.join()
println("end")
}
}
Complated Future 정리
- thenAccept 는 Future Stage 가 Done 일 경우 메인 Thread 에서 실행되어서 메인 Thread 가 Blocking 될 수 있음.
- thenAcceptAsync 를 사용해야 함.
- isCompletedExceptionally 로 오류가 났는지 확인할 수 있음
- 여러 Future 를 동시에 실행하고 allOf 를 사용하여 완료 되기를 기다렸다 체이닝 해서 처리할 수 있음
- 지연 로딩 기능을 제공 하지 않음 (hot, cold?)
- 지속 적으로 생성 되는 데이터를 처리하기 어려움
반응형
'개발 > Reactive Streams' 카테고리의 다른 글
| Reactor 의 Mono, Flux 는 Reactive Streams 의 구현체이다 (1) | 2025.05.14 |
|---|---|
| [Reactive Streams] Netflix 에서 RxJava 를 도입한 이유 (1) | 2025.04.12 |
