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

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

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'
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;
}
}

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();
}
}
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();
}
}
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());
}
}
메시지 처리 시간에 맞게 가시성 타임아웃을 설정하세요:
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);
}
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);
}
여러 메시지를 한 번에 처리하여 성능을 향상시키세요:
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()));
}
}
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()));
}
}

{
"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);
}
Amazon SQS는 현대적인 클라우드 애플리케이션에서 시스템 간 통신을 위한 강력하고 확장 가능한 솔루션입니다. 자바 애플리케이션에서 AWS SDK를 사용하여 쉽게 통합할 수 있으며, 적절한 설계 패턴과 모범 사례를 따르면 견고하고 효율적인 메시지 처리 시스템을 구축할 수 있습니다.
SQS를 활용하여 마이크로서비스 아키텍처에서 서비스 간 결합도를 낮추고, 시스템의 복원력과 확장성을 크게 향상시킬 수 있습니다.