# [Message Queue - 1부] Nestjs Message Queue & Redis

# 💁🏻‍♂️ Subject

* Message Queue에 대하여 알아보자.
    

# 🎯 Goal for Research

* \[x\] Message Queue의 개념을 이해.
    
* \[x\] Nestjs에서 Message Queue를 셋팅할 수 있다.
    
* \[x\] Message Queue에 Message 를 추가할 수 있다.
    
* \[x\] Message Queue에서 Message를 읽어 사용할 수 있다.
    

# 📃 Contents

### 개념

* 프로그램 간에 데이터를 교환할 때 사용하는 통신 방법
    
* 메시지 지향 미들웨어(Message Oriented Middleware:MOM)를 구현한 시스템
    
* Message Queue에는 3개의 이해관계자가 존재
    
    * Message Queue : Message 임시 저장 버퍼로 Consumer가 읽을 때까지 보관.
        
    * Producer : Message 생성자로 Message를 만들어 MQ에 삽입하는 역할.
        
    * Consumer : Message Queue에 있는 메세지를 읽고 처리하는 서버.
        
* Producer가 생성한 Message를 Message Queue는 Message 임시 저장 Buffer 역할을 수행하고, 이를 Consumer가 가져가서 실행하게된다.
    

### 사용이유

**안정성 향상**

* 트래픽 과다 현상으로 병목이 발생할 경우 Data Loss가 발생할 수 있는데 Message Queue에 저장함으로써 서버가 죽어도 재처리 가능.
    

**비동기처리**

* Producer & Comsumer간 Request -&gt; 대기 -&gt; Response를 기다리지 않고 독립적으로 작동할 수 있도록 함. 즉, 처리 속도가 다른 프로세스 간 병목 현상 방지.
    
* 이는 서비스간 의존성을 낮춤.
    

**서비스 확장성**

* Message Queue에 등록된 데이터는 여러 Consumer가 사용할 수 있다. 이를 통해 부하분산 효과가 있으며 트래픽 증가에도 안정적 동작이 가능하게 함.
    

### 사용 예시

**채팅 시스템**

* 사용자가 메세지 전송 시, Socket Emit이 되고 On한 쪽에서 리스닝하게되는데 이러한 메세지를 메세지 큐에 저장하고, On한 쪽에서 수신이 완료되면 메세지큐에서 삭제함으로써 메세지 손실을 방지.
    

**실시간 알림 서비스**

* 사용자가 특정 이벤트(쪽지 등)에 대하여 알림 구독 시, 메세지 큐에 정보를 저장 후 이벤트가 발생(ex&gt; 사용자 로그인)하면 메세지 큐를 통해 구독자에게 알림 전송.
    

**Web3 이벤트 수신**

* Smart Contract에서 발생하는 이벤트를 Message Queue에 저장하고, xxx이벤트가 발생했을 때 적절히 처리하는 서버(Consumer) 가 잡아서 처리를 했을 때 안정성 높게 Event 수신이 가능하다.
    

### Message Queue 기본 사용 옵션

**Visibility Timeout (제한시간 초과)**

* Consumer가 Message Queue 데이터를 읽은 뒤, 즉시 데이터를 삭제하지 않고 남겨두는 기간.
    
* Consumer가 Message 처리에 다소 시간이 소요되는 경우 유용.
    
* 장점 :
    
    * Consumer가 메세지 처리 실패 시, 다른 Consumer가 메세지 처리할 수 있음.
        
    * Message Data Loss 방지
        
* 단점 :
    
    * Message Queue 용량 소모량 증대
        

**Delay Queue (지연 큐)**

* 메시지가 큐에 저장된 후 특정 시간이 지나야 Consumer가 읽을 수 있는 큐.
    
* 메시지 전송을 일정 시간 지연시키거나, 특정 순서대로 메시지를 처리해야 하는 경우 유용.
    
* 장점 :
    
    * 부하 분산
        
    * 메세지 순서대로 처리 할 시 유용함.
        
* 단점 :
    
    * 처리순서가 중요하지 않은 경우 오버헤드로 작용.
        

**Exclusive Queue (독점 큐)**

* 한 번에 하나의 Consumer만 큐에 접근할 수 있도록 하는 큐
    
* 장점 :
    
    * 데이터 충돌 방지
        
* 단점 :
    
    * 비동기 처리 효율성 저하
        

### Message Queue 구현 데모

* Redis Docker로 실행
    

```plaintext
docker pull redis
sudo docker run -p 6379:6379 redis
```

* bull 설치
    

```bash
npm install --save @nestjs/bull bull
```

* 메세지를 등록하는 producer용 API를 생성합니다.
    
* app.module에 Bull Module 등록
    

```js
import { Module } from '@nestjs/common';
import { AppController } from './app.controller';
import { AppService } from './app.service';

// Nestjs Bull
import { BullModule } from '@nestjs/bull';
import { AppConsumer } from './consumer.service';

  

@Module({
	imports: [
		BullModule.forRoot({
		// Redis Host 설정 옵션
			redis: {
				host: '127.0.1',
				port: 6379,
			},
			limiter : {
				max : 3, // 초당 최대 1개의 작업을 허용
				duration : 1 // 1초동안 제한 적용
			}
		}),
	
		BullModule.registerQueue({
			name : "my-queue"
		}),
	],
	controllers: [AppController],
	providers: [AppService, AppConsumer],
})

export class AppModule {}
```

* Controller에 API 생성
    

```js
export class AppController {
	constructor(private readonly appService: AppService) {}
	@Get()
	addMessage(@Body() data : any) {
		return this.appService.addMessageQueue(data);
	}
}
```

* Service 생성
    

```js
export class AppService {
	constructor(
		@InjectQueue("my-queue") private testQueue : Queue
	){}

async addMessageQueue(data : any) {
	const job = await this.testQueue.add('my-queue', data, {	
		delay : 1500,
		// priority : 우선순위 설정
	});

	return job;
}
```

위 API를 통해 Message를 등록하면 Redis에 아래처럼 생성되는 것을 확인할 수 있다.

![](https://cdn.hashnode.com/res/hashnode/image/upload/v1709127256810/3ac4e834-3b21-478a-b37b-af1b870b2f15.png align="center")

* 이제 메세지를 소비할 Consumer를 생성한다.
    

```js
// consumer.service.ts
import { Injectable } from '@nestjs/common';
import { Process, Processor } from '@nestjs/bull';
import { Job } from 'bull';

@Processor('my-queue')
export class AppConsumer {
	@Process('my-queue')
	async processData(job: Job) {
		console.log('id', job.id);
		console.log('data :', job.data);
		
		// 완료처리를 진행함. 정상처리되었을 때 결과값을 아래에 넣어줘야함.
		job.moveToCompleted("잘처리 ✅");
	}
}
```

* 처리가 완료되었을 때, 결과값이 아래와 같이 업데이트 되는 것을 확인 할 수 있다.
    

![](https://cdn.hashnode.com/res/hashnode/image/upload/v1709126448388/ce24cff9-dbb2-4f43-ae98-1f2ffbc7185c.png align="center")
