diff --git a/ReactiveStreams/.gitignore b/ReactiveStreams/.gitignore new file mode 100644 index 0000000..3884e13 --- /dev/null +++ b/ReactiveStreams/.gitignore @@ -0,0 +1,38 @@ +HELP.md +.gradle +build/ +!gradle/wrapper/gradle-wrapper.jar +!**/src/main/**/build/ +!**/src/test/**/build/ + +### STS ### +.apt_generated +.classpath +.factorypath +.project +.settings +.springBeans +.sts4-cache +bin/ +!**/src/main/**/bin/ +!**/src/test/**/bin/ + +### IntelliJ IDEA ### +.idea +*.iws +*.iml +*.ipr +out/ +!**/src/main/**/out/ +!**/src/test/**/out/ + +### NetBeans ### +/nbproject/private/ +/nbbuild/ +/dist/ +/nbdist/ +/.nb-gradle/ + +### VS Code ### +.vscode/ +/gradle/ diff --git a/ReactiveStreams/build.gradle b/ReactiveStreams/build.gradle new file mode 100644 index 0000000..b1c3598 --- /dev/null +++ b/ReactiveStreams/build.gradle @@ -0,0 +1,29 @@ +plugins { + id 'org.springframework.boot' version '2.5.5' + id 'io.spring.dependency-management' version '1.0.11.RELEASE' + id 'java' +} + +group = 'com.example' +version = '0.0.1-SNAPSHOT' +sourceCompatibility = '1.8' + +repositories { + mavenCentral() +} + +dependencies { + implementation 'org.springframework.boot:spring-boot-starter' + implementation 'org.springframework.boot:spring-boot-starter-web' + testImplementation 'org.springframework.boot:spring-boot-starter-test' + testImplementation 'org.springframework.boot:spring-boot-starter-web' + implementation 'org.springframework.boot:spring-boot-starter-webflux' + compileOnly 'org.projectlombok:lombok:1.18.20' + annotationProcessor 'org.projectlombok:lombok:1.18.20' + testCompileOnly 'org.projectlombok:lombok:1.18.20' + testAnnotationProcessor 'org.projectlombok:lombok:1.18.20' +} + +test { + useJUnitPlatform() +} diff --git a/ReactiveStreams/gradle/wrapper/gradle-wrapper.jar b/ReactiveStreams/gradle/wrapper/gradle-wrapper.jar new file mode 100644 index 0000000..7454180 Binary files /dev/null and b/ReactiveStreams/gradle/wrapper/gradle-wrapper.jar differ diff --git a/ReactiveStreams/gradle/wrapper/gradle-wrapper.properties b/ReactiveStreams/gradle/wrapper/gradle-wrapper.properties new file mode 100644 index 0000000..ffed3a2 --- /dev/null +++ b/ReactiveStreams/gradle/wrapper/gradle-wrapper.properties @@ -0,0 +1,5 @@ +distributionBase=GRADLE_USER_HOME +distributionPath=wrapper/dists +distributionUrl=https\://services.gradle.org/distributions/gradle-7.2-bin.zip +zipStoreBase=GRADLE_USER_HOME +zipStorePath=wrapper/dists diff --git a/ReactiveStreams/gradlew b/ReactiveStreams/gradlew new file mode 100755 index 0000000..1b6c787 --- /dev/null +++ b/ReactiveStreams/gradlew @@ -0,0 +1,234 @@ +#!/bin/sh + +# +# Copyright © 2015-2021 the original authors. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +############################################################################## +# +# Gradle start up script for POSIX generated by Gradle. +# +# Important for running: +# +# (1) You need a POSIX-compliant shell to run this script. If your /bin/sh is +# noncompliant, but you have some other compliant shell such as ksh or +# bash, then to run this script, type that shell name before the whole +# command line, like: +# +# ksh Gradle +# +# Busybox and similar reduced shells will NOT work, because this script +# requires all of these POSIX shell features: +# * functions; +# * expansions «$var», «${var}», «${var:-default}», «${var+SET}», +# «${var#prefix}», «${var%suffix}», and «$( cmd )»; +# * compound commands having a testable exit status, especially «case»; +# * various built-in commands including «command», «set», and «ulimit». +# +# Important for patching: +# +# (2) This script targets any POSIX shell, so it avoids extensions provided +# by Bash, Ksh, etc; in particular arrays are avoided. +# +# The "traditional" practice of packing multiple parameters into a +# space-separated string is a well documented source of bugs and security +# problems, so this is (mostly) avoided, by progressively accumulating +# options in "$@", and eventually passing that to Java. +# +# Where the inherited environment variables (DEFAULT_JVM_OPTS, JAVA_OPTS, +# and GRADLE_OPTS) rely on word-splitting, this is performed explicitly; +# see the in-line comments for details. +# +# There are tweaks for specific operating systems such as AIX, CygWin, +# Darwin, MinGW, and NonStop. +# +# (3) This script is generated from the Groovy template +# https://github.com/gradle/gradle/blob/master/subprojects/plugins/src/main/resources/org/gradle/api/internal/plugins/unixStartScript.txt +# within the Gradle project. +# +# You can find Gradle at https://github.com/gradle/gradle/. +# +############################################################################## + +# Attempt to set APP_HOME + +# Resolve links: $0 may be a link +app_path=$0 + +# Need this for daisy-chained symlinks. +while + APP_HOME=${app_path%"${app_path##*/}"} # leaves a trailing /; empty if no leading path + [ -h "$app_path" ] +do + ls=$( ls -ld "$app_path" ) + link=${ls#*' -> '} + case $link in #( + /*) app_path=$link ;; #( + *) app_path=$APP_HOME$link ;; + esac +done + +APP_HOME=$( cd "${APP_HOME:-./}" && pwd -P ) || exit + +APP_NAME="Gradle" +APP_BASE_NAME=${0##*/} + +# Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script. +DEFAULT_JVM_OPTS='"-Xmx64m" "-Xms64m"' + +# Use the maximum available, or set MAX_FD != -1 to use that value. +MAX_FD=maximum + +warn () { + echo "$*" +} >&2 + +die () { + echo + echo "$*" + echo + exit 1 +} >&2 + +# OS specific support (must be 'true' or 'false'). +cygwin=false +msys=false +darwin=false +nonstop=false +case "$( uname )" in #( + CYGWIN* ) cygwin=true ;; #( + Darwin* ) darwin=true ;; #( + MSYS* | MINGW* ) msys=true ;; #( + NONSTOP* ) nonstop=true ;; +esac + +CLASSPATH=$APP_HOME/gradle/wrapper/gradle-wrapper.jar + + +# Determine the Java command to use to start the JVM. +if [ -n "$JAVA_HOME" ] ; then + if [ -x "$JAVA_HOME/jre/sh/java" ] ; then + # IBM's JDK on AIX uses strange locations for the executables + JAVACMD=$JAVA_HOME/jre/sh/java + else + JAVACMD=$JAVA_HOME/bin/java + fi + if [ ! -x "$JAVACMD" ] ; then + die "ERROR: JAVA_HOME is set to an invalid directory: $JAVA_HOME + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." + fi +else + JAVACMD=java + which java >/dev/null 2>&1 || die "ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH. + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." +fi + +# Increase the maximum file descriptors if we can. +if ! "$cygwin" && ! "$darwin" && ! "$nonstop" ; then + case $MAX_FD in #( + max*) + MAX_FD=$( ulimit -H -n ) || + warn "Could not query maximum file descriptor limit" + esac + case $MAX_FD in #( + '' | soft) :;; #( + *) + ulimit -n "$MAX_FD" || + warn "Could not set maximum file descriptor limit to $MAX_FD" + esac +fi + +# Collect all arguments for the java command, stacking in reverse order: +# * args from the command line +# * the main class name +# * -classpath +# * -D...appname settings +# * --module-path (only if needed) +# * DEFAULT_JVM_OPTS, JAVA_OPTS, and GRADLE_OPTS environment variables. + +# For Cygwin or MSYS, switch paths to Windows format before running java +if "$cygwin" || "$msys" ; then + APP_HOME=$( cygpath --path --mixed "$APP_HOME" ) + CLASSPATH=$( cygpath --path --mixed "$CLASSPATH" ) + + JAVACMD=$( cygpath --unix "$JAVACMD" ) + + # Now convert the arguments - kludge to limit ourselves to /bin/sh + for arg do + if + case $arg in #( + -*) false ;; # don't mess with options #( + /?*) t=${arg#/} t=/${t%%/*} # looks like a POSIX filepath + [ -e "$t" ] ;; #( + *) false ;; + esac + then + arg=$( cygpath --path --ignore --mixed "$arg" ) + fi + # Roll the args list around exactly as many times as the number of + # args, so each arg winds up back in the position where it started, but + # possibly modified. + # + # NB: a `for` loop captures its iteration list before it begins, so + # changing the positional parameters here affects neither the number of + # iterations, nor the values presented in `arg`. + shift # remove old arg + set -- "$@" "$arg" # push replacement arg + done +fi + +# Collect all arguments for the java command; +# * $DEFAULT_JVM_OPTS, $JAVA_OPTS, and $GRADLE_OPTS can contain fragments of +# shell script including quotes and variable substitutions, so put them in +# double quotes to make sure that they get re-expanded; and +# * put everything else in single quotes, so that it's not re-expanded. + +set -- \ + "-Dorg.gradle.appname=$APP_BASE_NAME" \ + -classpath "$CLASSPATH" \ + org.gradle.wrapper.GradleWrapperMain \ + "$@" + +# Use "xargs" to parse quoted args. +# +# With -n1 it outputs one arg per line, with the quotes and backslashes removed. +# +# In Bash we could simply go: +# +# readarray ARGS < <( xargs -n1 <<<"$var" ) && +# set -- "${ARGS[@]}" "$@" +# +# but POSIX shell has neither arrays nor command substitution, so instead we +# post-process each arg (as a line of input to sed) to backslash-escape any +# character that might be a shell metacharacter, then use eval to reverse +# that process (while maintaining the separation between arguments), and wrap +# the whole thing up as a single "set" statement. +# +# This will of course break if any of these variables contains a newline or +# an unmatched quote. +# + +eval "set -- $( + printf '%s\n' "$DEFAULT_JVM_OPTS $JAVA_OPTS $GRADLE_OPTS" | + xargs -n1 | + sed ' s~[^-[:alnum:]+,./:=@_]~\\&~g; ' | + tr '\n' ' ' + )" '"$@"' + +exec "$JAVACMD" "$@" diff --git a/ReactiveStreams/gradlew.bat b/ReactiveStreams/gradlew.bat new file mode 100644 index 0000000..ac1b06f --- /dev/null +++ b/ReactiveStreams/gradlew.bat @@ -0,0 +1,89 @@ +@rem +@rem Copyright 2015 the original author or authors. +@rem +@rem Licensed under the Apache License, Version 2.0 (the "License"); +@rem you may not use this file except in compliance with the License. +@rem You may obtain a copy of the License at +@rem +@rem https://www.apache.org/licenses/LICENSE-2.0 +@rem +@rem Unless required by applicable law or agreed to in writing, software +@rem distributed under the License is distributed on an "AS IS" BASIS, +@rem WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +@rem See the License for the specific language governing permissions and +@rem limitations under the License. +@rem + +@if "%DEBUG%" == "" @echo off +@rem ########################################################################## +@rem +@rem Gradle startup script for Windows +@rem +@rem ########################################################################## + +@rem Set local scope for the variables with windows NT shell +if "%OS%"=="Windows_NT" setlocal + +set DIRNAME=%~dp0 +if "%DIRNAME%" == "" set DIRNAME=. +set APP_BASE_NAME=%~n0 +set APP_HOME=%DIRNAME% + +@rem Resolve any "." and ".." in APP_HOME to make it shorter. +for %%i in ("%APP_HOME%") do set APP_HOME=%%~fi + +@rem Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script. +set DEFAULT_JVM_OPTS="-Xmx64m" "-Xms64m" + +@rem Find java.exe +if defined JAVA_HOME goto findJavaFromJavaHome + +set JAVA_EXE=java.exe +%JAVA_EXE% -version >NUL 2>&1 +if "%ERRORLEVEL%" == "0" goto execute + +echo. +echo ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH. +echo. +echo Please set the JAVA_HOME variable in your environment to match the +echo location of your Java installation. + +goto fail + +:findJavaFromJavaHome +set JAVA_HOME=%JAVA_HOME:"=% +set JAVA_EXE=%JAVA_HOME%/bin/java.exe + +if exist "%JAVA_EXE%" goto execute + +echo. +echo ERROR: JAVA_HOME is set to an invalid directory: %JAVA_HOME% +echo. +echo Please set the JAVA_HOME variable in your environment to match the +echo location of your Java installation. + +goto fail + +:execute +@rem Setup the command line + +set CLASSPATH=%APP_HOME%\gradle\wrapper\gradle-wrapper.jar + + +@rem Execute Gradle +"%JAVA_EXE%" %DEFAULT_JVM_OPTS% %JAVA_OPTS% %GRADLE_OPTS% "-Dorg.gradle.appname=%APP_BASE_NAME%" -classpath "%CLASSPATH%" org.gradle.wrapper.GradleWrapperMain %* + +:end +@rem End local scope for the variables with windows NT shell +if "%ERRORLEVEL%"=="0" goto mainEnd + +:fail +rem Set variable GRADLE_EXIT_CONSOLE if you need the _script_ return code instead of +rem the _cmd.exe /c_ return code! +if not "" == "%GRADLE_EXIT_CONSOLE%" exit 1 +exit /b 1 + +:mainEnd +if "%OS%"=="Windows_NT" endlocal + +:omega diff --git a/ReactiveStreams/settings.gradle b/ReactiveStreams/settings.gradle new file mode 100644 index 0000000..b979a36 --- /dev/null +++ b/ReactiveStreams/settings.gradle @@ -0,0 +1 @@ +rootProject.name = 'ReactiveStreams' diff --git a/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/ReactiveStreamsApplication.java b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/ReactiveStreamsApplication.java new file mode 100644 index 0000000..8adbd35 --- /dev/null +++ b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/ReactiveStreamsApplication.java @@ -0,0 +1,198 @@ +package io.pinest94.reactivestreams; + +import java.util.Queue; +import java.util.concurrent.Callable; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.function.Consumer; +import java.util.function.Function; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.ApplicationRunner; +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.context.annotation.Bean; +import org.springframework.http.ResponseEntity; +import org.springframework.http.client.Netty4ClientHttpRequestFactory; +import org.springframework.scheduling.annotation.Async; +import org.springframework.scheduling.annotation.AsyncResult; +import org.springframework.scheduling.annotation.EnableAsync; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import org.springframework.stereotype.Component; +import org.springframework.util.concurrent.ListenableFuture; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.RestController; +import org.springframework.web.client.AsyncRestTemplate; +import org.springframework.web.context.request.async.DeferredResult; +import org.springframework.web.servlet.mvc.method.annotation.ResponseBodyEmitter; + +import io.netty.channel.nio.NioEventLoopGroup; +import io.pinest94.reactivestreams.completion.Completion; +import lombok.extern.slf4j.Slf4j; + +@Slf4j +@SpringBootApplication +@EnableAsync +public class ReactiveStreamsApplication { + + /*** + * 해당 클래스에서 알아야할 포인트 + * 1. run 메소드가 어떻게 동작할지 생각하기 + * 2. 실제 서비스에서는 @Async만 절대 사용하지 않는다. 이유는 @Async는 요청이 들어오면 새로운 스레드를 생성하지만 스레드풀로 관리하거나 캐시가 따로 존재하지 않는다. + * 즉 요청이 들어오면 계속 생성을 하기 때문에 메모리와 CPU에 상당한 부담을 주게된다. + * 3. 실제 서비스에서는 @Bean으로 ThreadPoolTaskExecutor를 만들어서 Thread를 관리하여 @Async를 사용하도록 한다. + * 4. @Bean으로 정의된 ThreadPoolTaskExecutor가 있다면 @Async는 디폴트로 빈에 정의된 것을 사용한다. + * 5. 각 메서드마다 스레드 풀 정책을 달리하고 싶어서 정의된 것이 여러 개인 경우에는 @Async(value = "tp")으로 사용하면 된다. + */ + + @RestController + public static class MyController { + + Queue> results = new ConcurrentLinkedQueue<>(); + static String description = "Hi! My name is hansol kim"; + AsyncRestTemplate rt = new AsyncRestTemplate( + new Netty4ClientHttpRequestFactory(new NioEventLoopGroup(1))); + + @Autowired + MyService myService; + + @GetMapping("/callable") + public Callable callable() throws InterruptedException { + log.info("callable"); + return () -> { + log.info("async"); + Thread.sleep(5000); + return "hello"; + }; + } + + @GetMapping("/dr") + public DeferredResult deferred() { + log.info("dr"); + DeferredResult dr = new DeferredResult<>(); + results.add(dr); + return dr; + } + + @GetMapping("/dr/count") + public String drcount() { + return String.valueOf(results.size()); + } + + @GetMapping("/dr/event") + public String drevent(String msg) { + for (DeferredResult dr : results) { + dr.setResult("Hello " + msg); + results.remove(dr); + } + return "OK"; + } + + @GetMapping("/async") + public String async() throws InterruptedException { + log.info("async"); + Thread.sleep(5000); + return "hello"; + } + + @GetMapping("/emitter") + public ResponseBodyEmitter emitter() { + ResponseBodyEmitter emitter = new ResponseBodyEmitter(); + + Executors.newSingleThreadExecutor().submit(() -> { + try { + emitter.send("

