KubeMQ
ConnectorsAWS (SQS & SNS)Reference

Migrating from AWS SQS/SNS

Override the AWS SDK endpoint to KubeMQ: quick start, SQS, SNS fan-out, FIFO and redrive migrate; email, SMS and Lambda subscriptions do not.

Point your existing AWS SQS/SNS application at KubeMQ by changing only the endpoint URL. Same AWS SDK, same code, same SQS and SNS wire protocols. There is no SDK to adopt, no proto, no data migration — SQS data lives in normal KubeMQ Queue channels. Rollback is config-only (CONNECTORS_AWS_ENABLE=false).

If you already run an SDK against a LocalStack endpoint, the switch is the same single variable — point it at the KubeMQ AWS connector instead of LocalStack.

But several connector behaviors deviate from real AWS. Read the deviations below before you migrate; most are invisible until a corner case hits production.

Overview

The AWS connector exposes the real AWS SQS and SNS wire protocols (the AWS JSON and Query protocols) over HTTP, so unmodified AWS SDK clients — boto3, the AWS SDK for Go/JavaScript, the AWS CLI — talk to KubeMQ without any code changes. An endpoint-override environment variable is all that changes on the client side.

SQS queues map onto native KubeMQ Queue channels (sqs.{name}), making AWS producers and native gRPC/REST consumers interoperable on the same messages. SNS topics are virtual (registry-only) and fan out to SQS subscriptions and HTTP/HTTPS webhook endpoints. Requests are authenticated with AWS Signature V4.

  • Canonical client: AWS SDK boto3 1.x (Python).
  • Port: 4566 (HTTP), LocalStack's default. Change it with CONNECTORS_AWS_PORT.
  • Drop-in level: endpoint-only — change the SDK endpoint override; no application code changes.
  • Enable: off by default. Turn it on with CONNECTORS_AWS_ENABLE=true; on Kubernetes, spec.aws.enabled: true.

Quick start: repoint one client

Try the switch on your own machine before you plan the migration.

The Docker tab points one client at the KubeMQ you started with Try KubeMQ, turning its connector on where needed and keeping your evaluation or license and your data. The Kubernetes tab works on a cluster you installed with Install on Kubernetes. Use it to try the switch, not for production; see Plans compared for staging and production licenses. If your current broker runs on this machine, stop it first, because two programs cannot publish the same port. The kmq command runs natively on Windows; the other commands are bash, so on Windows run them in WSL (Windows Subsystem for Linux).

You started KubeMQ with Try KubeMQ or Install with Docker.

Turn the AWS connector on

The connector is off by default. This command turns it on and publishes port 4566 on the same container and volume. If you started KubeMQ with the Docker command, skip it and use Without kmq below.

Terminal
kmq deploy update --installation kubemq --connector aws --deadline 10m

Check that installation_id is the same as before and data_retained is true. You should see:

Output
{
  "installation_id": "…",
  "state": "…",
  "target": "…",
  "data_retained": true,
  "next_action": "…"
}

Point your client at KubeMQ

Before
# No endpoint override: the SDK uses the regional AWS endpoint.
After
export AWS_ENDPOINT_URL_SQS=http://localhost:4566
export AWS_ENDPOINT_URL_SNS=http://localhost:4566
export AWS_ACCESS_KEY_ID=test
export AWS_SECRET_ACCESS_KEY=test
export AWS_DEFAULT_REGION=us-east-1

The keys can be any value, but the SDK still needs them to sign the request.

Send and receive one message

Install the client with pip install boto3, then run the file with python3.

smoke_test.py
import boto3

ENDPOINT = "http://localhost:4566"
CREDS = dict(
    endpoint_url=ENDPOINT,
    region_name="us-east-1",
    aws_access_key_id="test",
    aws_secret_access_key="test",
)

sqs = boto3.client("sqs", **CREDS)

# 1. Ensure the queue exists
queue_url = sqs.create_queue(QueueName="smoke-test")["QueueUrl"]

# 2. Publish one message
sqs.send_message(QueueUrl=queue_url, MessageBody="smoke-test-payload")
print("Sent: smoke-test-payload")

