아마존 SQS란? 자바 사용가이드

GarionNachal·2025년 9월 10일

AWS

목록 보기
4/11

개요

Amazon Simple Queue Service(SQS)는 완전 관리형 메시지 큐 서비스로, 마이크로서비스, 분산 시스템 및 서버리스 애플리케이션을 분리하고 확장할 수 있게 해주는 AWS의 핵심 서비스입니다.

SQS란 무엇인가?

Amazon SQS는 애플리케이션 구성 요소 간에 메시지를 안전하게 전송할 수 있는 메시지 큐 서비스입니다. 이를 통해 시스템의 각 구성 요소를 분리하여 독립적으로 확장하고 운영할 수 있습니다.

주요 특징

  • 완전 관리형: 서버 프로비저닝, 패치 관리, 용량 계획이 불필요
  • 높은 가용성: 여러 가용 영역에 걸쳐 중복 저장
  • 무제한 확장: 초당 수백만 개의 메시지 처리 가능
  • 보안: 전송 중 및 저장 중 암호화 지원
  • 비용 효율적: 사용한 만큼만 지불하는 요금제

SQS 큐 유형

1. 표준 큐 (Standard Queue)

  • 무제한 처리량: 초당 거의 무제한의 API 호출 지원
  • 최소 한 번 전달: 메시지가 최소 한 번은 전달됨을 보장
  • 최선 노력 순서: 메시지가 전송된 순서와 동일한 순서로 전달하려고 노력

2. FIFO 큐 (First-In-First-Out Queue)

  • 순서 보장: 메시지가 전송된 정확한 순서로 전달
  • 정확히 한 번 처리: 중복 메시지 없이 정확히 한 번만 전달
  • 높은 처리량: 배치를 사용하면 초당 최대 3,000개의 메시지 처리

마이크로서비스에서의 메시지 큐 패턴

자바에서 SQS 사용하기

1. 의존성 설정

Maven을 사용하는 경우 pom.xml에 다음을 추가:

<dependency>
    <groupId>software.amazon.awssdk</groupId>
    <artifactId>sqs</artifactId>
    <version>2.20.162</version>
</dependency>

Gradle을 사용하는 경우:

implementation 'software.amazon.awssdk:sqs:2.20.162'

2. SQS 클라이언트 설정

import software.amazon.awssdk.regions.Region;
import software.amazon.awssdk.services.sqs.SqsClient;
import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider;

public class SQSClientConfig {
    
    private static SqsClient sqsClient;
    
    public static SqsClient getSqsClient() {
        if (sqsClient == null) {
            sqsClient = SqsClient.builder()
                .region(Region.AP_NORTHEAST_2) // 서울 리전
                .credentialsProvider(DefaultCredentialsProvider.create())
                .build();
        }
        return sqsClient;
    }
}

Java AWS SDK SQS 예제

3. 큐 생성

import software.amazon.awssdk.services.sqs.model.*;

public class QueueManager {
    
    private final SqsClient sqsClient;
    
    public QueueManager() {
        this.sqsClient = SQSClientConfig.getSqsClient();
    }
    
    // 표준 큐 생성
    public String createStandardQueue(String queueName) {
        try {
            CreateQueueRequest request = CreateQueueRequest.builder()
                .queueName(queueName)
                .build();
                
            CreateQueueResponse response = sqsClient.createQueue(request);
            return response.queueUrl();
        } catch (QueueNameExistsException e) {
            // 큐가 이미 존재하는 경우
            GetQueueUrlRequest getUrlRequest = GetQueueUrlRequest.builder()
                .queueName(queueName)
                .build();
            return sqsClient.getQueueUrl(getUrlRequest).queueUrl();
        }
    }
    
    // FIFO 큐 생성
    public String createFIFOQueue(String queueName) {
        Map<String, String> attributes = new HashMap<>();
        attributes.put("FifoQueue", "true");
        attributes.put("ContentBasedDeduplication", "true");
        
        CreateQueueRequest request = CreateQueueRequest.builder()
            .queueName(queueName + ".fifo")
            .attributes(attributes)
            .build();
            
        CreateQueueResponse response = sqsClient.createQueue(request);
        return response.queueUrl();
    }
}