"); + for (int i = 0; i <= description.length(); ++i) { + emitter.send(description.charAt(i)); + Thread.sleep(100); + } + emitter.send("

"); + } catch (Exception e) {} + }); + + return emitter; + } + + /*** + * 외부 API를 호출하는데 호출되는 곳에서 로직이 오래걸리는 경우(실습용으로 2초 설정) + * 현재 RestTemplate의 getForObject 메소드는 blocking 메소드이다. + * 즉, 요청을 보내고 응답이 올때까지 blocking이 된다는 뜻이다. + * 물론 요청을 보내고 바로 응답이 온다면 별 문제가 없겠지만 2초가 걸리게 되면(오래걸리게 되면) 해당 스레드는 요청을 보내고 2초간 놀게 된다. + * 2초동안 스레드가 놀게된다는 것은 다른 요청을 받지 못하고 대기상태가 되어 해당 서버컴퓨터의 CPU가 놀게된다는 뜻을 의미한다. 상당히 효율적이지 못하다. + * 그래서 non-blocking을 제공하는 AsyncRestTemplate을 사용하여 실습을 진행했다. + * @return + */ + + private static final String URL1 = "http://localhost:8081/service?req={req}"; + private static final String URL2 = "http://localhost:8081/service2?req={req}"; + + @GetMapping("/rest/{idx}") + public DeferredResult rest(@PathVariable String idx) { + DeferredResult dr = new DeferredResult<>(); + + toCF(rt.getForEntity( + URL1, String.class, "hello" + idx)) + .thenCompose(s -> toCF(rt.getForEntity(URL2, String.class, s.getBody()))) + .thenApplyAsync(s2 -> myService.work(s2.getBody())) + .thenAccept(s3 -> dr.setResult(s3)) + .exceptionally(e -> { + dr.setErrorResult(e.getMessage()); + return null; + }); + return null; + } + + CompletableFuture toCF(ListenableFuture lf) { + CompletableFuture cf = new CompletableFuture(); + lf.addCallback(s -> cf.complete(s), e -> cf.completeExceptionally(e)); + return cf; + } + } + + @Component + public static class MyService { + public String work(String req) { + return req + "/asyncwork"; + } + } + + @Bean + ThreadPoolTaskExecutor tp() { + /*** + * setCorePoolSize, setMaxPoolSize, setQueueCapacity를 알아보자 + * 기본적으로 스레드의 풀사이즈는 setCorePoolSize의 설정 값을 따른다. 여기까지는 어렵지 않다. + * 만약 10개의 스레드가 모두 사용 중이고 11번째 스레드 사용 요청이 들어오면 어떻게 될까? + * 정답은 setQueueCapacity으로 설정된 큐에 해당 요청을 대기시킨다. + * 그래서 해당 큐가 모두 꽉차게 될때(=setQueueCapacity) 스레드를 증가시키고 최대 setMaxPoolSize값만큼 증가시킨다. + * 주의! 11번째 스레드 요청이 들어올 때 setMaxPoolSize값 만큼 스레드가 증가하는 것이 아니다. + */ + ThreadPoolTaskExecutor threadPoolTaskExecutor = new ThreadPoolTaskExecutor(); + threadPoolTaskExecutor.setCorePoolSize(10); + threadPoolTaskExecutor.setMaxPoolSize(100); + threadPoolTaskExecutor.setQueueCapacity(200); + threadPoolTaskExecutor.setThreadNamePrefix("rx : "); + threadPoolTaskExecutor.initialize(); + return threadPoolTaskExecutor; + } + + public static void main(String[] args) { + SpringApplication.run(ReactiveStreamsApplication.class, args); + } + + @Autowired + MyService myService; + + @Bean + ApplicationRunner run() { + return args -> { + log.info("run()"); + // Future future = myService.hello(); + // log.info("result : " + future.isDone()); + // log.info("result : " + future.get()); + }; + } + +} diff --git a/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completableFuture/CFuture.java b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completableFuture/CFuture.java new file mode 100644 index 0000000..c773f61 --- /dev/null +++ b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completableFuture/CFuture.java @@ -0,0 +1,36 @@ +package io.pinest94.reactivestreams.completableFuture; + +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.ForkJoinPool; +import java.util.concurrent.TimeUnit; + +import lombok.extern.slf4j.Slf4j; + +@Slf4j +public class CFuture { + public static void main(String[] args) throws ExecutionException, InterruptedException { + ExecutorService es = Executors.newFixedThreadPool(10); + + CompletableFuture + .supplyAsync(() -> { + log.info("runAsync"); + return 1; + }, es) + .thenApply(s -> { + log.info("then Apply {}", s); + return s + 1; + }) + .thenApplyAsync(s2 -> { + log.info("then Apply {}", s2); + return s2 * 3; + }, es) + .exceptionally(e -> -10) + .thenAcceptAsync(s3 -> log.info("then Accept {}", s3), es); + log.info("exit"); + ForkJoinPool.commonPool().shutdown(); + ForkJoinPool.commonPool().awaitTermination(10, TimeUnit.SECONDS); + } +} diff --git a/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/AcceptCompletion.java b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/AcceptCompletion.java new file mode 100644 index 0000000..9edbacc --- /dev/null +++ b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/AcceptCompletion.java @@ -0,0 +1,18 @@ +package io.pinest94.reactivestreams.completion; + +import java.util.function.Consumer; + +import org.springframework.http.ResponseEntity; +import org.springframework.util.concurrent.ListenableFuture; + +public class AcceptCompletion extends Completion { + public Consumer con; + public AcceptCompletion(Consumer con) { + this.con = con; + } + + @Override + void run(S value) { + con.accept(value); + } +} diff --git a/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/ApplyCompletion.java b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/ApplyCompletion.java new file mode 100644 index 0000000..300f829 --- /dev/null +++ b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/ApplyCompletion.java @@ -0,0 +1,19 @@ +package io.pinest94.reactivestreams.completion; + +import java.util.function.Function; + +import org.springframework.http.ResponseEntity; +import org.springframework.util.concurrent.ListenableFuture; + +public class ApplyCompletion extends Completion { + Function> fn; + public ApplyCompletion(Function> fn) { + this.fn = fn; + } + + @Override + void run(S value) { + ListenableFuture lf = fn.apply(value); + lf.addCallback(s->complete(s), e->error(e)); + } +} diff --git a/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/Completion.java b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/Completion.java new file mode 100644 index 0000000..7ff8364 --- /dev/null +++ b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/Completion.java @@ -0,0 +1,48 @@ +package io.pinest94.reactivestreams.completion; + +import java.util.function.Consumer; +import java.util.function.Function; + +import org.springframework.http.ResponseEntity; +import org.springframework.util.concurrent.ListenableFuture; + +public class Completion { + Consumer> con; + Completion next; + + public Completion() {} + + public void andAccept(Consumer con) { + Completion c = new AcceptCompletion(con); + this.next = c; + } + + public Completion andApply( + Function> fn) { + Completion c = new ApplyCompletion(fn); + this.next = c; + return c; + } + + public Completion andError(Consumer econ) { + Completion c = new ErrorCompletion(econ); + this.next = c; + return c; + } + + public static Completion from(ListenableFuture listenableFuture) { + Completion c = new Completion(); + listenableFuture.addCallback(s -> c.complete(s), e -> c.error(e)); + return c; + } + + protected void error(Throwable e) { + if(next != null) next.error(e); + } + + protected void complete(T s) { + if (next != null) { next.run(s); } + } + + void run(S value) {} +} \ No newline at end of file diff --git a/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/ErrorCompletion.java b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/ErrorCompletion.java new file mode 100644 index 0000000..6eb80b1 --- /dev/null +++ b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/completion/ErrorCompletion.java @@ -0,0 +1,18 @@ +package io.pinest94.reactivestreams.completion; + +import java.util.function.Consumer; + +import org.springframework.http.ResponseEntity; + +public class ErrorCompletion extends Completion { + public Consumer econ; + + public ErrorCompletion(Consumer econ) { + this.econ = econ; + } + + @Override + void run(T value) { + if (next != null) { next.run(value); } + } +} diff --git a/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/futures/FutureEx.java b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/futures/FutureEx.java new file mode 100644 index 0000000..de9fb81 --- /dev/null +++ b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/futures/FutureEx.java @@ -0,0 +1,73 @@ +package io.pinest94.reactivestreams.futures; + +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.FutureTask; + +import lombok.extern.slf4j.Slf4j; + +@Slf4j +public class FutureEx { + + /*** + * 해당 클래스에서 알아야할 포인트 + * 1. future1,2를 실행하는 것과 job1,2를 실행하는 것의 차이 + * 2. future와 futureTask의 차이 + * 3. futureTask에서 done 메소드를 오버라이드한 익명 메소드내에서 get 메소드 호출이 시사하는 바를 알기 + * 4. + */ + + public static void main(String[] args) throws ExecutionException, InterruptedException { + ExecutorService executorService = Executors.newCachedThreadPool(); + + Future future = executorService.submit(() -> { + Thread.sleep(3000); + log.info("Async"); + return "Hello"; + }); + + FutureTask futureTask = new FutureTask(() -> { + Thread.sleep(2000); + log.info("Async2"); + return "Hello2"; + }) { + @Override + protected void done() { + try { + log.info(get()); + } catch (InterruptedException e) { + e.printStackTrace(); + } catch (ExecutionException e) { + e.printStackTrace(); + } + } + }; + + Long startTime = System.currentTimeMillis(); + + executorService.execute(futureTask); + log.info(future.get()); + // log.info(futureTask.get()); + // log.info(job()); + // log.info(job2()); + + Long endTime = System.currentTimeMillis(); + + log.info("elapsedTime : " + (endTime - startTime)); + executorService.shutdown(); + } + + public static String job() throws InterruptedException { + Thread.sleep(2000); + log.info("Async"); + return "Hello"; + } + + public static String job2() throws InterruptedException { + Thread.sleep(2000); + log.info("Async"); + return "Hello"; + } +} diff --git a/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/futures/LoadTest.java b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/futures/LoadTest.java new file mode 100644 index 0000000..d09a4fd --- /dev/null +++ b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/futures/LoadTest.java @@ -0,0 +1,57 @@ +package io.pinest94.reactivestreams.futures; + +import java.util.concurrent.BrokenBarrierException; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.springframework.util.StopWatch; +import org.springframework.web.client.RestTemplate; + +import lombok.extern.slf4j.Slf4j; + +@Slf4j +public class LoadTest { + + static AtomicInteger counter = new AtomicInteger(0); + + public static void main(String[] args) throws InterruptedException, BrokenBarrierException { + ExecutorService executorService = Executors.newFixedThreadPool(100 ); + + RestTemplate restTemplate = new RestTemplate(); + String url = "http://localhost:8080/rest"; + + CyclicBarrier barrier = new CyclicBarrier(101); + + StopWatch main = new StopWatch(); + main.start(); + + for(int i=0; i<100; i++) { + executorService.submit(() -> { + int idx = counter.addAndGet(1); + + barrier.await(); // 스레드가 100개 만들어질때까지 기다림 + + log.info("Thread : " + idx); + + StopWatch sw = new StopWatch(); + sw.start(); + + String res = restTemplate.getForObject(url+"/"+idx, String.class); + + sw.stop(); + log.info("Elapsed Time : {}, {} / {}", idx, sw.getTotalTimeSeconds(), res); + return null; + }); + } + + barrier.await(); + executorService.shutdown(); + executorService.awaitTermination(100, TimeUnit.SECONDS); + + main.stop(); + log.info("Total : {}", main.getTotalTimeSeconds()); + } +} diff --git a/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/futures/RemoteService.java b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/futures/RemoteService.java new file mode 100644 index 0000000..36dc98f --- /dev/null +++ b/ReactiveStreams/src/main/java/io/pinest94/reactivestreams/futures/RemoteService.java @@ -0,0 +1,32 @@ +package io.pinest94.reactivestreams.futures; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.RequestParam; +import org.springframework.web.bind.annotation.RestController; + +@SpringBootApplication +public class RemoteService { + + @RestController + public static class MyController { + @GetMapping("/service") + public String service(@RequestParam("req") String req) throws InterruptedException { + Thread.sleep(1000); + return req + "/service"; + } + + @GetMapping("/service2") + public String service2(@RequestParam("req") String req) throws InterruptedException { + Thread.sleep(1000); + return req + "/service2"; + } + } + + public static void main(String[] args) { + System.setProperty("server.port", "8081"); + System.setProperty("server.tomcat.threads.max", "1000"); + SpringApplication.run(RemoteService.class, args); + } +} diff --git a/ReactiveStreams/src/main/resources/application.yaml b/ReactiveStreams/src/main/resources/application.yaml new file mode 100644 index 0000000..98fce2c --- /dev/null +++ b/ReactiveStreams/src/main/resources/application.yaml @@ -0,0 +1,4 @@ +server: + tomcat: + threads: + max: 1 \ No newline at end of file diff --git a/ReactiveStreams/src/test/java/io/pinest94/reactivestreams/ReactiveStreamsApplicationTests.java b/ReactiveStreams/src/test/java/io/pinest94/reactivestreams/ReactiveStreamsApplicationTests.java new file mode 100644 index 0000000..9765029 --- /dev/null +++ b/ReactiveStreams/src/test/java/io/pinest94/reactivestreams/ReactiveStreamsApplicationTests.java @@ -0,0 +1,13 @@ +package io.pinest94.reactivestreams; + +import org.junit.jupiter.api.Test; +import org.springframework.boot.test.context.SpringBootTest; + +@SpringBootTest +class ReactiveStreamsApplicationTests { + + @Test + void contextLoads() { + } + +}