# 3. Receive and confirm arrival
resp = sqs.receive_message(QueueUrl=queue_url, WaitTimeSeconds=5)
msgs = resp.get("Messages", [])
assert msgs, "ERROR: no message received"
assert msgs[0]["Body"] == "smoke-test-payload", f"Unexpected body: {msgs[0]['Body']}"
print(f"Received: {msgs[0]['Body']}")

# 4. Acknowledge
sqs.delete_message(QueueUrl=queue_url, ReceiptHandle=msgs[0]["ReceiptHandle"])
print("Acknowledged. Smoke test PASSED.")

You should see:

Output
Sent: smoke-test-payload
Received: smoke-test-payload
Acknowledged. Smoke test PASSED.

To turn the connector off again, run the same kmq deploy update command without --connector. Your messages and your evaluation or license stay.

You installed KubeMQ with Install on Kubernetes.

Turn the connector on

Put the connector's key in cluster-values.yaml. If Helm installed your cluster, add it to the file you created on Install on Kubernetes; if kmq installed it, create the file with only this key.

cluster-values.yaml (add)
aws:
  enabled: true
  sessionAffinity: ClientIP

sessionAffinity: ClientIP keeps each client on one server; see Service exposure & session affinity.

If kmq installed your cluster, apply the file with kmq:

Terminal
kmq deploy update --installation messaging --kube-context YOUR_KUBE_CONTEXT --values cluster-values.yaml --deadline 10m

Replace:

  • YOUR_KUBE_CONTEXT — your cluster's kubectl context; kubectl config get-contexts lists them.

Check that data_retained is true. You should see:

Output
{
  "installation_id": "…",
  "state": "…",
  "target": "…",
  "data_retained": true,
  "next_action": "…"
}

If Helm installed your cluster, apply the file with the command in Apply a settings change.

Either way, check that the connector's Service exists:

Terminal
kubectl --context YOUR_KUBE_CONTEXT -n kubemq get svc messaging-aws

Check that the row lists 4566/TCP. You should see:

Output
NAME             TYPE        CLUSTER-IP   EXTERNAL-IP   PORT(S)     AGE
messaging-aws    ClusterIP   …            <none>        4566/TCP    …

Reach the connector from your machine

Keep this running in a second terminal:

Terminal 2
kubectl --context YOUR_KUBE_CONTEXT -n kubemq port-forward svc/messaging-aws 4566:4566

You should see:

Output
Forwarding from 127.0.0.1:4566 -> 4566
Forwarding from [::1]:4566 -> 4566

Send and receive one message

Run steps 2 and 3 of the Docker tab unchanged. The port-forward serves the connector at the same localhost:4566 address.

Before you move production traffic

  • Read What does not migrate before you commit.
  • Plan how traffic moves and how you roll back: run both brokers, move one queue or topic at a time, and keep AWS until you have verified the move.
  • Install KubeMQ for real: Choose your path.

Compatibility Matrix

The column below is the AWS SQS/SNS slice of the master cross-protocol matrix.

DimensionAWS SQS/SNS
Drop-in levelendpoint-only
Point-to-point queues✅ SQS
Pub/sub (non-durable)✅ SNS
Durable / persistent subscriptions✅ (SQS durable)
Request/reply (RPC)N/A (no RPC)
Ordering guarantee✅ FIFO
TransactionsN/A
Dead-letter / redrive✅ redrive + move-task
Selectors / filtering / wildcards✅ SNS filter policies¹
Auth modelSigV4 / accept-any
TLS / mTLS❌ connector HTTP-only²
Top unsupportedSNS email/SMS/Lambda/push; queue/topic IAM policies; TLS at connector

¹ SNS filter policies work with the MessageAttributes scope only; the MessageBody scope is rejected (InvalidParameter).

² The connector listens on plain HTTP. Terminate TLS at a reverse proxy (see Security).

Connection / Endpoint Migration

No code changes are required. Set the AWS SDK endpoint-override environment variables to point at the KubeMQ host:

Terminal
# Before (real AWS): no override; the SDK uses the regional AWS endpoint.

# After (KubeMQ):
export AWS_ENDPOINT_URL_SQS=http://kubemq-host:4566
export AWS_ENDPOINT_URL_SNS=http://kubemq-host:4566
export AWS_ACCESS_KEY_ID=AKIAEXAMPLE
export AWS_SECRET_ACCESS_KEY=secret
export AWS_DEFAULT_REGION=kubemq    # any region value works; the region segment is not enforced

