mqutils
GitLab ↗

Docs / Backends / AWS SQS

AWS SQS

Standard and FIFO queues, dead-letter queues via RedrivePolicy, visibility timeout and long polling

go
import "go.digitalxero.dev/mq-aws/v2"

Index

funcNewMessageBuilder

go
func NewMessageBuilder() mqtypes.MessageBuilder

NewMessageBuilder creates a new SQS message builder. The builder provides a fluent interface for constructing SQS messages with all supported properties and attributes.

SQS has specific limitations:

  • Message body: max 256KB
  • Message attributes: max 10 attributes
  • Attribute name: max 256 characters
  • Attribute value: max 256KB for binary, unlimited for string

Example:

go
msg := aws.NewMessageBuilder().
    WithBody([]byte(`{"orderId": "12345", "status": "pending"}`)).
    WithContentType("application/json").
    WithCorrelationId("order-12345").
    WithRoutingKey("orders"). // Used as MessageGroupId for FIFO
    WithHeaders(map[string]interface{}{
        "priority": "high",
        "source": "web-api",
    }).
    Build()

funcNewSQSConsumer

go
func NewSQSConsumer(config *viper.Viper) (types.Consumer, error)

NewSQSConsumer creates a new AWS SQS consumer with the provided configuration. This function is typically called by mqutils.NewConsumer when it detects an SQS URL.

keydefaultdescription
urlrequiredSQS queue URL e.g., “sqs://us-east-1/123456789012/my-queue”
queueQueue name (extracted from URL if not provided)
handler“sqsLogger”Name of registered handler function or “sqsBatchLogger” when batch processing is enabled
regionAWS region (extracted from URL or uses default)
access_key_idAWS access key ID (uses default credentials if not set)
secret_access_keyAWS secret access key (uses default credentials if not set)
session_tokenoptionalAWS session token for temporary credentials
max_concurrent_handlerseffective message_channel_buffer; must be positiveMaximum in-flight handlers or batches
max_retriesMaximum retry attempts for failed messages (default: 50); also used as the RedrivePolicy maxReceiveCount when a dead-letter queue is configured
dead_letter_queue_urloptionalQueue URL failed messages are redriven to; the main queue’s RedrivePolicy is pointed at it on connect
skip_verifyfalseSkip TLS certificate verification for LocalStack-style endpoints, never disables TLS
visibility_timeout30Message visibility timeout in seconds
wait_time_seconds20Long polling wait time max: 20
max_messages10Max messages per receive max: 10
fifo_queueauto-detect from nameWhether this is a FIFO queue
content_based_deduplicationfalseEnable deduplication for FIFO
batch_size5Number of messages per batch
batch_timeout“100ms”Batch collection timeout as a duration
enable_batch_processingfalseEnable batch message processing

Returns an error if configuration validation fails or transport creation fails.

funcNewSQSProducer

go
func NewSQSProducer(config *viper.Viper) (types.Producer, error)

NewSQSProducer creates a new AWS SQS producer with the provided configuration. This function is typically called by mqutils.NewProducer when it detects an SQS URL.

keydefaultdescription
urlrequiredSQS queue URL e.g., “sqs://us-east-1/123456789012/my-queue”
queue_urlAlternative to url for direct queue URL specification
regionAWS region (extracted from URL or uses default)
access_key_idAWS access key ID (uses default credentials if not set)
secret_access_keyAWS secret access key (uses default credentials if not set)
session_tokenoptionalAWS session token for temporary credentials
max_message_size262144/256KBMaximum message size in bytes
message_group_idDefault FIFO message group ID (required for .fifo queues)
message_deduplication_idoptionalDefault FIFO deduplication ID; falls back to the message’s correlation ID when unset
delay_seconds0Per-message delivery delay for standard queues, max 900; FIFO queues do not support per-message delays
skip_verifyfalseSkip TLS certificate verification for LocalStack-style endpoints, never disables TLS

Call Close on the returned producer to release transport resources after publishing has finished. Close is idempotent.

Returns an error if configuration validation fails or transport creation fails. The producer must be started with Start() before publishing messages.

funcNewSQSTransport

go
func NewSQSTransport() mqtypes.Transport

NewSQSTransport creates a new AWS SQS transport implementation. The transport provides low-level SQS operations including queue management, message sending/receiving, and health monitoring.

The transport manages:

  • AWS SDK v2 client with automatic credential resolution
  • Queue URL validation and attribute retrieval
  • Message batching for efficient operations
  • FIFO queue support with deduplication
  • Dead letter queue configuration
  • Connection health monitoring

Unlike message brokers with persistent connections, SQS uses HTTP-based API calls, so there’s no connection to maintain.

Generated by gomarkdoc

move open/ opens search anywhere