4. 메시지 전송

import software.amazon.awssdk.services.sqs.model.SendMessageRequest;
import software.amazon.awssdk.services.sqs.model.SendMessageResponse;
import software.amazon.awssdk.services.sqs.model.MessageAttributeValue;
import java.util.Map;
import java.util.HashMap;

public class MessageSender {
    
    private final SqsClient sqsClient;
    
    public MessageSender() {
        this.sqsClient = SQSClientConfig.getSqsClient();
    }
    
    // 표준 큐에 메시지 전송
    public String sendMessage(String queueUrl, String messageBody) {
        SendMessageRequest request = SendMessageRequest.builder()
            .queueUrl(queueUrl)
            .messageBody(messageBody)
            .build();
            
        SendMessageResponse response = sqsClient.sendMessage(request);
        return response.messageId();
    }
    
    // 메시지 속성과 함께 전송
    public String sendMessageWithAttributes(String queueUrl, String messageBody, 
                                          Map<String, String> attributes) {
        
        Map<String, MessageAttributeValue> messageAttributes = new HashMap<>();
        
        for (Map.Entry<String, String> entry : attributes.entrySet()) {
            messageAttributes.put(entry.getKey(), 
                MessageAttributeValue.builder()
                    .stringValue(entry.getValue())
                    .dataType("String")
                    .build());
        }
        
        SendMessageRequest request = SendMessageRequest.builder()
            .queueUrl(queueUrl)
            .messageBody(messageBody)
            .messageAttributes(messageAttributes)
            .build();
            
        SendMessageResponse response = sqsClient.sendMessage(request);
        return response.messageId();
    }
    
    // FIFO 큐에 메시지 전송
    public String sendFIFOMessage(String queueUrl, String messageBody, 
                                String messageGroupId, String deduplicationId) {
        SendMessageRequest request = SendMessageRequest.builder()
            .queueUrl(queueUrl)
            .messageBody(messageBody)
            .messageGroupId(messageGroupId)
            .messageDeduplicationId(deduplicationId)
            .build();
            
        SendMessageResponse response = sqsClient.sendMessage(request);
        return response.messageId();
    }
}

5. 메시지 수신 및 처리

import software.amazon.awssdk.services.sqs.model.*;
import java.util.List;

public class MessageReceiver {
    
    private final SqsClient sqsClient;
    
    public MessageReceiver() {
        this.sqsClient = SQSClientConfig.getSqsClient();
    }
    
    // 메시지 수신
    public List<Message> receiveMessages(String queueUrl, int maxMessages) {
        ReceiveMessageRequest request = ReceiveMessageRequest.builder()
            .queueUrl(queueUrl)
            .maxNumberOfMessages(maxMessages)
            .waitTimeSeconds(20) // 롱 폴링 설정
            .messageAttributeNames("All")
            .build();
            
        ReceiveMessageResponse response = sqsClient.receiveMessage(request);
        return response.messages();
    }
    
    // 메시지 처리 및 삭제
    public void processMessages(String queueUrl) {
        while (true) {
            List<Message> messages = receiveMessages(queueUrl, 10);
            
            for (Message message : messages) {
                try {
                    // 메시지 처리 로직
                    processMessage(message);
                    
                    // 처리 완료 후 메시지 삭제
                    deleteMessage(queueUrl, message.receiptHandle());
                    
                } catch (Exception e) {
                    System.err.println("메시지 처리 실패: " + e.getMessage());
                    // 필요시 DLQ로 전송하거나 재처리 로직 구현
                }
            }
            
            if (messages.isEmpty()) {
                try {
                    Thread.sleep(1000); // 1초 대기
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                    break;
                }
            }
        }
    }
    
    private void processMessage(Message message) {
        System.out.println("메시지 처리: " + message.body());
        System.out.println("메시지 ID: " + message.messageId());
        
        // 메시지 속성 출력
        message.messageAttributes().forEach((key, value) -> 
            System.out.println("속성 " + key + ": " + value.stringValue()));
    }
    