All current AWS SDKs and the AWS CLI honor AWS_ENDPOINT_URL_SQS / AWS_ENDPOINT_URL_SNS. If your SDK version predates these variables, use the per-client endpoint override instead:

import boto3

sqs = boto3.client(
    "sqs",
    endpoint_url="http://kubemq-host:4566",
    region_name="kubemq",
    aws_access_key_id="AKIAEXAMPLE",
    aws_secret_access_key="secret",
)

In the default accept-any mode the connector does not cryptographically verify the signature — but the SDK must still form a syntactically valid SigV4 request whose credential-scope service is sqs or sns. So you must give the SDK a dummy access key, secret, and region (any values); omitting them yields a "missing credentials" SDK error. The only SigV4-exempt action is SNS ConfirmSubscription.

Credential mapping

For each AccessKeyId your application uses, add a credential entry to the KubeMQ configuration. The ClientID field maps that key to a KubeMQ identity for authorization and audit:

config.toml
[Connectors.Aws]
Enable = true
Port   = "4566"

[[Connectors.Aws.Credentials]]
AccessKeyId     = "AKIAEXAMPLE"
SecretAccessKey = "secret"
ClientID        = "billing-service"   # optional; defaults to AccessKeyId

Recreate resources

Queues, topics, subscriptions, and their attributes/tags must be recreated through the AWS API against KubeMQ. The registry is authoritative — existing SQS/SNS resources from real AWS are not migrated automatically.

Concept & Destination Mapping

AWS conceptKubeMQ patternKubeMQ channel
SQS standard queueQueuessqs.{queue-name}
SQS FIFO queueQueues (per group)sqs.{queue-name} per MessageGroupId
SNS topicVirtual (registry-only)sns.{topic-name} (authorization pseudo-channel)
SNS subscription → SQSQueues (fan-out)sqs.{target-queue-name}

FIFO ordering: each MessageGroupId maps to its own KubeMQ Queue channel, preserving per-group ordering. MessageGroupId is required on every FIFO send.

Message attributes round-trip losslessly through KubeMQ message Tags (sqs_attr_{Name} = {DataType}|{value}). The connector also stamps sqs_message_id, sqs_sender_id (the authenticated ClientID), and sqs_trace_header when present.

Native interop: native KubeMQ gRPC/REST clients can produce and consume on sqs.* channels directly. Natively produced messages lack sqs_* tags; MessageId falls back to the broker MessageID and no per-message redrive policy is stamped (see deviations).

ARNs use the configured Region (default kubemq) and AccountId (default 000000000000). Set Connectors.Aws.Region / Connectors.Aws.AccountId if your tooling validates ARN format.

Canonical Client Example

The examples use the AWS SDK boto3 1.x: boto3.client, create_queue, send_message, receive_message, delete_message, create_topic, subscribe, and publish.

SQS — create, send, receive, delete

import boto3

# boto3 1.x — point the SDK at KubeMQ
sqs = boto3.client(
    "sqs",
    endpoint_url="http://kubemq-host:4566",
    region_name="kubemq",
    aws_access_key_id="AKIAEXAMPLE",
    aws_secret_access_key="secret",
)

# Create a standard queue (idempotent — returns the URL if it already exists)
resp = sqs.create_queue(QueueName="orders")
queue_url = resp["QueueUrl"]

# Send a message
sqs.send_message(
    QueueUrl=queue_url,
    MessageBody='{"orderId": "A-001", "amount": 99.95}',
    MessageAttributes={
        "source": {"DataType": "String", "StringValue": "checkout-service"},
    },
)

# Receive (long-poll up to 20 s)
resp = sqs.receive_message(
    QueueUrl=queue_url,
    MaxNumberOfMessages=1,
    WaitTimeSeconds=20,
    MessageAttributeNames=["All"],
)
for msg in resp.get("Messages", []):
    print(f"Received: {msg['Body']}")
    # Acknowledge by deleting
    sqs.delete_message(QueueUrl=queue_url, ReceiptHandle=msg["ReceiptHandle"])

FIFO queue

# Create a FIFO queue (name must end in .fifo)
resp = sqs.create_queue(
    QueueName="orders.fifo",
    Attributes={
        "FifoQueue": "true",
        "ContentBasedDeduplication": "true",
    },
)
fifo_url = resp["QueueUrl"]

