1. (Stream)Stream을 병렬로 처리하기

이월(0216tw)·2024년 5월 16일

스트림의 장점 중 하나는 실행할 스트림을 병렬로 처리할 수 있다는 것이다.

아래와 같이 1부터 10까지 저장된 List 가 있습니다.
각각의 데이터에 대해 100을 곱해서 100 , 200 , 300 , ... 처럼 만들려고 합니다.
이때 병렬로 처리되도록 코드를 구현해주세요.

beforeEach 세팅

@BeforeEach
public void setUp() {
    int intArray[] = new int[10];
    for(int i =1 ; i<intArray.length; i++) {
        intArray[i] = i;
    }

    //int[] 을 ArrayList 로 변환
    numbers = Arrays.stream(intArray).boxed().collect(Collectors.toCollection(ArrayList::new));     
}

[코드 설명]

  • Arrays.stream()
    int[] 배열을 IntStream 객체로 변환한다.
  • boxed()
    변환된 IntStream(기본타입스트림) 을 Stream<Integer> 인 박싱된 객체 스트림으로 변환
  • collect()
    위의 스트림 요소(Element)를 수집(콜렉팅)하는 메서드
  • Collectors.toCollection(ArrayList::new)
    콜렉팅한 대상을 새로 생성한 ArrayList인스턴스에 수집하도록 설정

단일 처리 예시와 걸린 시간

@Test
public void 단일로_처리하기() {
    long start = System.currentTimeMillis();

    numbers.stream().map(number -> {
        try {
            Thread.sleep(1000);
            System.out.println("Processing number: " + number + " on thread: " + Thread.currentThread().getName());
            return number * 100;
         } catch (InterruptedException e) {
            e.printStackTrace();
            return null;
         }
     }).forEach(result -> System.out.println("Result: " + result + " on thread: " + Thread.currentThread().getName()));

     long end = System.currentTimeMillis();

     System.out.println("스트림 단일 처리 시간 => " + (end-start));
   }
   
출력결과 
Processing number: 0 on thread: main
Result: 0 on thread: main
Processing number: 1 on thread: main
Result: 100 on thread: main
Processing number: 2 on thread: main
Result: 200 on thread: main
Processing number: 3 on thread: main
Result: 300 on thread: main
Processing number: 4 on thread: main
Result: 400 on thread: main
Processing number: 5 on thread: main
Result: 500 on thread: main
Processing number: 6 on thread: main
Result: 600 on thread: main
Processing number: 7 on thread: main
Result: 700 on thread: main
Processing number: 8 on thread: main
Result: 800 on thread: main
Processing number: 9 on thread: main
Result: 900 on thread: main
스트림 단일 처리 시간 => 10055

단일 스레드이므로 main스레드가 해당 작업을 모두 처리했다.
각 작업이 1초 걸린다는 가정을 만들기 위해 sleep(1000)을 추가했으며
10번의 작업이 실행되므로 약 10초 정도가 걸렸다.

[코드 설명]

  • numbers.stream()
    현재 numbers 는 ArrayList<Integer>이다. 이를 스트림객체로 변환한다. (ArrayList<Integer> -> Stream<Integer>)

  • map( el -> 실행할함수or코드 )
    각 요소(el)에 대해 처리할 메서드나 코드로 변환해 새로운 스트림 객체를 반환한다.
    여기서는 1초 대기후 출력 , 그리고 요소에 대해 100을 곱한 값을 다시 리턴하고 있다.
    map은 중간연산이므로 다른 중간 연산과 체이닝이 가능하다.

  • forEach( el -> 실행할함수or코드)
    각 요소(el)에 대해 주어진 동작을 수행한다.
    최종연산이기 때문에 스트림의 처리를 종료한다.


잠깐! map() 과 forEach() 의 차이

map()은 중간연산이기에 다른 중간연산(filter() , map() , sorted() 등..) 과 체이닝 하여 데이터 처리를 할 수 있다.

forEach은 최종연산이기에 이후에 중간연산을 추가하면 컴파일 에러가 발생한다.



병렬 처리 예시와 걸린 시간

@Test
    public void 스트림_병렬처리하기() {

        long start = System.currentTimeMillis();

        numbers.parallelStream()
                // 각 요소를 제곱하여 병렬로 처리
                .map(number -> {
                    try {
                        Thread.sleep(1000);
                        System.out.println("Processing number: " + number + " on thread: " + Thread.currentThread().getName());
                        return number * 100;
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                        return null;
                    }
                })
                // 결과를 출력
                .forEach(result -> System.out.println("Result: " + result + " on thread: " + Thread.currentThread().getName()));
        long end = System.currentTimeMillis();
        System.out.println("스트림 병렬 처리 시간 => " + (end - start));
    }
    
   
출력결과 
Processing number: 5 on thread: ForkJoinPool-1-worker-5
Processing number: 1 on thread: ForkJoinPool-1-worker-4
Processing number: 4 on thread: ForkJoinPool-1-worker-8
Processing number: 6 on thread: ForkJoinPool-1-worker-1
Processing number: 7 on thread: ForkJoinPool-1-worker-6
Result: 100 on thread: ForkJoinPool-1-worker-4
Result: 400 on thread: ForkJoinPool-1-worker-8
Result: 700 on thread: ForkJoinPool-1-worker-6
Processing number: 9 on thread: ForkJoinPool-1-worker-7
Result: 900 on thread: ForkJoinPool-1-worker-7
Processing number: 2 on thread: ForkJoinPool-1-worker-2
Result: 200 on thread: ForkJoinPool-1-worker-2
Processing number: 8 on thread: ForkJoinPool-1-worker-3
Result: 600 on thread: ForkJoinPool-1-worker-1
Result: 500 on thread: ForkJoinPool-1-worker-5
Result: 800 on thread: ForkJoinPool-1-worker-3
Processing number: 3 on thread: ForkJoinPool-1-worker-6
Result: 300 on thread: ForkJoinPool-1-worker-6
Processing number: 0 on thread: ForkJoinPool-1-worker-4
Result: 0 on thread: ForkJoinPool-1-worker-4
스트림 병렬 처리 시간 => 2023

병렬 스레드이므로 다수의 스레드가 사용된 것을 볼 수 있다. (총 8개의 스레드 사용됨)
동일하게 각 작업이 1초 걸린다는 가정을 만들기 위해 sleep(1000)을 추가했으며
worker-4와 worker-6이 한번 더 실행을 했기 때문에 결과적으로 약 2초가 걸렸다.

(1초 : worker -> 1,2,3,4,5,6,7,8 이 작업 한번씩 총 8번 실행)
(2초 : worker -> 4 , 6 이 작업 한번씩 총 2번 실행)


[코드 설명]

  • numbers.parallelStream()
    병렬처리가 가능한 스트림객체를 반환한다.

  • 그 외 작업은 단일 처리 예시와 동일

profile
#SQLD강사 #AI개발 #AI강사 #개발자 개발도 하고 강의도 하지만 고민을 제일 많이 합니다

0개의 댓글