Docs / Backends / AWS SQS
AWS SQS
Standard and FIFO queues, dead-letter queues via RedrivePolicy, visibility timeout and long polling
import "go.digitalxero.dev/mq-aws/v2"Index
- func NewMessageBuilder() mqtypes.MessageBuilder
- func NewSQSConsumer(config *viper.Viper) (types.Consumer, error)
- func NewSQSProducer(config *viper.Viper) (types.Producer, error)
- func NewSQSTransport() mqtypes.Transport
funcNewMessageBuilder
func NewMessageBuilder() mqtypes.MessageBuilderNewMessageBuilder 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:
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
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.
| key | default | description |
|---|---|---|
| url | required | SQS queue URL e.g., “sqs://us-east-1/123456789012/my-queue” |
| queue | — | Queue name (extracted from URL if not provided) |
| handler | “sqsLogger” | Name of registered handler function or “sqsBatchLogger” when batch processing is enabled |
| region | — | AWS region (extracted from URL or uses default) |
| access_key_id | — | AWS access key ID (uses default credentials if not set) |
| secret_access_key | — | AWS secret access key (uses default credentials if not set) |
| session_token | optional | AWS session token for temporary credentials |
| max_concurrent_handlers | effective message_channel_buffer; must be positive | Maximum in-flight handlers or batches |
| max_retries | — | Maximum retry attempts for failed messages (default: 50); also used as the RedrivePolicy maxReceiveCount when a dead-letter queue is configured |
| dead_letter_queue_url | optional | Queue URL failed messages are redriven to; the main queue’s RedrivePolicy is pointed at it on connect |
| skip_verify | false | Skip TLS certificate verification for LocalStack-style endpoints, never disables TLS |
| visibility_timeout | 30 | Message visibility timeout in seconds |
| wait_time_seconds | 20 | Long polling wait time max: 20 |
| max_messages | 10 | Max messages per receive max: 10 |
| fifo_queue | auto-detect from name | Whether this is a FIFO queue |
| content_based_deduplication | false | Enable deduplication for FIFO |
| batch_size | 5 | Number of messages per batch |
| batch_timeout | “100ms” | Batch collection timeout as a duration |
| enable_batch_processing | false | Enable batch message processing |
Returns an error if configuration validation fails or transport creation fails.
funcNewSQSProducer
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.
| key | default | description |
|---|---|---|
| url | required | SQS queue URL e.g., “sqs://us-east-1/123456789012/my-queue” |
| queue_url | — | Alternative to url for direct queue URL specification |
| region | — | AWS region (extracted from URL or uses default) |
| access_key_id | — | AWS access key ID (uses default credentials if not set) |
| secret_access_key | — | AWS secret access key (uses default credentials if not set) |
| session_token | optional | AWS session token for temporary credentials |
| max_message_size | 262144/256KB | Maximum message size in bytes |
| message_group_id | — | Default FIFO message group ID (required for .fifo queues) |
| message_deduplication_id | optional | Default FIFO deduplication ID; falls back to the message’s correlation ID when unset |
| delay_seconds | 0 | Per-message delivery delay for standard queues, max 900; FIFO queues do not support per-message delays |
| skip_verify | false | Skip 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
func NewSQSTransport() mqtypes.TransportNewSQSTransport 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
mqutils