# Send to a specific message group (preserves ordering per group)
sqs.send_message(
    QueueUrl=fifo_url,
    MessageBody='{"orderId": "A-002"}',
    MessageGroupId="region-us-east",
)

Dead-letter queue (redrive)

import json

# 1. Create the DLQ
dlq_resp = sqs.create_queue(QueueName="orders-dlq")
dlq_url = dlq_resp["QueueUrl"]
dlq_arn = sqs.get_queue_attributes(
    QueueUrl=dlq_url, AttributeNames=["QueueArn"]
)["Attributes"]["QueueArn"]

# 2. Attach a redrive policy to the source queue
sqs.set_queue_attributes(
    QueueUrl=queue_url,
    Attributes={
        "RedrivePolicy": json.dumps({
            "deadLetterTargetArn": dlq_arn,
            "maxReceiveCount": "3",
        }),
    },
)
# After 3 failed receives the broker automatically moves the message to orders-dlq.

SNS fan-out (topic → SQS)

sns = boto3.client(
    "sns",
    endpoint_url="http://kubemq-host:4566",
    region_name="kubemq",
    aws_access_key_id="AKIAEXAMPLE",
    aws_secret_access_key="secret",
)

# Create the topic and subscribe the SQS queue
topic_arn = sns.create_topic(Name="product-events")["TopicArn"]

orders_arn = sqs.get_queue_attributes(
    QueueUrl=queue_url, AttributeNames=["QueueArn"]
)["Attributes"]["QueueArn"]

sns.subscribe(
    TopicArn=topic_arn,
    Protocol="sqs",          # only 'sqs' and 'http'/'https' are supported
    Endpoint=orders_arn,
)

# Publish — the message fans out to all confirmed subscriptions
sns.publish(
    TopicArn=topic_arn,
    Message='{"event": "product.created", "id": "P-100"}',
    MessageAttributes={
        "category": {"DataType": "String", "StringValue": "electronics"},
    },
)

Security

Authentication — SigV4

The connector verifies AWS Signature V4 (Authorization: AWS4-HMAC-SHA256 …). Header-style and query-string-style signatures are both supported, with a ±15-minute clock-skew window.

  • Credentials are configured under Connectors.Aws.Credentials (or CredentialsData for environment / Kubernetes Secret injection).
  • Accept-any mode: if no credentials are configured the connector only parses the AccessKeyId and uses it as the ClientID. This is intended for local development only and is logged once at startup.
  • Failure cases: unknown key → HTTP 403 InvalidClientTokenId; bad signature / clock skew → HTTP 403 SignatureDoesNotMatch; malformed header → HTTP 400 IncompleteSignature.
  • ConfirmSubscription requests with no Authorization header bypass SigV4 — the single-purpose confirmation token is the authenticator (AWS parity). No other action is exempt.

TLS — terminate at a reverse proxy

The connector listens on plain HTTP only. It does not support TLS natively. To secure traffic in transit, place a TLS-terminating reverse proxy (nginx, Envoy, HAProxy, an AWS ALB, …) in front of the connector port, and configure clients to target the proxy's HTTPS endpoint.

Optional SNS message signing (off by default)

By default (Connectors.Aws.MessageSigning = false) SNS notification envelopes are unsigned. Set Connectors.Aws.MessageSigning = true to emit SigV2 RSA-SHA256 signatures. Note that the signing certificate is self-signed and not Amazon-rooted; SDK verifiers that pin Amazon's cert chain will still reject the signature. See also SNS notification signatures are unsigned by default.

Authorization (Casbin)

When the KubeMQ authorization service is enabled, every channel-mapped operation is enforced against the existing Casbin policy:

Operation classCasbin check
SendMessage, SendMessageBatch, SNS Publish (per matched SQS target)write on sqs.{queue}
ReceiveMessage, DeleteMessage, ChangeMessageVisibilityread on sqs.{queue}
Queue management (Create / Delete / Purge / Set / Tag)write on sqs.{queue}
Topic & subscription managementwrite on sns.{topic}

ListQueues, ListTopics, ListSubscriptions, and GetQueueUrl are allowed for any authenticated principal and return unfiltered results. See Authentication & security for policy configuration.

What Does NOT Migrate / Deviations

Unsupported SNS subscription protocols