    // 메시지 삭제
    public void deleteMessage(String queueUrl, String receiptHandle) {
        DeleteMessageRequest request = DeleteMessageRequest.builder()
            .queueUrl(queueUrl)
            .receiptHandle(receiptHandle)
            .build();
            
        sqsClient.deleteMessage(request);
    }
}

실제 사용 예제

주문 처리 시스템

public class OrderProcessingSystem {
    
    private final MessageSender messageSender;
    private final MessageReceiver messageReceiver;
    private final String orderQueueUrl;
    
    public OrderProcessingSystem() {
        QueueManager queueManager = new QueueManager();
        this.orderQueueUrl = queueManager.createStandardQueue("order-processing-queue");
        this.messageSender = new MessageSender();
        this.messageReceiver = new MessageReceiver();
    }
    
    // 주문 접수
    public void submitOrder(Order order) {
        String orderJson = convertToJson(order);
        
        Map<String, String> attributes = new HashMap<>();
        attributes.put("OrderType", order.getType());
        attributes.put("Priority", order.getPriority());
        attributes.put("CustomerId", order.getCustomerId());
        
        String messageId = messageSender.sendMessageWithAttributes(
            orderQueueUrl, orderJson, attributes);
            
        System.out.println("주문 접수 완료. 메시지 ID: " + messageId);
    }
    
    // 주문 처리 워커 시작
    public void startOrderProcessor() {
        new Thread(() -> {
            System.out.println("주문 처리 워커 시작...");
            
            while (true) {
                List<Message> messages = messageReceiver.receiveMessages(orderQueueUrl, 5);
                
                for (Message message : messages) {
                    try {
                        processOrder(message);
                        messageReceiver.deleteMessage(orderQueueUrl, message.receiptHandle());
                    } catch (Exception e) {
                        System.err.println("주문 처리 실패: " + e.getMessage());
                    }
                }
            }
        }).start();
    }
    
    private void processOrder(Message message) {
        Order order = parseOrderFromJson(message.body());
        
        // 주문 처리 로직
        System.out.println("주문 처리 중: " + order.getOrderId());
        
        // 결제 처리
        processPayment(order);
        
        // 재고 확인
        checkInventory(order);
        
        // 배송 준비
        prepareShipment(order);
        
        System.out.println("주문 처리 완료: " + order.getOrderId());
    }
}

모범 사례

1. 가시성 타임아웃 설정

메시지 처리 시간에 맞게 가시성 타임아웃을 설정하세요:

public void setVisibilityTimeout(String queueUrl, int timeoutSeconds) {
    Map<String, String> attributes = new HashMap<>();
    attributes.put("VisibilityTimeout", String.valueOf(timeoutSeconds));
    
    SetQueueAttributesRequest request = SetQueueAttributesRequest.builder()
        .queueUrl(queueUrl)
        .attributes(attributes)
        .build();
        
    sqsClient.setQueueAttributes(request);
}

2. 데드 레터 큐 (DLQ) 설정

public void setupDeadLetterQueue(String mainQueueUrl, String dlqArn, int maxReceiveCount) {
    Map<String, String> attributes = new HashMap<>();
    
    // 리드라이브 정책 설정
    String redrivePolicy = String.format(
        "{\"deadLetterTargetArn\":\"%s\",\"maxReceiveCount\":%d}",
        dlqArn, maxReceiveCount
    );
    
    attributes.put("RedrivePolicy", redrivePolicy);
    
    SetQueueAttributesRequest request = SetQueueAttributesRequest.builder()
        .queueUrl(mainQueueUrl)
        .attributes(attributes)
        .build();
        
    sqsClient.setQueueAttributes(request);
}

3. 배치 처리

여러 메시지를 한 번에 처리하여 성능을 향상시키세요:

import software.amazon.awssdk.services.sqs.model.*;
import java.util.ArrayList;
import java.util.List;

