API 서버:
사용자 요청 처리
빠른 응답 반환
Worker:
Queue에 쌓인 작업 처리
외부 API 호출
파일 생성
재시도 처리
사용자 응답 속도 개선
실패 재시도 가능
서버 부하 분산
운영 추적
상담 신청 저장
↓
API 서버는 즉시 성공 응답
↓
Queue에 알림톡 발송 Job 등록
↓
Worker가 알림톡 발송 처리
↓
성공/실패 이력 저장
Producer:
작업을 Queue에 넣는 쪽
Queue:
작업 대기열
Worker:
작업을 꺼내 처리하는 쪽
상담 신청 API
↓
notificationQueue.add('send-alimtalk', { consultId: 123 })
↓
Redis Queue에 Job 저장
↓
Worker가 Job 수신
↓
알림톡 발송
↓
성공/실패 기록
| 기능 | 설명 |
|---|---|
| Queue | 백그라운드 Job 관리 |
| Cache | 상품/배너/FAQ 등 빠른 조회 |
| Rate Limit | 과도한 요청 제한 |
| Session | 로그인 세션 저장 |
| Lock | 배치 중복 실행 방지 |
| Temporary Data | 인증번호, 임시 토큰 저장 |
API 서버가 Job 등록
↓
Redis에 Job 저장
↓
Worker가 Redis에서 Job 조회
↓
처리 결과와 재시도 상태 저장
API 서버가 사용자 요청 처리
동시에 Worker가 대량 엑셀 생성
↓
CPU/메모리 사용량 증가
↓
API 응답 느려짐
↓
사용자 경험 악화
API Process:
HTTP 요청 처리
Worker Process:
Queue Job 처리
Redis:
API와 Worker 사이의 Job 저장소
상담 신청 완료
↓
Queue에 알림톡 Job 등록
↓
Worker가 발송
↓
NotificationLog 저장
관리자 엑셀 다운로드 요청
↓
ExportJob 생성
↓
Queue에 Export Job 등록
↓
Worker가 파일 생성
↓
S3 업로드
↓
ExportJob COMPLETED
Webhook 수신
↓
핵심 상태 변경
↓
Queue에 후속 CRM 연동 Job 등록
↓
Worker가 CRM API 호출
실패 알림 재발송
실패 Webhook 재처리
실패 Export 재생성
외부 API timeout 재시도
EP 파일 재생성
Worker는 실패한 작업을 재시도하고 이력을 관리하는 데 적합합니다.
| 작업 | 이유 |
|---|---|
| 로그인 검증 | 즉시 성공/실패가 필요 |
| 주문 생성 핵심 저장 | 저장 성공 여부가 바로 필요 |
| 결제 승인 결과 확인 | 사용자가 결과를 즉시 알아야 함 |
| 중복 신청 검증 | 저장 전에 막아야 함 |
| 권한 검사 | API 처리 전에 즉시 필요 |
| 필수 DB 트랜잭션 | 데이터 정합성에 직접 영향 |
상담 신청 저장 자체:
API에서 동기 처리
신청 완료 알림톡:
Worker에서 비동기 처리
@nestjs/bullmq를 사용해 Queue와 Worker를 구성할 수 있습니다.npm install @nestjs/bullmq bullmq ioredis
import { BullModule } from '@nestjs/bullmq';
@Module({
imports: [
BullModule.forRoot({
connection: {
host: process.env.REDIS_HOST,
port: Number(process.env.REDIS_PORT),
},
}),
],
})
export class AppModule {}
@Module({
imports: [
BullModule.registerQueue({
name: 'notification',
}),
],
providers: [NotificationProducer, NotificationWorker],
})
export class NotificationQueueModule {}
notification, export, webhook, external-api, batch@Injectable()
export class NotificationProducer {
constructor(
@InjectQueue('notification')
private readonly notificationQueue: Queue,
) {}
async addConsultCreatedJob(consultId: number) {
await this.notificationQueue.add(
'send-consult-created-alimtalk',
{
consultId,
},
{
attempts: 3,
backoff: {
type: 'exponential',
delay: 3000,
},
removeOnComplete: true,
removeOnFail: false,
},
);
}
}
async createConsult(dto: CreateConsultDto) {
const consult = await this.prisma.consult.create({
data: {
name: dto.name,
phone: dto.phone,
status: 'PENDING',
},
});
await this.notificationProducer.addConsultCreatedJob(consult.id);
return consult;
}
@Processor('notification')
export class NotificationWorker extends WorkerHost {
constructor(
private readonly notificationService: NotificationService,
) {
super();
}
async process(job: Job) {
switch (job.name) {
case 'send-consult-created-alimtalk':
return this.notificationService.sendConsultCreatedAlimtalk(
job.data.consultId,
);
default:
throw new Error(`Unknown job name: ${job.name}`);
}
}
}
async sendConsultCreatedAlimtalk(consultId: number) {
const consult = await this.prisma.consult.findUnique({
where: {
id: consultId,
},
});
if (!consult) {
throw new Error('상담 신청 정보를 찾을 수 없습니다.');
}
const alreadySent = await this.prisma.notificationLog.findFirst({
where: {
relatedType: 'CONSULT',
relatedId: consultId,
templateCode: 'CONSULT_CREATED',
status: 'SUCCESS',
},
});
if (alreadySent) {
return {
skipped: true,
reason: 'already_sent',
};
}
const result = await this.alimtalkClient.send({
phone: consult.phone,
templateCode: 'CONSULT_CREATED',
});
await this.prisma.notificationLog.create({
data: {
relatedType: 'CONSULT',
relatedId: consultId,
templateCode: 'CONSULT_CREATED',
status: result.success ? 'SUCCESS' : 'FAILED',
providerMessageId: result.providerMessageId,
errorCode: result.errorCode,
errorMessage: result.errorMessage,
},
});
if (!result.success) {
throw new Error(result.errorMessage ?? '알림톡 발송 실패');
}
return result;
}
| 상태 | 의미 |
|---|---|
| waiting | 처리 대기 |
| active | 처리 중 |
| completed | 처리 완료 |
| failed | 처리 실패 |
| delayed | 지연 후 실행 대기 |
| paused | Queue 일시 정지 |
waiting
↓
active
↓
completed
또는
waiting
↓
active
↓
failed
↓
delayed
↓
active
↓
completed
TIMEOUT
NETWORK_ERROR
PROVIDER_500
PROVIDER_502
PROVIDER_503
RATE_LIMIT
INVALID_PHONE
INVALID_TEMPLATE
AUTH_FAILED
PERMISSION_DENIED
INVALID_PAYLOAD
NOT_FOUND_REQUIRED_DATA
await queue.add(
'send-alimtalk',
{ consultId },
{
attempts: 3,
backoff: {
type: 'exponential',
delay: 5000,
},
},
);
알림톡 발송 성공
↓
성공 로그 저장 전 Worker 종료
↓
Job 재시도
↓
알림톡 중복 발송 가능
notification:consult:123:CONSULT_CREATED
export:consults:job:55
webhook:event:payment:evt_1234
model NotificationLog {
id Int @id @default(autoincrement())
idempotencyKey String @unique
relatedType String
relatedId Int
templateCode String
status String
createdAt DateTime @default(now())
}
await this.notificationQueue.add(
'send-consult-created-alimtalk',
{ consultId },
{
jobId: `notification:consult:${consultId}:created`,
attempts: 3,
backoff: {
type: 'exponential',
delay: 3000,
},
},
);
removeOnComplete 정책도 함께 고려해야 합니다.@Processor('notification', {
concurrency: 5,
})
export class NotificationWorker extends WorkerHost {
async process(job: Job) {
// Job 처리
}
}
| 작업 종류 | 추천 방향 |
|---|---|
| 알림톡/SMS | 외부 API rate limit 고려 |
| 엑셀 생성 | 낮은 concurrency 권장 |
| S3 업로드 | 중간 수준 가능 |
| Webhook 후속 처리 | 중요도에 따라 조절 |
| 광고 전환 API | 외부 API 제한 고려 |
엑셀 생성 Worker concurrency=10
↓
대용량 파일 10개 동시 생성
↓
메모리 급증
↓
서버 불안정
새 API 서버 배포
↓
새로운 jobName으로 Queue 등록
↓
Worker는 아직 구버전
↓
Unknown job name 에러 발생
1차 배포:
Worker가 oldJobName + newJobName 모두 처리
2차 배포:
API가 newJobName 등록
3차 배포:
oldJobName 제거
상담 신청은 되는데 알림톡이 안 감
엑셀 다운로드가 계속 처리 중
Webhook 후속 처리가 안 됨
외부 API 재시도가 멈춤
Queue waiting 수가 계속 증가
pm2 list
pm2 logs worker --lines 100
docker compose logs -f worker
module.exports = {
apps: [
{
name: 'togethermall-api',
script: 'dist/main.js',
instances: 2,
exec_mode: 'cluster',
},
{
name: 'togethermall-worker',
script: 'dist/worker.js',
instances: 1,
},
],
};
services:
api:
image: togethermall-api:20260708
command: node dist/main.js
ports:
- "3000:3000"
env_file:
- .env.production
depends_on:
- redis
worker:
image: togethermall-api:20260708
command: node dist/worker.js
env_file:
- .env.production
depends_on:
- redis
redis:
image: redis:7-alpine
restart: always
async function bootstrap() {
const app = await NestFactory.create(AppModule);
app.enableCors();
await app.listen(process.env.PORT ?? 3000);
}
bootstrap();
async function bootstrap() {
await NestFactory.createApplicationContext(WorkerModule);
}
bootstrap();
main.ts는 HTTP 서버를 띄웁니다.worker.ts는 HTTP 서버 없이 Queue Worker만 실행합니다.| 지표 | 의미 |
|---|---|
| waiting count | 처리 대기 Job 수 |
| active count | 처리 중 Job 수 |
| completed count | 완료 Job 수 |
| failed count | 실패 Job 수 |
| delayed count | 재시도 대기 Job 수 |
| 처리 시간 | Job 하나 처리에 걸리는 시간 |
| 실패율 | 전체 대비 실패 비율 |
waiting count가 계속 증가
failed count가 급증
active Job이 오래 멈춤
delayed Job이 계속 쌓임
Worker 로그에 Redis 연결 오류
외부 API timeout 증가
작업 종류
연관 데이터 ID
상태
실패 사유
시도 횟수
마지막 시도 시간
재시도 가능 여부
수동 재처리 버튼
model JobFailureLog {
id Int @id @default(autoincrement())
queueName String
jobName String
jobId String?
relatedType String?
relatedId Int?
status String
attempts Int @default(0)
errorCode String?
errorMessage String?
retryable Boolean @default(false)
createdAt DateTime @default(now())
updatedAt DateTime @updatedAt
@@index([queueName, status])
@@index([relatedType, relatedId])
}
{
"name": "홍길동",
"phone": "01012345678",
"memo": "상담 메모 전체",
"address": "서울시 ..."
}
{
"consultId": 123
}
Unknown job name이 없는가?NestJS + Prisma + Redis + BullMQ로 Queue Worker 구조를 설계하고 싶어.
상황:
1. 상담 신청 저장 후 알림톡을 발송해야 함
2. 알림톡 실패해도 상담 신청 저장은 성공이어야 함
3. 엑셀 다운로드는 대용량이라 Worker에서 생성 후 S3에 업로드하고 싶음
4. Webhook 수신 후 CRM 연동도 Worker로 분리하고 싶음
5. 같은 알림톡이 중복 발송되면 안 됨
6. 실패한 Job은 최대 3번 재시도하고 싶음
7. 실패 Job은 관리자 페이지에서 확인하고 수동 재처리하고 싶음
8. API 서버와 Worker를 Docker Compose에서 분리 실행할 예정
9. Job data에는 개인정보를 최소화하고 싶음
요청:
- Queue 종류 설계
- Producer/Worker 구조
- BullMQ 설정
- 재시도/backoff 기준
- 멱등성 처리
- NotificationLog/ExportJob/JobFailureLog 모델
- API 서버와 Worker 분리 실행 방식
- 배포 순서와 운영 체크리스트
를 실무 기준으로 정리해줘.
consultId, exportJobId처럼 식별자만 넣고 Worker가 DB에서 필요한 정보를 조회하는 방식이 안전합니다.