API 요청
↓
DB에 작업 생성
↓
빠르게 응답
↓
Worker가 뒤에서 작업 처리
↓
작업 결과 DB에 저장
API 응답이 느려짐
timeout 발생 가능
외부 API 장애가 고객 요청에 직접 영향
엑셀 생성 중 서버 메모리 사용 증가
작업 실패 시 재시도 어려움
작업 진행 상태를 관리자에게 보여주기 어려움
Producer:
작업 생성
Queue:
작업 대기
Worker:
작업 실행
Job Table:
작업 상태 기록
| 역할 | 설명 |
|---|---|
| Producer | 작업을 생성하는 API Use Case |
| Queue | 작업을 대기시키는 공간 |
| Worker | 작업을 처리하는 실행자 |
| Job | 처리해야 할 작업 단위 |
| Retry | 실패한 작업 재시도 |
| Dead Letter | 반복 실패한 작업 격리 |
상담 접수 알림톡 발송
상태 변경 알림톡 발송
관리자 엑셀 Export
네이버/당근 EP 생성
대량 상품 이미지 처리
유입 분석 집계
외부 Webhook 재시도
요청
↓
상담 저장
↓
알림톡 API 호출
↓
알림톡 응답 대기
↓
응답 반환
장점:
흐름이 단순함
결과를 바로 알 수 있음
초기 구현이 쉬움
단점:
외부 API가 느리면 전체 요청이 느림
외부 API 실패가 원본 저장 실패로 이어질 수 있음
재시도 구조가 약함
요청
↓
상담 저장
↓
notification_job 생성
↓
응답 반환
↓
Worker가 알림톡 발송
↓
결과 저장
장점:
API 응답이 빨라짐
외부 API 장애와 원본 저장 분리
재시도 가능
작업 상태 추적 가능
단점:
구조가 조금 복잡해짐
Worker 운영 필요
작업 상태 관리 필요
처리 시간이 오래 걸림
외부 API를 호출함
실패 시 재시도가 필요함
작업 진행 상태를 보여줘야 함
대량 데이터를 처리함
파일 생성/업로드가 필요함
API timeout 가능성이 있음
엑셀 파일 생성
알림톡/SMS 발송
대량 상품 업데이트
외부 API 동기화
EP 파일 생성
통계 집계
S3 파일 정리
만료된 Export 파일 삭제
단순 상담 저장
단순 목록 조회
단순 상태 변경
단순 메모 저장
짧은 validation
권한 확인
jobs 또는 도메인별 Job 테이블을 만들고, Worker가 주기적으로 처리하는 방식으로 시작할 수 있습니다.notification_jobs
export_jobs
ep_generation_jobs
cleanup_jobs
구현이 비교적 단순함
현재 PostgreSQL만으로 시작 가능
관리자 화면에서 작업 상태 확인 쉬움
트랜잭션과 함께 Job 생성 가능
작업 이력이 DB에 남음
처리량이 많아지면 DB 부하
Worker 간 동시 처리 제어 필요
실시간성이 Redis Queue보다 약함
polling 구조가 필요할 수 있음
PENDING
↓
PROCESSING
↓
DONE
PROCESSING
↓
FAILED
FAILED
↓
RETRYING
↓
PROCESSING
| 상태 | 의미 |
|---|---|
PENDING | 처리 대기 |
PROCESSING | 처리 중 |
DONE | 처리 완료 |
FAILED | 처리 실패 |
RETRYING | 재시도 대기 |
CANCELED | 취소됨 |
EXPIRED | 만료됨 |
작업 흐름을 추적할 수 있어야 함
실패 사유를 저장해야 함
재시도 횟수를 저장해야 함
처리 시작/완료 시간을 저장해야 함
Worker가 죽어도 복구 가능해야 함
limit=100000으로 한 번에 내려주는 방식은 위험합니다.관리자 엑셀 다운로드 요청
↓
export_jobs row 생성
↓
API는 jobId 반환
↓
Worker가 파일 생성
↓
S3 업로드
↓
export_jobs.status = DONE
↓
관리자 다운로드
id
type
status
requestedByAdminId
filterSnapshot
fileKey
fileUrl
errorCode
errorMessage
retryCount
requestedAt
startedAt
finishedAt
expiresAt
createdAt
updatedAt
model ExportJob {
id Int @id @default(autoincrement())
type ExportJobType
status JobStatus @default(PENDING)
requestedByAdminId Int
requestedByAdmin AdminUser @relation(fields: [requestedByAdminId], references: [id])
filterSnapshot Json
fileKey String?
fileUrl String?
errorCode String?
errorMessage String?
retryCount Int @default(0)
requestedAt DateTime @default(now())
startedAt DateTime?
finishedAt DateTime?
expiresAt DateTime?
createdAt DateTime @default(now())
updatedAt DateTime @updatedAt
@@index([status, createdAt])
@@index([requestedByAdminId, createdAt])
}
enum ExportJobType {
CONSULT_LIST
ORDER_LIST
AUDIT_LOG
}
enum JobStatus {
PENDING
PROCESSING
DONE
FAILED
RETRYING
CANCELED
EXPIRED
}
검색 조건을 filterSnapshot으로 저장
요청 관리자 ID 저장
상태와 실패 사유 저장
파일 만료 시간 저장
status + createdAt 인덱스 필요
권한 확인
↓
검색 조건 검증
↓
filterSnapshot 생성
↓
export_jobs 생성
↓
audit_logs 생성
↓
jobId 반환
@Injectable()
export class RequestConsultExportUseCase {
constructor(
private readonly prisma: PrismaService,
private readonly exportJobRepository: ExportJobRepository,
private readonly auditLogRepository: AuditLogRepository,
private readonly exportPermissionPolicy: ExportPermissionPolicy,
) {}
async execute(command: RequestConsultExportCommand) {
if (!this.exportPermissionPolicy.canExportConsults(command.adminRole)) {
throw new ForbiddenException('엑셀 다운로드 권한이 없습니다.');
}
return this.prisma.$transaction(async (tx) => {
const job = await this.exportJobRepository.create(
{
type: 'CONSULT_LIST',
status: 'PENDING',
requestedByAdminId: command.adminId,
filterSnapshot: command.filter,
expiresAt: addDays(new Date(), 3),
},
tx,
);
await this.auditLogRepository.create(
{
actorType: 'ADMIN',
actorId: command.adminId,
action: 'EXPORT_DOWNLOAD_REQUEST',
targetType: 'EXPORT_JOB',
targetId: String(job.id),
afterValue: {
type: 'CONSULT_LIST',
},
requestId: command.requestId,
ipAddress: command.ipAddress,
userAgent: command.userAgent,
},
tx,
);
return {
jobId: job.id,
status: job.status,
};
});
}
}
API에서 엑셀 파일을 바로 만들지 않기
검색 조건을 그대로 문자열 로그에 남기지 않기
개인정보 필터값은 Audit Log에 최소화
권한 없는 관리자 다운로드 차단
PENDING job 조회
↓
PROCESSING으로 변경
↓
filterSnapshot 기준 데이터 조회
↓
엑셀 파일 생성
↓
S3 업로드
↓
DONE 처리
async processNextExportJob() {
const job = await this.exportJobRepository.claimNextPendingJob();
if (!job) {
return;
}
try {
const rows = await this.consultRepository.findForExport(
job.filterSnapshot,
);
const file = await this.excelService.createConsultExcel(rows);
const uploaded = await this.s3Adapter.upload({
key: `exports/consults/${job.id}.xlsx`,
body: file,
});
await this.exportJobRepository.markDone({
jobId: job.id,
fileKey: uploaded.key,
fileUrl: uploaded.url,
finishedAt: new Date(),
});
} catch (error) {
await this.exportJobRepository.markFailed({
jobId: job.id,
errorCode: 'EXPORT_FAILED',
errorMessage: getSafeErrorMessage(error),
});
}
}
Job 선점 처리
중복 처리 방지
실패 사유 저장
재시도 횟수 관리
파일 생성 중 메모리 사용 주의
S3 업로드 실패 처리
개인정보 파일 보안
PENDING 상태를 PROCESSING으로 바꾸는 과정을 안전하게 처리해야 합니다.Worker A:
PENDING job 조회
Worker B:
같은 PENDING job 조회
문제:
둘 다 같은 파일 생성 가능
const result = await prisma.exportJob.updateMany({
where: {
id: jobId,
status: 'PENDING',
},
data: {
status: 'PROCESSING',
startedAt: new Date(),
},
});
if (result.count === 0) {
return null;
}
FOR UPDATE SKIP LOCKEDSELECT id
FROM export_jobs
WHERE status = 'PENDING'
ORDER BY created_at ASC
LIMIT 1
FOR UPDATE SKIP LOCKED;
초기:
status 조건 update로 시작 가능
Worker가 여러 개:
FOR UPDATE SKIP LOCKED 고려
처리량 증가:
Redis Queue/BullMQ 고려
상담 신청 저장
↓
notification_jobs 생성
↓
API 응답
↓
Worker가 알림톡 발송
↓
notification_logs 업데이트
id
type
status
targetType
targetId
channel
templateCode
recipientMasked
recipientEncrypted
payload
providerMessageId
errorCode
errorMessage
retryCount
scheduledAt
startedAt
sentAt
createdAt
updatedAt
model NotificationJob {
id Int @id @default(autoincrement())
type String
status JobStatus @default(PENDING)
targetType String
targetId String
channel NotificationChannel
templateCode String
recipientMasked String?
recipientEncrypted String?
payload Json
providerMessageId String?
errorCode String?
errorMessage String?
retryCount Int @default(0)
scheduledAt DateTime?
startedAt DateTime?
sentAt DateTime?
createdAt DateTime @default(now())
updatedAt DateTime @updatedAt
@@index([status, scheduledAt])
@@index([targetType, targetId, createdAt])
}
enum NotificationChannel {
ALIMTALK
SMS
EMAIL
}
전화번호 원본 저장 주의
recipientEncrypted 또는 별도 보안 기준 필요
recipientMasked는 관리자 화면 표시용
payload에 개인정보 과다 저장 금지
외부 API Secret 저장 금지
Transaction 시작
↓
consults 생성
↓
notification_jobs 생성
↓
commit
↓
Worker가 발송
await this.prisma.$transaction(async (tx) => {
const consult = await this.consultRepository.createWithSnapshot(
{
customerName: command.customerName,
phoneMasked: command.phoneMasked,
phoneEncrypted: command.phoneEncrypted,
productId: product.id,
productNameSnapshot: product.modelName,
carrierSnapshot: product.carrier,
source: command.source,
visitorId: command.visitorId,
},
tx,
);
await this.notificationJobRepository.create(
{
type: 'CONSULT_RECEIVED',
targetType: 'CONSULT',
targetId: String(consult.id),
channel: 'ALIMTALK',
templateCode: 'CONSULT_RECEIVED',
recipientMasked: command.phoneMasked,
recipientEncrypted: command.phoneEncrypted,
payload: {
consultId: consult.id,
productName: product.modelName,
},
},
tx,
);
return consult;
});
상담 저장 실패 시 Job도 생성되지 않음
Job 생성 실패 시 상담 저장도 rollback 여부 결정
실제 발송은 transaction 밖에서 Worker가 처리
payload에는 필요한 값만 저장
PENDING 상태의 알림 Job을 가져와 외부 알림톡/SMS API를 호출합니다.DONE 또는 SENT, 실패하면 FAILED 또는 RETRYING으로 바꿉니다.PENDING notification job 조회
↓
PROCESSING으로 변경
↓
recipient 복호화
↓
알림톡 API 호출
↓
성공/실패 결과 저장
async processNextNotificationJob() {
const job = await this.notificationJobRepository.claimNextPendingJob();
if (!job) {
return;
}
try {
const recipient = this.cryptoService.decrypt(job.recipientEncrypted);
const result = await this.alimtalkAdapter.send({
phone: recipient,
templateCode: job.templateCode,
payload: job.payload,
});
await this.notificationJobRepository.markSent({
jobId: job.id,
providerMessageId: result.messageId,
sentAt: new Date(),
});
} catch (error) {
await this.notificationJobRepository.markFailedOrRetry({
jobId: job.id,
errorCode: getProviderErrorCode(error),
errorMessage: getSafeErrorMessage(error),
});
}
}
수신 전화번호 로그 출력 금지
외부 API 응답에 개인정보 포함 여부 확인
실패 메시지에 Secret 포함되지 않게 필터링
재시도 횟수 제한
동일 알림 중복 발송 방지
1차 실패
↓
retryCount 증가
↓
일정 시간 후 재시도
↓
최대 횟수 초과 시 FAILED
외부 API 일시 장애
네트워크 timeout
rate limit
일시적인 서버 오류
잘못된 전화번호
잘못된 템플릿 코드
권한/인증 정보 오류
필수 변수 누락
수신 거부
최대 재시도:
3회
간격:
1분 → 5분 → 15분
최종 실패:
FAILED 상태 저장
관리자 화면에서 실패 사유 표시
재시도할수록 간격을 늘림
외부 API 장애 중 무리한 반복 호출 방지
Worker가 알림톡 발송 성공
↓
DB markSent 전에 서버 죽음
↓
Job은 PROCESSING 또는 PENDING 상태
↓
재시도 시 같은 알림 다시 발송 가능
provider idempotency key 사용 가능 여부 확인
Job ID를 외부 요청 고유키로 사용
발송 전 중복 여부 확인
targetType + targetId + type unique 고려
상태 transition 조건 update 사용
CREATE UNIQUE INDEX uniq_notification_target_type
ON notification_jobs (target_type, target_id, type)
WHERE status IN ('PENDING', 'PROCESSING', 'DONE');
같은 상담에 같은 알림을 여러 번 보내야 하는 경우도 있음
재발송 기능과 자동 중복 방지를 구분
수동 재발송은 별도 type 또는 resendGroup 사용 고려
PROCESSING 상태에 영원히 남을 수 있습니다.Job status = PROCESSING
startedAt = 3시간 전
Worker는 이미 종료됨
결과:
작업이 다시 처리되지 않음
PROCESSING 상태가 일정 시간 이상 지속되면 감지
retryCount 확인
RETRYING 또는 FAILED로 변경
관리자 화면에 표시
await prisma.exportJob.updateMany({
where: {
status: 'PROCESSING',
startedAt: {
lt: subMinutes(new Date(), 30),
},
},
data: {
status: 'RETRYING',
errorCode: 'JOB_TIMEOUT',
errorMessage: '작업 시간이 초과되어 재시도 대기 상태로 변경되었습니다.',
retryCount: {
increment: 1,
},
},
});
긴 작업과 멈춘 작업을 구분해야 함
작업별 timeout 기준 다르게 설정
엑셀 대량 생성은 더 긴 timeout 가능
알림톡 발송은 짧은 timeout 가능
API 서버 안에서 setInterval 또는 Schedule 사용
장점:
구현이 쉬움
별도 배포가 필요 없음
초기 단계에 적합
단점:
API 서버와 Worker 부하가 섞임
서버 여러 대일 때 중복 처리 위험
장기 작업이 API 서버에 영향
API 서버:
요청 처리
Worker 서버:
Job 처리
장점:
API 부하와 작업 부하 분리
Worker만 재시작 가능
작업 종류별 확장 가능
단점:
배포/운영 복잡도 증가
로그/모니터링 추가 필요
초기:
NestJS Schedule 또는 간단한 Worker 프로세스
작업량 증가:
별도 Worker 프로세스 분리
처리량 증가:
Redis Queue/BullMQ 도입 검토
API
↓
BullMQ Queue
↓
Redis
↓
Worker
Queue 기능이 풍부함
retry/backoff 지원
delayed job 지원
concurrency 설정 가능
job progress 관리 가능
Redis 운영 필요
장애 지점 증가
초기 구조 복잡도 증가
DB Job과 상태 중복 관리 가능성
DB Job으로 충분:
작업량 적음
관리자 Export 위주
알림 발송량 적음
BullMQ 고려:
작업량 많음
재시도/예약/동시성 제어가 중요
여러 Worker 확장이 필요
엑셀 다운로드 요청
↓
작업 목록에 PENDING 표시
↓
PROCESSING 표시
↓
DONE이면 다운로드 버튼 활성화
↓
FAILED면 실패 사유 표시
작업 ID
작업 유형
상태
요청 관리자
요청 시간
완료 시간
만료 시간
다운로드 버튼
실패 사유
발송 대상
채널
템플릿
상태
요청 시간
발송 시간
실패 사유
재시도 횟수
수동 재시도 버튼
상태를 한글 라벨로 표시
실패 사유를 운영자가 이해할 수 있게 표시
완료된 Export는 다운로드 버튼 제공
만료된 파일은 재요청 안내
수동 재시도는 권한 제한
전화번호/이름 포함
상담 메모 포함
S3 URL 외부 노출
파일 장기 보관
권한 없는 관리자 다운로드
다운로드 이력 없음
다운로드 권한 별도 관리
S3 private bucket 사용
pre-signed URL 사용
URL 만료 시간 설정
Export 파일 1~7일 후 삭제
다운로드 요청 Audit Log 기록
전화번호 마스킹 여부 정책화
fileUrl을 영구 public URL로 저장하지 않기
관리자 화면에서 권한 확인 후 URL 발급
파일 만료 후 S3 object 삭제
ExportJob row는 남기되 fileKey/fileUrl 처리 기준 정하기
Job claim 성공
Job processing 시작
Job done
Job failed
retry count 증가
stuck job 감지
외부 API 실패 코드
S3 업로드 실패
jobId
jobType
status
durationMs
errorCode
requestId 또는 traceId
targetType
targetId
전화번호 원본
고객 이름 원본
상담 메모 전체
외부 API Secret
Authorization header
S3 pre-signed URL 전체
# Runbook: Worker 작업 실패 대응
## 상황
- ExportJob 실패
- NotificationJob 실패
- Job이 PROCESSING 상태로 오래 유지
- Worker 프로세스 중단
## 즉시 확인
- 실패한 jobId:
- job type:
- status:
- retryCount:
- errorCode:
- errorMessage:
- startedAt:
- finishedAt:
- 최근 배포:
- 외부 API 장애 여부:
## 확인 절차
1. Worker 로그 확인
2. Job row 상태 확인
3. 같은 유형의 실패가 여러 건인지 확인
4. 외부 API 상태 확인
5. S3 업로드/권한 문제 확인
6. retry 가능한 오류인지 판단
7. 수동 retry 또는 FAILED 유지 결정
8. 필요 시 Worker 재시작
9. 관리자 화면 안내
## 주의
- 알림 중복 발송 여부 확인
- Export 파일 중복 생성 여부 확인
- 개인정보 로그 노출 여부 확인
- 수동 재시도 권한 확인
이유:
엑셀 다운로드는 데이터량이 커질수록 위험
개인정보 파일 보안 필요
API timeout 방지
관리자 작업 상태 표시 가능
적용 범위:
ExportJob 테이블
RequestConsultExportUseCase
ExportWorker
S3 upload
관리자 작업 목록
파일 만료 처리
이유:
상담 저장과 알림 발송 분리
외부 API 실패로 상담 저장 실패 방지
재시도 가능
발송 이력 추적 가능
적용 범위:
NotificationJob 테이블
CreateConsultUseCase에서 Job 생성
NotificationWorker
AlimtalkAdapter
발송 실패/재시도 관리
이유:
Export 파일 만료 삭제
오래된 임시 데이터 정리
보안/비용 관리
적용 범위:
만료된 Export 파일 조회
S3 object 삭제
ExportJob EXPIRED 처리
작업 결과 로그
dry run 지원
관리자 상담 목록 엑셀 다운로드를 ExportJob + Worker 구조로 분리해줘.
조건:
1. 일반 목록 API에서 대량 엑셀을 바로 생성하지 마
2. RequestConsultExportUseCase를 만들고 ExportJob row만 생성하게 해줘
3. 검색 조건은 filterSnapshot으로 저장해줘
4. requestedByAdminId, status, retryCount, expiresAt을 저장해줘
5. ExportJob 생성과 Audit Log 저장은 하나의 transaction으로 묶어줘
6. ExportWorker는 PENDING job을 PROCESSING으로 claim한 뒤 처리하게 해줘
7. 같은 job이 중복 처리되지 않게 status 조건 update를 사용해줘
8. 엑셀 생성 후 S3 private bucket에 업로드하고 fileKey를 저장해줘
9. 다운로드는 권한 확인 후 pre-signed URL을 발급하는 구조로 해줘
10. 전화번호 원본이나 개인정보가 로그에 남지 않게 해줘
11. 실패 시 errorCode, errorMessage, retryCount를 저장해줘
12. 변경 후 DB migration, 테스트 케이스, QA 체크리스트를 정리해줘
상담 접수 알림톡 발송을 NotificationJob + Worker 구조로 분리해줘.
조건:
1. CreateConsultUseCase에서 알림톡 API를 직접 호출하지 마
2. 상담 저장과 NotificationJob 생성은 같은 transaction으로 묶어줘
3. Worker가 PENDING NotificationJob을 claim해서 알림톡 API를 호출하게 해줘
4. 전화번호 원본은 로그와 응답에 남기지 마
5. recipientMasked와 recipientEncrypted를 구분해줘
6. 실패 시 retryCount를 증가시키고 재시도 가능/불가능 오류를 구분해줘
7. 동일 상담에 같은 알림이 중복 발송되지 않도록 idempotency 기준을 제안해줘
8. providerMessageId, sentAt, errorCode, errorMessage를 저장해줘
9. 외부 API Secret은 DB나 로그에 저장하지 마
10. 테스트 케이스와 Worker 실패 대응 Runbook을 정리해줘
API에서 긴 작업을 바로 처리하지 않는가?
Job 생성과 원본 DB 변경이 필요한 경우 같은 transaction인가?
Worker가 같은 Job을 중복 처리하지 않는가?
실패 사유와 retryCount가 저장되는가?
외부 API 호출이 transaction 안에 있지 않은가?
개인정보가 로그에 남지 않는가?
Export 파일 보안 기준이 적용되는가?
관리자 화면에서 Job 상태를 확인할 수 있는가?
NestJS + Prisma + PostgreSQL 기반 온라인 휴대폰 판매몰에서 Queue/Worker 구조를 설계하려고 해.
서비스 상황:
1. 고객 상담 신청 후 알림톡을 발송해야 함
2. 알림톡 외부 API 실패 때문에 상담 저장이 실패하면 안 됨
3. 관리자는 상담 목록을 검색한 조건 그대로 엑셀 다운로드할 수 있어야 함
4. 엑셀 생성은 오래 걸릴 수 있고 개인정보가 포함될 수 있음
5. Export 파일은 S3 private bucket에 저장하고 만료 후 삭제하고 싶음
6. API 요청에서는 Job만 생성하고 Worker가 실제 작업을 처리하게 하고 싶음
7. Job 상태는 PENDING, PROCESSING, DONE, FAILED, RETRYING, EXPIRED로 관리하려고 함
8. Worker가 같은 Job을 중복 처리하지 않게 해야 함
9. 실패 시 retryCount와 errorMessage를 저장하고 싶음
10. 전화번호 원본, Secret, pre-signed URL은 로그에 남기면 안 됨
요청:
- DB 기반 Job Table로 시작하는 구조
- ExportJob 테이블 설계
- NotificationJob 테이블 설계
- RequestConsultExportUseCase 구조
- CreateConsultUseCase에서 NotificationJob 생성 구조
- Worker claim 로직
- retry/backoff 설계
- 중복 발송 방지 idempotency 기준
- Export 파일 보안 기준
- Worker 로그/모니터링 기준
- 관리자 Job 상태 UI
- 실패 대응 Runbook
을 실무 기준으로 정리해줘.
외부 API 호출을 transaction 밖으로 분리하는가?
상담 저장과 NotificationJob 생성은 같은 transaction으로 묶는가?
엑셀 생성은 API에서 직접 하지 않게 하는가?
Export 파일 개인정보 보안 기준을 다루는가?
Worker 중복 처리 방지 로직을 설명하는가?
retry 가능한 오류와 불가능한 오류를 구분하는가?
전화번호 원본과 Secret 로그 노출을 경고하는가?
BullMQ/Redis를 무조건 도입하라고 하지 않는가?
현재 규모에서는 DB Job부터 시작 가능하다고 보는가?
export_jobs, notification_jobs 테이블로 시작하는 것이 현실적입니다.PENDING, PROCESSING, DONE, FAILED, RETRYING, EXPIRED 같은 상태를 가져야 하며, 실패 사유와 retryCount를 저장해야 합니다.ExportJob을 생성한 뒤 Worker가 검색 조건 snapshot 기준으로 파일을 생성하는 구조가 안전합니다.NotificationJob row만 생성하고, 실제 외부 API 호출은 Worker가 처리해야 합니다.FOR UPDATE SKIP LOCKED, 또는 Queue 시스템의 claim 기능을 사용해야 합니다.