public void sendMessageBatch(String queueUrl, List<String> messages) {
    List<SendMessageBatchRequestEntry> entries = new ArrayList<>();
    
    for (int i = 0; i < messages.size(); i++) {
        entries.add(SendMessageBatchRequestEntry.builder()
            .id("msg_" + i)
            .messageBody(messages.get(i))
            .build());
    }
    
    SendMessageBatchRequest request = SendMessageBatchRequest.builder()
        .queueUrl(queueUrl)
        .entries(entries)
        .build();
        
    SendMessageBatchResponse response = sqsClient.sendMessageBatch(request);
    
    // 실패한 메시지 확인
    if (!response.failed().isEmpty()) {
        System.err.println("전송 실패한 메시지들:");
        response.failed().forEach(failed -> 
            System.err.println("ID: " + failed.id() + ", 오류: " + failed.message()));
    }
}

모니터링 및 CloudWatch 지표

import software.amazon.awssdk.services.cloudwatch.CloudWatchClient;
import software.amazon.awssdk.services.cloudwatch.model.*;

public class SQSMonitoring {
    
    private final CloudWatchClient cloudWatchClient;
    
    public SQSMonitoring() {
        this.cloudWatchClient = CloudWatchClient.builder()
            .region(Region.AP_NORTHEAST_2)
            .build();
    }
    
    // 큐의 메시지 수 확인
    public void getQueueMetrics(String queueName) {
        GetMetricStatisticsRequest request = GetMetricStatisticsRequest.builder()
            .namespace("AWS/SQS")
            .metricName("ApproximateNumberOfMessages")
            .dimensions(Dimension.builder()
                .name("QueueName")
                .value(queueName)
                .build())
            .startTime(Instant.now().minus(1, ChronoUnit.HOURS))
            .endTime(Instant.now())
            .period(300) // 5분 간격
            .statistics(Statistic.AVERAGE)
            .build();
            
        GetMetricStatisticsResponse response = cloudWatchClient.getMetricStatistics(request);
        
        System.out.println("큐 메트릭 정보:");
        response.datapoints().forEach(datapoint -> 
            System.out.println("시간: " + datapoint.timestamp() + 
                             ", 평균 메시지 수: " + datapoint.average()));
    }
}

SQS 아키텍처 다이어그램

보안 고려사항

IAM 정책 예제

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Action": [
                "sqs:SendMessage",
                "sqs:ReceiveMessage",
                "sqs:DeleteMessage",
                "sqs:GetQueueAttributes"
            ],
            "Resource": "arn:aws:sqs:ap-northeast-2:123456789012:my-queue"
        }
    ]
}

서버 사이드 암호화 설정

public void enableEncryption(String queueUrl) {
    Map<String, String> attributes = new HashMap<>();
    
    // KMS 키 사용하여 암호화 활성화
    attributes.put("KmsMasterKeyId", "alias/aws/sqs");
    attributes.put("KmsDataKeyReusePeriodSeconds", "300");
    
    SetQueueAttributesRequest request = SetQueueAttributesRequest.builder()
        .queueUrl(queueUrl)
        .attributes(attributes)
        .build();
        
    sqsClient.setQueueAttributes(request);
}

성능 최적화 팁

  1. 롱 폴링 사용: 빈 응답을 줄이고 비용을 절약
  2. 배치 작업: 여러 메시지를 한 번에 처리
  3. 적절한 가시성 타임아웃: 메시지 처리 시간에 맞게 설정
  4. DLQ 활용: 실패한 메시지를 별도로 처리
  5. 메시지 크기 최적화: 큰 페이로드는 S3 사용 고려

결론

Amazon SQS는 현대적인 클라우드 애플리케이션에서 시스템 간 통신을 위한 강력하고 확장 가능한 솔루션입니다. 자바 애플리케이션에서 AWS SDK를 사용하여 쉽게 통합할 수 있으며, 적절한 설계 패턴과 모범 사례를 따르면 견고하고 효율적인 메시지 처리 시스템을 구축할 수 있습니다.

SQS를 활용하여 마이크로서비스 아키텍처에서 서비스 간 결합도를 낮추고, 시스템의 복원력과 확장성을 크게 향상시킬 수 있습니다.

profile
AI를 꿈꾸는 BackEnd개발자

0개의 댓글