블로그

  • [Reactive Extensions] Hot 변환은 어떤 때에 필요한가?

    환경

    UniRx 입문 시리즈 목차는 이쪽

    Hot 변환 한 포인트​

    여러 상황이 있지만 가장 Hot 변환이 중요해지는 상황은 하나의 스트림을 여러 번 Subscribe하는 경우 입니다.

    예) 입력된 문자열이 특정 키워드와 일치하는지 검사​

    Hot 변환이 필요한 예로 “입력된 키 입력을 감지하고 4 문자의 특정 키워드가 입력되었는지를 알아 내는 스트림”을 만들어 보겠습니다.

    준비​

    우선은 준비 단계로 입력된 키 정보를 4글자씩 내놓는 스트림을 만듭니다.

    var keyBufferStream
        = Observable.FromEvent<KeyEventHandler, KeyEventArgs>(
            h => (sender, e) => h(e),
            h => KeyDown += h,
            h => KeyDown -= h)
            .Select(x => x.Key.ToString())// 입력 키를 문자로 변환
            .Buffer(4, 1)// 4 개씩 정리
            .Select(x => x.Aggregate((p, c) => p + c));// 문자에서 문자열로 변환
    
    // 결과를 표시하고 보니
    keyBufferStream.Subscribe(Console.WriteLine);
    

    실행결과 예(ABCDEFGH 키 입력 결과)

    ABCD
    BCDE
    CDEF
    DEFG
    EFGH
    

    이같이 keyBufferStream은 입력 키가 4자씩으로 뭉치고 흐르는 스트림입니다.

    역주: 아래는 유니티에서 실행 가능한 예제 입니다.

    var keyBufferStream = this.UpdateAsObservable()
        .Where(_ => Input.anyKeyDown)// 아무 버튼 눌렀을 때
        .Where(_ => !(Input.GetMouseButtonDown(0) || Input.GetMouseButtonDown(1) || Input.GetMouseButtonDown(2)))// 마우스는 무시
        .Select(_ => Input.inputString)// 버튼 스트링
        .Buffer(4, 1)// 4 개씩 정리
        .Select(x => x.Aggregate((p, c) => p + c));// 문자에서 문자열로 변환
    
    // 결과 표시
    keyBufferStream.Subscribe(Debug.Log);
    

    Aggregate는 Linq 에서 지원 하는 메서드이며, 집계 연산자 입니다. 마지막 값을 돌려주는 메소드 입니다.

    keyBufferStream을 사용하여 “HOGE” 또는 “FUGA”의 입력을 감시하자​

    그럼 이 keyBufferStream을 사용하여 “HOGE”와 “FUGA”를 감시해 봅시다.

    Where 사이 HOGE와 FUGA에서 2회 Subscribe 합니다.

    keyBufferStream.Where(x => x == "HOGE")
        .Subscribe(_ => Debug.Log("Input HOGE"));
    
    keyBufferStream.Where(x => x == "FUGA")
        .Subscribe(_ => Debug.Log("Input FUGA"));
    

    실행 결과 (HOGEFUGA 입력 한 결과)

    Input HOGE
    Input FUGA
    

    각각의 문자열에 반응하는 스트림을 만들어 Subscribe 할 수 있었습니다.

    만.. 이 스트림에는 커다란 문제가 있습니다.

    무엇이 문제인가?​

    상기 스트림은 무엇이 문제인가? 그것은 keyBufferStream이 Cold Observable로 형성되는 것이 문제 입니다. 이전의 포스트에서도 설명했지만, (Cold Observable은 분기하지 않습니다.) Subscribe 할 때마다 매번 새로운 스트림을 생성하는 특성이 있습니다.

    따라서 상기와 같은 작성을 해 버리면 다음과 같은 문제가 발생할 수 있습니다.

    • 뒤에서 다중 스트림이 생성되어 버립니다. 메모리와 CPU를 낭비합니다.
    • Subscribe 한 시점에서 따라 흘러 나오는 결과가 다릅니다. (참고 Cold Observable의 성질)
      • 역주: Cold Observable은 Subscribe 한 순간부터 오퍼레이터가 작동하게 됩니다. Subscribe 전에 온 메시지는 모든 처리 조차 되지 않고 소멸 됩니다.

    스트림이 2개로 흐른다는 증거

    var keyBufferStream
        = Observable.FromEvent<KeyEventHandler, KeyEventArgs>(
            h => (sender, e) => h(e),
            h => KeyDown += h,
            h => KeyDown -= h)
            .Select(x => x.Key.ToString())
            .Buffer(4, 1)
            .Do(_=> Console.WriteLine("Buffered")) // Buffer가 OnNext를 방출한 타이밍에 출력된다.
            .Select(x => x.Aggregate((p, c) => p + c));
    
    keyBufferStream
        .Where(x => x == "HOGE")
        .Subscribe(_ => Console.WriteLine("Input HOGE"));
    
    keyBufferStream
        .Where(x => x == "FUGA")
        .Subscribe(_ => Console.WriteLine("Input FUGA"));

    실행 결과(AAAA와 Buffer가 1번만 움직이도록 키 입력)

    Buffered
    Buffered // Buffer는 1회만 흐르고 있을텐데 2번 출력되고 있다 = 스트림이 2개로 흐르고 있다.

    Hot Observable이 스트림의 근원인 FromEvent 밖에 없기 때문에, Subscribe 할 때마다 FromEvent로부터 새롭게 스트림이 생성되어 버리는 움직임이 되고 있습니다.

    문제의 해결책 “Hot 변환”​

    여기에서 첫번째 “Hot 변환은 하나의 스트림을 동시에 여러 Subscribe하는 경우에 사용한다” 라는 이야기로 돌아갑니다.

    즉 Hot 변환하여 스트림의 분기점을 만들어 여러 Subscribe 했을 때 스트림을 하나로 통합 할 수 있게 되는 것입니다.

    Hot 변환 된 예

    var keyBufferStream
        = Observable.FromEvent<KeyEventHandler, KeyEventArgs>(
            h => (sender, e) => h(e),
            h => KeyDown += h,
            h => KeyDown -= h)
            .Select(x => x.Key.ToString())
            .Buffer(4, 1)
            .Select(x => x.Aggregate((p, c) => p + c))
            .Publish() // Publish에서 Hot 변환(Publish가 대표하여 Subscribe 해 준다)
            .RefCount(); // RefCount은 Observer가 추가되었을 때 자동 Connect 해 주는 오퍼레이터.
    
    keyBufferStream
        .Where(x => x == "HOGE")
        .Subscribe(_ => Console.WriteLine("Input HOGE"));
    
    keyBufferStream
        .Where(x => x == "FUGA")
        .Subscribe(_ => Console.WriteLine("Input FUGA"));

    실행 결과(HOGEFUGA 입력)

    Input HOGE
    Input FUGA

    Hot 변환 방식에는 여러 가지가 있지만 가장 쉬운 것이 Publish()와 RefCount()를 결합 하는 것 입니다.

    이번에는 Hot 변환의 필요성에 대해 설명하고 싶기 때문에 Publish와 RefCount의 상세한 설명은 생략하겠습니다.(자세한 설명은 여기)

    역주: UniRx에서는 Publish()와 RefCount()의 결합인 Share() 오퍼레이터를 제공 합니다. 두개를 사용해야 될 경우에는 Share()를 사용하시면 됩니다.

    정리​

    • 스트림을 의도적으로 분기 하고 싶을 때 Hot 변환을 수행 한다.
    • 스트림을 생성하여 반환하는 속성과 함수를 정의하면 끝에 Hot 변환을 하는 것이 안전하다.
    • Hot 변환을 잊어 버리면 메모리나 CPU가 낭비되거나 Subscribe 타이밍이 어긋날 수 있다.
    • Hot 변환 오퍼레이터는 몇 개 있지만, Publish() + RefCount()의 조합이 편리하다 (만능은 아니다)
  • Hot과 Cold 대해


    환경

    UniRx 입문 시리즈 목차는 이쪽

    Cold Observable

    • 자발적으로 아무것도 하지 않는 수동적인 Observable
    • Observer이 등록되어 (Subscribe 되고) 처음 일을 시작한다.
    • 스트림의 전후를 그냥 연결만 한다. 스트림을 분기시키는 기능은 없다.

    Hot Observable

    • 자신이 값을 발행하는 능동적인 Observable
    • 후속 Observer의 존재에 관계없이 메시지를 발행한다
    • 자신보다 상류의 Cold Observable을 시작하고 값의 발행을 요구하는 기능을 가진다
    • 하류의 Observer를 모두 묶어, 정리해 같은 값을 발행한다 (스트림을 분기시킨다)

    Hot과 Cold 구분법​

    대부분의 오퍼레이터는 Cold인 성질이며, 자신이 명시적으로 스트림(stream)을 Hot으로 변환하지 않는 한 Cold 그대로 입니다.

    Hot 변환용 오퍼레이터는 이른바 Publish 계의 오퍼레이터가 해당 합니다.

    Hot에 대해​

    Hot Observable의 성질​

    스트림을 가동시키는 성질​

    Rx 스트림은 기본적으로 Subscribe가 된 순간에 각 오퍼레이터의 작동이 시작하게 되어 있습니다. 하지만 Hot Observable을 스트림 중간에 끼우는 것으로, Subscribe를 실행 이전에 스트림을 실행시킬 수 있습니다.

    스트림을 분기하는 성질​

    Hot Observable은 스트림을 분기 할 수 있습니다.

    Cold에 대해​

    Subscribe 될 때까지 작동하지 않는 성질​

    Cold Observable은 Subscribe될때 (또는 Hot 변환될때)까지 작동하지 않습니다. 마음이 없는 Observable 입니다.

    작동하지 않는 Cold Observable에 전달된 메시지는 모두 처리 조차 되지 않고 소멸 됩니다.

    특히 값의 발행 타이밍이나 전후 관계가 중요한 오퍼레이터를 사용하는 경우는, 어느 타이밍부터 처리가 시작되는지를 충분히 인지하고 사용하지 않으면 안됩니다. 같은 스트림 정의라도, Subscribe 한 타이밍에 따라서 동작이 바뀌어 버립니다. 아래가 그 예 입니다.

    var subject = new Subject<string>();
            
    // subject에서 생성된 Observable은 [Hot]
    var sourceObservable = subject.AsObservable();
    
    // 스트림에 흘러 들어온 문자열을 연결하여 새로운 문자열로 만드는 스트림
    // Scan()은 [Cold]
    var stringObservable = sourceObservable.Scan((p, c) => p + c);
    
    // 스트림에 값을 흘린다
    subject.OnNext("A");
    subject.OnNext("B");
    
    // 스트림에 값을 흘린 후 Subscribe 한다.
    stringObservable.Subscribe(Debug.Log);
    
    // Subscribe 후 스트림에 값을 흘린다.
    subject.OnNext("C");
    
    // 완료
    subject.OnCompleted();

    실행결과

    C

    위의 코드를 실행한 결과 C가 출력 될 것입니다.

    이것은 Scan 오퍼레이터가 Cold이기 때문에 Subscribe 전에 발행된 A 그리고 B 가 처리되지 않았기 때문입니다.

    만약 여기에서 “Subscribe하기 이전에 발급 된 값을 처리 했으면 좋겠다”는 경우는 어떻게 하면 좋을까요. 이 경우 Hot 변환 오퍼레이터를 끼워 Subscribe하기 이전에 스트림을 시작하면 좋을 것입니다.

    var subject = new Subject<string>();
            
    // subject에서 생성된 Observable은 [Hot]
    var sourceObservable = subject.AsObservable();
    
    // 스트림에 흘러 들어온 문자열을 연결하여 새로운 문자열로 만드는 스트림
    // Scan()은 [Cold]
    var stringObservable = sourceObservable
        .Scan((p, c) => p + c)
        .Publish(); // Hot 변환 오퍼레이터
    
    stringObservable.Connect(); // 스트림 가동 개시
    
    // 스트림에 값을 흘린다
    subject.OnNext("A");
    subject.OnNext("B");
    
    // 스트림에 값을 흘린 후 Subscribe 한다.
    stringObservable.Subscribe(Debug.Log);
    
    // Subscribe 후 스트림에 값을 흘린다.
    subject.OnNext("C");
    
    // 완료
    subject.OnCompleted();

    실행 결과

    ABC

    Publish 라는 Hot 변환 연산자를 사이에 끼우는 것으로, Subscribe하는 이전에 스트림을 강제로 실행시킬 수 있습니다.

    각각의 Observer에 대해 별도의 처리를 한다 (스트림의 분기점이 되지 않는다.)​

    Cold Observable은 스트림을 분기시키는 성질을 가지고 있지 않습니다.

    따라서 Cold Observable을 여러 Subscribe하는 경우 각각 별도의 스트림이 생성되고 할당 될 것입니다.


    그러나 스트림에 Hot Observable이 존재하는 경우 가장 말단에 가까운 Hot Observable로 스트림이 분기되어, 또 다른 별도의 스트림이 생성됩니다.

    정리​

    스트림이 어디에서 분기 하는가? 항상 의식하고 설계하는 것이 중요합니다.

    “클래스 외에 Observable을 공개 할 때는 Cold인 채로 공개하지 말고, 반드시 말단에서 Hot으로 변환한 후 공개한다”등 의도하지 않은 곳에서 스트림이 분기되어 버리지 않도록 확실하게 제어해 줄 필요가 있습니다.

    또한 Cold → Hot 변환에는 적용 오퍼레이터가 준비되어 있으므로, 그 쪽을 이용하면 좋을 것 같습니다.

    (Publish, PublishLast, Multicast 등)

  • 洗衣服

    一对年轻的夫妇对面搬来一户新邻居。

    한 젊은 부부가 길 건너편에 새 이웃으로 이사했습니다.

    第二天早上,当他们吃早饭的时候,年轻的妻子看到了新搬来的邻居正在外面洗衣服。

    다음 날 아침, 아침 식사를 하던 중 젊은 아내는 새 이웃이 밖에서 빨래를 하는 모습을 보았습니다.

    妻子对丈夫说道:“那些衣服洗得不干净,也许那个邻居不知道如何清洗。也许她需要好一点的洗衣粉。”

    아내는 남편에게 “저 옷들은 잘 안 빨리는데, 저 이웃은 세탁 방법을 모르는 것 같아요. 더 좋은 세탁 세제를 써야 할 것 같아요.”

    丈夫看了看了妻子,沉默不语。

    남편은 아내를 바라보며 침묵했습니다.

    就这样每次邻居洗衣服,妻子都会这样评论对方一番。 

    그 후로 아내는 이웃이 옷을 세탁할 때마다 이런 식으로 상대방에 대해 언급했습니다.

    大概一个月后,年轻的妻子惊奇地发现,邻居的晾衣绳上居然悬挂着一件干净的衣服,她大叫着对丈夫说:“快看!她学会洗衣服了。我想知道是谁教会她这个的呢?”

    약 한 달 후, 젊은 아내는 이웃집 빨랫줄에 깨끗한 셔츠가 걸려 있는 것을 보고 깜짝 놀라며 남편에게 이렇게 외쳤습니다. 아내는 빨래하는 법을 배웠어요. 누가 이런 걸 가르쳐줬을까?”라고 물었습니다.

    她的丈夫却回答到:“我今天早上一大早起来,然后我把玻璃悬擦干净了。”

    그러자 남편은 “오늘 아침 일찍 일어나서 유리 서스펜션을 닦았어요”라고 대답했습니다.

    在我们作出判断之前,首先要看一下你的“窗户”是否干净。

    판단을 내리기 전에 먼저 ‘창문’이 깨끗한지 확인해야 합니다.