Only sqs and http/https subscription protocols are supported. The following AWS SNS delivery targets are not supported and return InvalidParameter:

  • Email / email-JSON
  • SMS
  • AWS Lambda
  • Mobile push (APNs, GCM, ADM, Baidu)

FIFO topics additionally restrict subscriptions to sqs only (http/https are rejected for .fifo topics).

No TLS at the connector

The connector is HTTP-only (see TLS — terminate at a reverse proxy). In particular, do not configure clients to send to https://kubemq-host:4566 directly.

Queue/topic IAM policies ignored

Policy fields on SetQueueAttributes and SetTopicAttributes are accepted but not enforced. Authorization is handled by KubeMQ's Casbin engine (see above).

FIFO SequenceNumber is broker-derived

FIFO SequenceNumber is derived from the broker-assigned timestamp and zero-padded to 20 digits. On send it reflects the broker send-timestamp; on receive it reflects the true broker sequence (still per-group increasing). It is not a monotonic counter matching AWS semantics — do not use it for ordering comparisons across producers.

FIFO topic ContentBasedDeduplication not supported

Topic-level ContentBasedDeduplication on a FIFO topic is unsupported: setting it returns InvalidParameter and getting it always reads "false". Deduplication is enforced on the target FIFO queues, not at the topic level — pass an explicit MessageDeduplicationId instead.

Message retention clamped

MessageRetentionPeriod is clamped to Connectors.Aws.MaxExpirationSeconds at send time. Changes to SetQueueAttributes retention are not retroactive to messages already in the queue (AWS applies them retroactively).

SNS notification signatures are unsigned by default

By default (Connectors.Aws.MessageSigning = false) the SNS notification envelope is unsigned (SignatureVersion: "1", empty Signature / SigningCertURL) — webhook consumers must skip verification. Set Connectors.Aws.MessageSigning = true to emit SignatureVersion: "2" SigV2 RSA-SHA256 signatures (self-signed cert; not Amazon-rooted, so SDK verifiers that pin Amazon's cert chain still won't validate it).

Receipt handles and the in-flight tracker are node-local

In a cluster, DeleteMessage and ChangeMessageVisibility must reach the same node that served the ReceiveMessage. The receipt handle encodes the node ID; other nodes return ReceiptHandleIsInvalid. Use session-sticky (source-IP or connection-sticky) load balancing, or pin each consumer to one node.

Hard crash loses in-flight messages and pending webhook retries

KubeMQ's downstream read is destructive. Only the graceful-shutdown path returns in-flight messages to their queues; SNS delivery/retry state is in-memory on the publishing node, so a hard kill loses it.

Native producers bypass SQS policy stamping

Messages published to sqs.* channels by native KubeMQ gRPC/REST clients carry no per-message retention or redrive policy and no sqs_* identity tags.

Empty-queue short-poll latency

A ReceiveMessage with WaitTimeSeconds=0 on an empty queue returns within ~1 second instead of immediately (the broker's downstream wait granularity has a 1-second minimum).

Other out-of-scope operations

The following simply won't work (see Capabilities): KMS / SSE, AddPermission / RemovePermission, SQS message-move tasks, signed-notification verification, extended-client messages over 256 KiB, and CloudWatch metrics emulation. Cross-account is unsupported — QueueOwnerAWSAccountId is accepted and ignored; there is a single configurable AccountId.

Verification Smoke Test

For the basic round trip, see the quick start; this section adds deeper checks.

The same round trip with the AWS CLI, using the variables you exported in the quick start's step 2:

Terminal
aws --endpoint-url http://localhost:4566 sqs create-queue --queue-name smoke-test

You should see:

Output
{
    "QueueUrl": "http://localhost:4566/000000000000/smoke-test"
}
Terminal
aws --endpoint-url http://localhost:4566 sqs send-message --queue-url http://localhost:4566/000000000000/smoke-test --message-body smoke-test-payload --query MD5OfMessageBody --output text

You should see:

Output
1ae3878661ac01788a3c42b3643c1b49
Terminal
aws --endpoint-url http://localhost:4566 sqs receive-message --queue-url http://localhost:4566/000000000000/smoke-test --wait-time-seconds 5 --query 'Messages[0].Body' --output text

You should see:

Output
smoke-test-payload

See Also

Was this page helpful?

On this page