Compare AWS Messaging with CDK
Build a CDK lab to compare SQS buffering, SNS fanout, and EventBridge routing.
Introduction
30 Second Summary
A new order can wait for one worker, alert several teams, or reach only the team that handles risky purchases. Without visible evidence, choosing the right path can feel like guesswork.
In this project, you will build an order-processing decision lab with Python and AWS CDK. Controlled messages expose the buffering of Amazon SQS, the fanout of Amazon SNS, and the content-based routing of Amazon EventBridge directly in your terminal.
What You'll Build
When your lab runs, your terminal shows each order taking a distinct delivery path so you can defend every service choice with evidence.
By the end of this project, you'll have:
- Clear SQS buffering evidence showing two different order types retrieved from the same pull-based queue.
- Visible fanout proof when one SNS publication reaches two independent queues. You will also have content-based routing proof when EventBridge sends only high-value orders to fraud review.
- A reusable architecture decision report grounded in your observations. A synthesized CloudFormation template exposes the resources generated by your Python constructs.
- Secret Mission: Pause an EventBridge Subscriber while an order is published. Resume from the last processed position to recover the retained event.
Are there any prerequisites?
Basic Python syntax is enough. The lab assumes sandbox or administrator access, configured AWS CLI credentials, and VS Code.
Before We Start
Before any hands-on work begins, define the decision this lab should help you make. You are committing to an evidence-based AWS CDK messaging lab that replaces guesswork about choosing Amazon SQS, Amazon SNS, or Amazon EventBridge with behavior you observe yourself.
Set Up the CDK Workspace
Your messaging lab depends on two runtimes. The AWS CDK CLI runs on Node.js. Your infrastructure code runs with Python.
A mismatched runtime can turn a simple deployment into a confusing failure. This step verifies each tool before you create any lab resources.
In this step, get ready to:
- Verify Node.js 24.x. Install the AWS CDK CLI at 2.1144.0.
- Create the messaging-decision-lab workspace. Install its pinned Python packages inside .venv.
- Confirm your AWS identity. Bootstrap the CDK environment in us-east-1.
Verify the required runtimes
The CDK CLI requires Node.js. This project uses the Node.js 24 LTS line so everyone works with the same runtime.
- Check your Node.js version by running this command:
node --version
What does this check?
This prints the Node.js version currently available in your terminal. The first number must be 24 for this project.
✔️ I see version 24 or higher
Your Node.js runtime is ready. You can use its package manager to install the CDK CLI.
ⓧ I see an older version
An older Node.js runtime does not match this workspace. nvm lets you install the required version without replacing every Node.js installation on your machine.
- Install nvm 0.40.8 by running this command:
curl -o- https://raw.githubusercontent.com/nvm-sh/nvm/v0.40.8/install.sh | bash
What does this command do?
This downloads the nvm 0.40.8 installation script from its official repository. The script adds nvm to your shell environment.
- Close the current terminal window.
- Open a new terminal window from your Linux application launcher.
- Install Node.js 24.14.0 by running this command:
nvm install 24.14.0
What does this command do?
nvm downloads Node.js 24.14.0. It also selects that runtime for your current shell.
- Confirm the selected Node.js version by running this command:
node --version
What should I see?
You should see a version beginning with v24.
Still seeing the older version?
Close the terminal once more. Start a fresh terminal so your updated shell configuration loads.
If the old version remains active, help me check which Node.js executable my Linux shell is using.
ⓧ Command not found
Node.js is unavailable in your current shell. Install nvm first so it can manage the required runtime.
- Install nvm 0.40.8 by running this command:
curl -o- https://raw.githubusercontent.com/nvm-sh/nvm/v0.40.8/install.sh | bash
What does this command do?
This downloads the official nvm installation script. The script configures your shell to load nvm.
- Close the current terminal window.
- Open a new terminal window from your Linux application launcher.
- Install Node.js 24.14.0 by running this command:
nvm install 24.14.0
What does this command do?
nvm installs the required Node.js release. It selects that release for the new shell.
- Verify the installation by running this command:
node --version
What should I see?
You should see a version beginning with v24.
Still missing the Node.js command?
Your new terminal may not have loaded the nvm configuration. Restart your terminal before trying the version check again.
If the command remains unavailable, help me load nvm in my Linux shell.
The Python CDK library requires Python 3.10 or newer. The version check confirms that your interpreter can install the pinned project packages.
- Check your Python version by running this command:
python --version
What does this check?
This prints the Python interpreter selected by your shell. The version must be 3.10 or newer.
✔️ I see version 3.10 or higher
Your Python interpreter supports every pinned package in this workspace.
ⓧ I see an older version
The selected Python interpreter is below the packages' minimum requirement. Your Linux distribution may provide a newer interpreter through its supported package manager.
- Follow the official Python downloads guidance to install Python 3.10 or newer.
- Update your shell so the python command selects the newer interpreter.
- Confirm the active version by running this command:
python --version
What should I see?
You should see Python 3.10 or a higher version.
Still seeing the older Python version?
Your shell may still resolve python to the older interpreter. Use your Linux distribution's documented alternative-selection method to update that command.
If you need help identifying the active interpreter, help me check which Python executable my Linux shell is using.
ⓧ Command not found
Python is unavailable through the command this project uses. Install a supported release before creating the virtual environment.
- Follow the official Python downloads guidance to install Python 3.10 or newer.
- Confirm the installation by running this command:
python --version
What should I see?
You should see Python 3.10 or a higher version.
Python still unavailable?
Restart your terminal after installation so the updated command path loads.
If the command still fails, help me make Python available in my Linux shell.
The CDK CLI converts your infrastructure code into deployment instructions. Pinning it to 2.1144.0 keeps the commands consistent with this project.
- Check the installed CDK CLI version by running this command:
cdk --version
What does this check?
This prints the globally available CDK CLI version. This project expects 2.1144.0.
✔️ I see version 2.1144.0
The matching CDK CLI is ready. Your successful check also confirms that npm is available.
ⓧ I see another version
A different CDK CLI version can produce different command behavior. Install the pinned release globally to align this workspace.
- Install AWS CDK CLI 2.1144.0 by running this command:
npm install -g aws-cdk@2.1144.0
What does this command do?
npm installs the specified AWS CDK CLI release globally. The cdk command then becomes available across your terminal sessions.
- Verify the installed release by running this command:
cdk --version
What should I see?
You should see 2.1144.0 in the output.
CDK version did not change?
Your shell may be finding another global npm installation first. Restart the terminal before checking again.
If the mismatch remains, help me find which CDK executable my Linux shell is using.
ⓧ Command not found
The CDK CLI is not installed globally. The verified Node.js runtime gives npm the environment it needs for the installation.
- Install AWS CDK CLI 2.1144.0 by running this command:
npm install -g aws-cdk@2.1144.0
What does this command do?
npm installs the pinned AWS CDK CLI as a global command.
- Verify the installation by running this command:
cdk --version
What should I see?
You should see 2.1144.0 in the output.
CDK command still unavailable?
Restart the terminal so it reloads the global npm command path.
If the command remains unavailable, help me expose my global npm commands on Linux.
Create the Python workspace
A virtual environment isolates this lab's packages from other Python projects. The requirements.txt file records the exact library versions that build the stack.
- Move to your Desktop by running this command:
cd ~/Desktop
What does this command do?
This selects your Desktop as the location for the new project directory.
- Create the messaging-decision-lab directory by running these commands:
mkdir messaging-decision-lab
cd messaging-decision-lab
What do these commands do?
- The first command creates the lab directory on your Desktop.
- The second command makes that directory your terminal's current location.
- Confirm your current location by running this command:
pwd
What should I see?
The printed path should end with Desktop/messaging-decision-lab.
- Create the .venv virtual environment by running these commands:
python -m venv .venv
. .venv/bin/activate
What do these commands do?
- The first command creates an isolated Python environment inside .venv.
- The second command activates that environment for the current terminal.
You should see (.venv) near the start of your terminal prompt. That marker confirms future package installations stay inside this project.
- In the Explorer sidebar in VS Code, create requirements.txt inside the messaging-decision-lab directory.
- Paste this package list into requirements.txt.
aws-cdk-lib==2.272.0
constructs==10.8.1
boto3==1.43.109
What does this file control?
- The aws-cdk-lib pin provides the Python constructs used to define AWS infrastructure.
- The constructs pin provides the base construct model used by CDK stacks.
- The boto3 pin provides the AWS SDK used by the experiment harness.
- Save requirements.txt.
- Confirm that the file shows exactly three pinned packages in the VS Code editor.
Package list looks different?
Check each package name for a typing error. Confirm that every version uses two equals signs.
If the file still differs, help me compare my requirements file with the three required package pins.
✔️ Awesome, I've got everything!
Your dependency file is saved with all three required package pins.
ⓧ I'd like to double check the full code
aws-cdk-lib==2.272.0
constructs==10.8.1
boto3==1.43.109
The virtual environment is active. Installing from the saved file now gives the lab its infrastructure library plus Boto3.
- Install the pinned dependencies by running this command:
pip install -r requirements.txt
What does this command do?
pip reads each package pin from requirements.txt. It installs those versions inside the active virtual environment.
Dependency installation failed?
Confirm that (.venv) appears in the terminal prompt. Check that your Python version is 3.10 or newer.
For an environment-specific failure, help me diagnose my pinned Python package installation.
Before you test the imports, do you expect both libraries to load from the active virtual environment?
- Confirm the CDK library and Boto3 imports by running this command:
python -c "import aws_cdk, boto3; print('CDK and Boto3 ready')"
What should I see?
Python imports both installed libraries before printing CDK and Boto3 ready. Seeing that message proves the environment can load the packages.
Import check failed?
Confirm that your terminal remains inside messaging-decision-lab. Confirm that (.venv) still appears in the prompt.
If an import remains unavailable, help me trace which Python environment contains my CDK and Boto3 packages.
Verify AWS access and bootstrap CDK
The AWS CLI credentials choose the account where your lab runs. Checking the identity first keeps the shared bootstrap resources in the intended account.
- Display the identity attached to your configured AWS credentials by running this command:
aws sts get-caller-identity
What should I see?
The response shows your account identifier plus the authenticated identity ARN. This confirms that the AWS CLI can reach your account.
- Record the account identifier shown in the response: 123456789012.
Identity check failed?
Your configured credentials may have expired. Your active identity may also lack permission to call AWS Security Token Service.
If you cannot retrieve the identity, help me diagnose my configured AWS CLI credentials.
Before You Bootstrap
Bootstrapping creates or updates the shared CDKToolkit stack in your AWS account. That stack contains deployment resources managed through AWS CloudFormation.
This project is expected to remain under $1 for its intended workload. The cleanup section removes the application stack while preserving the shared toolkit for future CDK projects.
The first bootstrap can take a few minutes while AWS creates the shared deployment resources. A stream of progress messages means the operation is still working.
- Bootstrap your AWS account in us-east-1 by running this command:
cdk bootstrap aws://[[AWS_ACCOUNT_ID="123456789012"]]/us-east-1
What does this command do?
The command prepares the selected account and Region for CDK deployments. It deploys or updates the shared CDKToolkit stack.
Bootstrap did not complete?
Confirm that the account identifier matches the value returned by the identity check. Confirm that your sandbox permits CloudFormation plus the bootstrap resources.
If AWS reports a permissions failure, help me identify the missing permission in my CDK bootstrap output.
Before the final check, which four details do you expect the terminal to confirm about your toolchain and AWS identity?
- Run the final toolchain check with these commands:
node --version
cdk --version
python --version
aws sts get-caller-identity
What should I see?
- The Node.js output begins with v24.
- The CDK output includes 2.1144.0.
- The Python output shows 3.10 or newer.
- The AWS response shows the same account identifier you used to bootstrap us-east-1.
That is the foundation locked in. Your local workspace can load CDK plus Boto3. Your AWS account is ready to accept the first messaging stack.
Your CDK workspace is ready. Next up, you will deploy one SQS queue and watch it buffer two jobs with different processing needs.
Explore One SQS Queue
Your toolchain is ready, so the lab can move from setup to observable AWS behavior. Orders often arrive faster than workers can process them.
An Amazon SQS queue can hold that work until a consumer pulls it. In this step, you will test whether one queue can separate low-value fulfillment work from high-value fraud review work.
In this step, get ready to:
- Define a pull-based queue in an AWS CDK stack.
- Build a Python harness that sends two different order types.
- Deploy the stack to observe how one queue handles both orders.
Create the queue stack
The CDK application translates your Python stack into an AWS CloudFormation template. The first stack contains one queue so the experiment has a single delivery path.
- Switch back to VS Code with the messaging-decision-lab workspace.
- Use the file sidebar to create cdk.json inside messaging-decision-lab.
- Configure the CDK application entry point by pasting this code into cdk.json:
{
"app": "python3 app.py"
}
What does this configuration do?
- The app setting tells the CDK CLI how to start the application.
- The application starts from app.py using Python.
- Save cdk.json.
- Confirm that cdk.json appears beside requirements.txt in the file sidebar.
Is the CDK configuration invalid?
- Check that cdk.json contains double quotes around both strings.
- Remove any comma after the app value.
Ask for help checking the file if the JSON remains invalid: Help me find the JSON syntax problem in my cdk.json file.
The application file creates the stack in us-east-1. This keeps the deployment environment aligned with the region you bootstrapped earlier.
- Use the file sidebar to create app.py inside messaging-decision-lab.
- Define the CDK application by pasting this code into app.py:
#!/usr/bin/env python3
import aws_cdk as cdk
from messaging_decision_lab_stack import MessagingDecisionLabStack
app = cdk.App()
MessagingDecisionLabStack(
app,
"MessagingDecisionLabStack",
env=cdk.Environment(region="us-east-1"),
)
app.synth()
How does the application start?
- The CDK application object holds the infrastructure definitions that will be synthesized.
- The MessagingDecisionLabStack constructor creates the project stack in us-east-1.
- The app.synth() call produces the cloud assembly.
- Save app.py.
- Confirm that the editor shows MessagingDecisionLabStack as the stack created by the application.
Does the stack import look unresolved?
The imported stack file is created next. A temporary unresolved-import marker is expected until messaging_decision_lab_stack.py exists.
Ask for help if the marker remains after creating the stack file: Help me diagnose why app.py cannot import MessagingDecisionLabStack.
The stack uses one helper to apply the same queue settings consistently. Its first queue becomes the direct path for both test orders.
- Use the file sidebar to create messaging_decision_lab_stack.py inside messaging-decision-lab.
- Define the queue stack by pasting this code into messaging_decision_lab_stack.py:
from typing import Any
from aws_cdk import CfnOutput, Duration, Stack
from aws_cdk import aws_sqs as sqs
from constructs import Construct
class MessagingDecisionLabStack(Stack):
def __init__(
self,
scope: Construct,
construct_id: str,
**kwargs: Any,
) -> None:
super().__init__(scope, construct_id, **kwargs)
def create_queue(construct_id: str) -> sqs.Queue:
return sqs.Queue(
self,
construct_id,
retention_period=Duration.days(1),
visibility_timeout=Duration.minutes(15),
)
direct_queue = create_queue("DirectQueue")
CfnOutput(self, "DirectQueueUrl", value=direct_queue.queue_url)
What does the queue stack define?
- The create_queue() helper creates an SQS queue with one day of message retention.
- The queue uses a visibility timeout of 15 minutes.
- The DirectQueue construct gives both test orders one destination.
- The DirectQueueUrl output gives the experiment harness the deployed queue URL.
- Save messaging_decision_lab_stack.py.
- Confirm that the editor now recognizes MessagingDecisionLabStack.
Before you synthesize the application, what infrastructure do you expect the generated template to contain?
- Synthesize the CDK application by running this command in the activated virtual environment:
cdk synth
What does synthesis prove?
The CDK CLI runs app.py and converts the stack into a CloudFormation template. It also writes the cloud assembly to cdk.out.
A successful synthesis proves that the application can load its dependencies. It also proves that the stack definition is structurally valid.
You will see a generated template containing the queue resource and the stack output. The terminal returns to its prompt without a synthesis error.
Does synthesis fail?
- Confirm that the activated environment indicator still shows .venv.
- Check that all three Python filenames match the names shown above.
- Compare the indentation inside MessagingDecisionLabStack with the reference code.
Ask for help with the terminal output: Help me troubleshoot my cdk synth failure for this SQS stack.
Build the SQS experiment
The experiment harness uses Boto3 to send two orders with different processing needs. Its receiver prints the messages that come back from the deployed queue.
- Use the file sidebar to create lab.py inside messaging-decision-lab.
- Add the imports and output loader by pasting this first code section into lab.py:
#!/usr/bin/env python3
import json
import sys
import time
from pathlib import Path
from typing import Any
import boto3
REGION = "us-east-1"
STACK_NAME = "MessagingDecisionLabStack"
OUTPUTS_FILE = Path("cdk-outputs.json")
def load_outputs() -> dict[str, str]:
if not OUTPUTS_FILE.exists():
raise SystemExit(
"cdk-outputs.json is missing. Run cdk deploy --outputs-file cdk-outputs.json first."
)
document = json.loads(OUTPUTS_FILE.read_text(encoding="utf-8"))
return document[STACK_NAME]
How does the harness find the queue?
- The constants keep the region and stack name consistent with app.py.
- The OUTPUTS_FILE path points to the deployment output created later.
- The load_outputs() function reads the deployed stack outputs.
- The missing-file check stops the experiment with a direct instruction if deployment has not created the output file.
- Save lab.py.
- Confirm that the editor now shows the load_outputs() function below the three constants.
Is the output loader marked as invalid?
- Check that load_outputs() is aligned with the left edge of the file.
- Check that the lines inside the function use consistent indentation.
Ask for help reviewing this section: Help me fix the load_outputs function in lab.py.
Queue message bodies are strings when Boto3 receives them. The next helper converts JSON strings back into Python data while preserving any body that is not valid JSON.
- Place the cursor two blank lines below load_outputs().
- Add the message decoder by pasting this code:
def decode_body(body: str) -> Any:
try:
return json.loads(body)
except json.JSONDecodeError:
return body
Why decode the message body?
The decode_body() helper turns a JSON message body into Python data. If decoding fails, it returns the original string so the harness can still display it.
- Save lab.py.
- Confirm that decode_body() appears directly after load_outputs().
Is the decoder nested inside the loader?
Move decode_body() back to the left edge if it appears inside load_outputs(). Top-level functions must share the same indentation.
Ask for help checking the function boundary: Help me separate decode_body from load_outputs in lab.py.
A pull-based consumer asks the queue for available messages. The collector repeats that request until it has the expected count or reaches its deadline.
- Place the cursor two blank lines below decode_body().
- Add the message collector by pasting this code:
def collect_messages(
sqs_client: Any,
queue_url: str,
label: str,
expected: int,
) -> list[Any]:
collected: list[Any] = []
deadline = time.monotonic() + 15
while len(collected) < expected and time.monotonic() < deadline:
response = sqs_client.receive_message(
QueueUrl=queue_url,
MaxNumberOfMessages=min(10, expected - len(collected)),
VisibilityTimeout=900,
WaitTimeSeconds=3,
)
for message in response.get("Messages", []):
collected.append(decode_body(message["Body"]))
print(f"{label}: received {len(collected)} message(s)")
for message in collected:
print(json.dumps(message, indent=2, sort_keys=True))
return collected
How does the collector work?
- The deadline prevents the receiver from waiting forever.
- The receive request uses long polling to wait briefly for messages.
- The visibility timeout hides received messages during this experiment.
- The final loop prints each decoded body so the queue behavior is visible in the terminal.
- Save lab.py.
- Confirm that the editor shows collect_messages() as a complete function ending with return collected.
Does the collector show indentation errors?
- Align the receive request inside the while block.
- Align the append operation inside the for block.
- Align the final print loop with the while statement.
Ask for help with the nested blocks: Help me fix the control-flow indentation in collect_messages.
The SQS experiment creates one fulfillment order and one fraud-review order. Both send operations use the deployed DirectQueueUrl output.
- Place the cursor two blank lines below collect_messages().
- Add the SQS experiment by pasting this code:
def run_sqs(outputs: dict[str, str]) -> None:
sqs_client = boto3.client("sqs", region_name=REGION)
orders = [
{"orderId": "order-1001", "total": 50, "work": "fulfill"},
{"orderId": "order-1002", "total": 800, "work": "fraud-review"},
]
for order in orders:
sqs_client.send_message(
QueueUrl=outputs["DirectQueueUrl"],
MessageBody=json.dumps(order),
)
collect_messages(
sqs_client,
outputs["DirectQueueUrl"],
"One SQS consumer path",
expected=2,
)
print("Observation: SQS buffered both jobs, but it did not choose a consumer by content.")
What does the experiment control?
- The first order has a total of 50 with fulfillment work.
- The second order has a total of 800 with fraud-review work.
- Each send operation uses the same queue URL.
- The collector requests two messages from that queue URL.
- Save lab.py.
- Confirm that the editor shows both order-1001 and order-1002 inside run_sqs().
Is the SQS experiment incomplete?
- Check that both order dictionaries sit inside the orders list.
- Check that both queue references use DirectQueueUrl with matching capitalization.
Ask for help comparing the experiment: Help me check the run_sqs function in lab.py.
The report mode records the first service-selection rule from this experiment. Later experiments will extend the same report.
- Place the cursor two blank lines below run_sqs().
- Add the first report entry by pasting this code:
def print_report() -> None:
print(
"\n".join(
[
"Observed architecture decision guide",
"SQS | Pull work when consumers need buffering and processing-rate control.",
]
)
)
Why keep a report mode?
The print_report() function turns the observed queue behavior into a reusable decision rule. It keeps the recommendation grounded in the experiment.
- Save lab.py.
- Confirm that print_report() appears after run_sqs().
Does the report string look broken?
- Check that the join string contains \n.
- Check that both report lines remain inside the list.
Ask for help reviewing the nested brackets: Help me fix the print_report function in lab.py.
The final dispatcher reads one command-line mode and selects the matching experiment. It also prevents unsupported modes from running silently.
- Place the cursor two blank lines below print_report().
- Complete the harness by pasting this code:
def main() -> None:
if len(sys.argv) != 2:
raise SystemExit("Usage: python lab.py [sqs|report]")
mode = sys.argv[1]
if mode == "report":
print_report()
return
outputs = load_outputs()
actions = {"sqs": run_sqs}
if mode not in actions:
raise SystemExit("Usage: python lab.py [sqs|report]")
actions[mode](outputs)
if __name__ == "__main__":
main()
How does the dispatcher choose an action?
- The argument check requires exactly one mode after the filename.
- The report mode prints the decision guide without loading deployment outputs.
- The sqs mode loads the stack outputs before calling run_sqs().
- The final guard starts main() when the file runs as a script.
- Save lab.py.
- Confirm that main() is the final function in the file.
Does the dispatcher show a syntax error?
- Check that the actions dictionary closes before the next if statement.
- Check that the final if __name__ guard starts at the left edge.
Ask for help checking the final section: Help me fix the main dispatcher in lab.py.
Use the comparison below to catch missing lines before the deployment. Each reference file matches the complete state for this step.
✔️ Awesome, I've got everything!
Your five project files are ready. Save every open file before deploying the stack.
ⓧ I'd like to double check the full code
Your requirements.txt file should match:
aws-cdk-lib==2.272.0
constructs==10.8.1
boto3==1.43.109
Your cdk.json file should match:
{
"app": "python3 app.py"
}
Your app.py file should match:
#!/usr/bin/env python3
import aws_cdk as cdk
from messaging_decision_lab_stack import MessagingDecisionLabStack
app = cdk.App()
MessagingDecisionLabStack(
app,
"MessagingDecisionLabStack",
env=cdk.Environment(region="us-east-1"),
)
app.synth()
Your messaging_decision_lab_stack.py file should match:
from typing import Any
from aws_cdk import CfnOutput, Duration, Stack
from aws_cdk import aws_sqs as sqs
from constructs import Construct
class MessagingDecisionLabStack(Stack):
def __init__(
self,
scope: Construct,
construct_id: str,
**kwargs: Any,
) -> None:
super().__init__(scope, construct_id, **kwargs)
def create_queue(construct_id: str) -> sqs.Queue:
return sqs.Queue(
self,
construct_id,
retention_period=Duration.days(1),
visibility_timeout=Duration.minutes(15),
)
direct_queue = create_queue("DirectQueue")
CfnOutput(self, "DirectQueueUrl", value=direct_queue.queue_url)
Your lab.py file should match:
#!/usr/bin/env python3
import json
import sys
import time
from pathlib import Path
from typing import Any
import boto3
REGION = "us-east-1"
STACK_NAME = "MessagingDecisionLabStack"
OUTPUTS_FILE = Path("cdk-outputs.json")
def load_outputs() -> dict[str, str]:
if not OUTPUTS_FILE.exists():
raise SystemExit(
"cdk-outputs.json is missing. Run cdk deploy --outputs-file cdk-outputs.json first."
)
document = json.loads(OUTPUTS_FILE.read_text(encoding="utf-8"))
return document[STACK_NAME]
def decode_body(body: str) -> Any:
try:
return json.loads(body)
except json.JSONDecodeError:
return body
def collect_messages(
sqs_client: Any,
queue_url: str,
label: str,
expected: int,
) -> list[Any]:
collected: list[Any] = []
deadline = time.monotonic() + 15
while len(collected) < expected and time.monotonic() < deadline:
response = sqs_client.receive_message(
QueueUrl=queue_url,
MaxNumberOfMessages=min(10, expected - len(collected)),
VisibilityTimeout=900,
WaitTimeSeconds=3,
)
for message in response.get("Messages", []):
collected.append(decode_body(message["Body"]))
print(f"{label}: received {len(collected)} message(s)")
for message in collected:
print(json.dumps(message, indent=2, sort_keys=True))
return collected
def run_sqs(outputs: dict[str, str]) -> None:
sqs_client = boto3.client("sqs", region_name=REGION)
orders = [
{"orderId": "order-1001", "total": 50, "work": "fulfill"},
{"orderId": "order-1002", "total": 800, "work": "fraud-review"},
]
for order in orders:
sqs_client.send_message(
QueueUrl=outputs["DirectQueueUrl"],
MessageBody=json.dumps(order),
)
collect_messages(
sqs_client,
outputs["DirectQueueUrl"],
"One SQS consumer path",
expected=2,
)
print("Observation: SQS buffered both jobs, but it did not choose a consumer by content.")
def print_report() -> None:
print(
"\n".join(
[
"Observed architecture decision guide",
"SQS | Pull work when consumers need buffering and processing-rate control.",
]
)
)
def main() -> None:
if len(sys.argv) != 2:
raise SystemExit("Usage: python lab.py [sqs|report]")
mode = sys.argv[1]
if mode == "report":
print_report()
return
outputs = load_outputs()
actions = {"sqs": run_sqs}
if mode not in actions:
raise SystemExit("Usage: python lab.py [sqs|report]")
actions[mode](outputs)
if __name__ == "__main__":
main()
Deploy and test the queue
Synthesis proved that the stack can become a CloudFormation template. Deployment now creates the queue in your AWS account and writes its URL to the local outputs file.
This action creates a live AWS resource. The experiment sends only a handful of small messages, with cleanup included later in the project.
- Deploy MessagingDecisionLabStack by running this command:
cdk deploy --outputs-file cdk-outputs.json
What does deployment create?
- The CDK CLI submits the synthesized template through CloudFormation.
- The deployment creates DirectQueue in us-east-1.
- The outputs option writes DirectQueueUrl into cdk-outputs.json.
The first deployment can take a few minutes while CloudFormation creates the stack. A quiet terminal during this wait does not mean the deployment has stalled.
You will see a successful stack deployment summary. The cdk-outputs.json file will appear in the workspace.
Does the deployment fail?
- Confirm that your activated environment still shows .venv.
- Confirm that your configured AWS identity has permission to create CloudFormation and SQS resources.
- Check that the deployment region remains us-east-1.
Ask for help interpreting the deployment event: Help me troubleshoot this CDK deployment failure.
Before you run the experiment, do you think the queue will separate the orders according to their work values?
- Send both orders through the deployed queue by running this command:
python lab.py sqs
What does the experiment reveal?
The harness sends order-1001 and order-1002 to DirectQueue. Its pull-based consumer retrieves messages from that same queue.
The queue buffers the work without assigning messages to fulfillment or fraud-review consumers. The consumer still needs to inspect each body and make that routing decision.
You will see One SQS consumer path: received 2 message(s) followed by both order bodies. The final observation states that SQS did not choose a consumer by content.
You have your first messaging result: one queue reliably buffered two different jobs, while the routing responsibility stayed with the consumer.
Did the harness receive fewer than two messages?
- Confirm that cdk-outputs.json contains DirectQueueUrl under MessagingDecisionLabStack.
- Confirm that the terminal is using the configured AWS identity from earlier.
- Run the experiment once more if the first receive attempt ended before both messages arrived.
Ask for help with the observed output: Help me diagnose why my SQS experiment received fewer than two messages.
Your first messaging model is now backed by terminal evidence. Next, you will test a delivery model designed to give separate consumers their own copies.
Broadcast an Order with SNS
The previous experiment showed that Amazon SQS can buffer different jobs. One consumer path still had to inspect every message.
Audit needs a copy of each order notification. Analytics needs its own copy. In this step, your AWS CDK stack uses Amazon SNS to push one publication to both subscribers.
In this step, get ready to:
- Add an SNS topic with two independent SQS subscriptions.
- Publish one JSON order through the topic.
- Confirm that both subscriber queues receive the same order.
Build the fanout topology
A topic accepts one publication. Each subscription sends its own copy to a separate queue.
- In messaging_decision_lab_stack.py from earlier, locate the import group at the top of the file.
- Replace the import group by copying the expanded imports below:
from typing import Any
from aws_cdk import CfnOutput, Duration, Stack
from aws_cdk import aws_sns as sns
from aws_cdk import aws_sns_subscriptions as subscriptions
from aws_cdk import aws_sqs as sqs
from constructs import Construct
What do these imports provide?
- The sns module provides the topic construct that receives each publication.
- The subscriptions module connects each queue to the topic.
- The existing sqs module continues to create the queues that retain each subscriber copy.
- Locate direct_queue = create_queue("DirectQueue") inside MessagingDecisionLabStack.
- Add the fanout resources below that line by copying this code:
sns_audit_queue = create_queue("SnsAuditQueue")
sns_analytics_queue = create_queue("SnsAnalyticsQueue")
broadcast_topic = sns.Topic(self, "BroadcastTopic")
broadcast_topic.add_subscription(
subscriptions.SqsSubscription(
sns_audit_queue,
raw_message_delivery=True,
)
)
broadcast_topic.add_subscription(
subscriptions.SqsSubscription(
sns_analytics_queue,
raw_message_delivery=True,
)
)
How does this create fanout?
- The SnsAuditQueue construct gives the audit consumer an independent destination.
- The SnsAnalyticsQueue construct gives the analytics consumer another destination.
- The BroadcastTopic construct accepts the order publication.
- Each SqsSubscription connects one queue to the topic.
- The raw_message_delivery=True setting keeps each queue body equal to the published JSON message.
- Save messaging_decision_lab_stack.py.
- Synthesize the fanout resources by running this command:
cdk synth
What does synthesis prove?
The command converts your CDK constructs into an AWS CloudFormation template. Successful synthesis confirms that the topic plus both subscriptions form a valid infrastructure definition.
You should see the synthesized template in your terminal. It now includes an SNS topic plus two new subscriber queues.
Seeing a synthesis problem?
Check that the new imports sit above the existing SQS import. Confirm that the fanout resources use the same indentation as direct_queue.
Make sure each closing parenthesis matches its add_subscription() call.
Help me diagnose my SNS fanout synthesis problem.
The experiment needs live resource addresses after deployment. Stack outputs provide the topic ARN plus both subscriber queue URLs.
- Locate the existing DirectQueueUrl output near the bottom of messaging_decision_lab_stack.py.
- Replace the output group by copying this expanded version:
CfnOutput(self, "DirectQueueUrl", value=direct_queue.queue_url)
CfnOutput(self, "SnsAuditQueueUrl", value=sns_audit_queue.queue_url)
CfnOutput(self, "SnsAnalyticsQueueUrl", value=sns_analytics_queue.queue_url)
CfnOutput(self, "BroadcastTopicArn", value=broadcast_topic.topic_arn)
Why expose these outputs?
- The queue URL outputs tell the experiment where to retrieve each subscriber copy.
- The topic ARN output tells the publisher which topic should receive the order.
- The direct queue output preserves the SQS experiment from the previous step.
- Save messaging_decision_lab_stack.py.
- Confirm that the complete stack still synthesizes by running:
cdk synth
What does this check confirm?
Successful synthesis confirms that every output references a construct defined in the stack. The next deployment can export the three new values.
You should see the CloudFormation template finish generating without a synthesis failure.
Seeing an undefined name?
Confirm that the output lines remain inside __init__. Check each output reference against the resource definitions above it.
Help me fix an undefined construct in my CDK outputs.
✔️ Awesome, I've got everything!
Your stack file now defines one broadcast topic with two independent queue subscriptions. Make sure messaging_decision_lab_stack.py is saved.
ⓧ I'd like to double check the full code
from typing import Any
from aws_cdk import CfnOutput, Duration, Stack
from aws_cdk import aws_sns as sns
from aws_cdk import aws_sns_subscriptions as subscriptions
from aws_cdk import aws_sqs as sqs
from constructs import Construct
class MessagingDecisionLabStack(Stack):
def __init__(
self,
scope: Construct,
construct_id: str,
**kwargs: Any,
) -> None:
super().__init__(scope, construct_id, **kwargs)
def create_queue(construct_id: str) -> sqs.Queue:
return sqs.Queue(
self,
construct_id,
retention_period=Duration.days(1),
visibility_timeout=Duration.minutes(15),
)
direct_queue = create_queue("DirectQueue")
sns_audit_queue = create_queue("SnsAuditQueue")
sns_analytics_queue = create_queue("SnsAnalyticsQueue")
broadcast_topic = sns.Topic(self, "BroadcastTopic")
broadcast_topic.add_subscription(
subscriptions.SqsSubscription(
sns_audit_queue,
raw_message_delivery=True,
)
)
broadcast_topic.add_subscription(
subscriptions.SqsSubscription(
sns_analytics_queue,
raw_message_delivery=True,
)
)
CfnOutput(self, "DirectQueueUrl", value=direct_queue.queue_url)
CfnOutput(self, "SnsAuditQueueUrl", value=sns_audit_queue.queue_url)
CfnOutput(self, "SnsAnalyticsQueueUrl", value=sns_analytics_queue.queue_url)
CfnOutput(self, "BroadcastTopicArn", value=broadcast_topic.topic_arn)
How to use this reference
Compare this reference with your saved stack file. Every construct name must match because the experiment reads these outputs by name.
Add the SNS experiment
The infrastructure creates the broadcast path. The experiment needs one publisher plus two queue reads to reveal what each subscriber receives.
- In lab.py from earlier, locate the end of run_sqs().
- Add run_sns() below that function by copying this code:
def run_sns(outputs: dict[str, str]) -> None:
sns_client = boto3.client("sns", region_name=REGION)
sqs_client = boto3.client("sqs", region_name=REGION)
order = {"orderId": "order-2001", "total": 125, "event": "OrderPlaced"}
response = sns_client.publish(
TopicArn=outputs["BroadcastTopicArn"],
Message=json.dumps(order),
)
print(f"SNS message ID: {response['MessageId']}")
collect_messages(
sqs_client,
outputs["SnsAuditQueueUrl"],
"Audit subscription",
expected=1,
)
collect_messages(
sqs_client,
outputs["SnsAnalyticsQueueUrl"],
"Analytics subscription",
expected=1,
)
print("Observation: SNS pushed one publication to both subscriptions.")
What does this experiment do?
- The Boto3 SNS client publishes order-2001 to the topic exported by the stack.
- The SQS client long-polls the audit queue for one subscriber copy.
- The same client long-polls the analytics queue for its independent copy.
- The labels keep the terminal evidence separate even though both queues receive the same order.
- Locate def main() -> None: near the bottom of lab.py.
- Replace the complete main() function by copying this version:
def main() -> None:
if len(sys.argv) != 2:
raise SystemExit("Usage: python lab.py [sqs|sns|report]")
mode = sys.argv[1]
if mode == "report":
print_report()
return
outputs = load_outputs()
actions = {"sqs": run_sqs, "sns": run_sns}
if mode not in actions:
raise SystemExit("Usage: python lab.py [sqs|sns|report]")
actions[mode](outputs)
How does the dispatcher work?
- The usage text now lists sns as a supported experiment mode.
- The actions dictionary maps that mode to run_sns.
- The final function call runs the action selected by the terminal argument.
- Save lab.py.
Seeing a syntax warning?
Check the indentation inside both collect_messages() calls. Confirm that the quotes around MessageId remain inside the formatted string.
Check that both usage messages include sns. Confirm that the actions dictionary maps that mode to run_sns.
Help me find the problem in my SNS experiment.
Deploy and prove fanout
This update creates additional AWS resources. The intended project workload is expected to remain under $1. AWS charges remain usage-based.
Expect the stack update to take a few minutes while CloudFormation creates the topic plus its subscriber resources. A quiet terminal during part of that update does not mean the deployment has stalled.
- Review the security change summary if CDK requests approval.
- Confirm the deployment when prompted.
- Deploy the updated stack by running:
cdk deploy --outputs-file cdk-outputs.json
What does this deployment change?
- CloudFormation adds the broadcast topic to the existing stack.
- It creates one audit queue.
- It creates one analytics queue.
- It connects both queues to the topic as separate subscriptions.
- The deployment refreshes cdk-outputs.json with the topic ARN plus both queue URLs.
You should see the stack update complete. The refreshed outputs file now contains BroadcastTopicArn, SnsAuditQueueUrl, and SnsAnalyticsQueueUrl.
Deployment not completing?
Confirm that your .venv remains active. Check that your AWS credentials still point to the account bootstrapped in us-east-1.
Review the failed CloudFormation resource in the deployment output. Compare its construct name with messaging_decision_lab_stack.py before retrying.
Help me diagnose my CDK deployment failure for the SNS fanout resources.
The decision report should capture the behavior this experiment demonstrates. Its SNS entry turns the terminal evidence into a reusable service-selection rule.
- Locate print_report() below run_sns().
- Replace that function by copying this expanded report:
def print_report() -> None:
print(
"\n".join(
[
"Observed architecture decision guide",
"SQS | Pull work when consumers need buffering and processing-rate control.",
"SNS | Push one publication to multiple subscribed endpoints.",
]
)
)
Why update the report?
The report records the delivery behavior demonstrated by each experiment. The SNS line captures push-based fanout as an architecture decision.
- Save lab.py.
Before you run the experiment, do you expect the two consumers to compete for one message or receive separate copies?
- Publish the order through SNS by running:
python lab.py sns
What does this command prove?
The publisher sends one JSON order to BroadcastTopic. SNS pushes a separate copy to each subscribed SQS queue.
The harness long-polls both queues. Raw message delivery lets you compare the JSON bodies directly.
You should first see an SNS message ID. The audit result should report one message containing order-2001.
The analytics result should also report one message containing order-2001. The final observation confirms that one publication reached both subscriptions.
What fanout proves
SNS broadcasts one publication to every matching subscription. Each queue retains its own copy.
This experiment has no subscription filter policy. Selective delivery would require filters on the subscriptions.
Missing a subscriber message?
Confirm that the deployment refreshed cdk-outputs.json before you ran the experiment. Check that both queue URL keys appear under MessagingDecisionLabStack.
Compare both SqsSubscription constructs with the queue URL references inside run_sns().
Help me find why only one SNS subscriber received the order.
✔️ Awesome, I've got everything!
Your experiment harness publishes one order to SNS. It retrieves an independent copy from each subscriber queue.
ⓧ I'd like to double check the full code
#!/usr/bin/env python3
import json
import sys
import time
from pathlib import Path
from typing import Any
import boto3
REGION = "us-east-1"
STACK_NAME = "MessagingDecisionLabStack"
OUTPUTS_FILE = Path("cdk-outputs.json")
def load_outputs() -> dict[str, str]:
if not OUTPUTS_FILE.exists():
raise SystemExit(
"cdk-outputs.json is missing. Run cdk deploy --outputs-file cdk-outputs.json first."
)
document = json.loads(OUTPUTS_FILE.read_text(encoding="utf-8"))
return document[STACK_NAME]
def decode_body(body: str) -> Any:
try:
return json.loads(body)
except json.JSONDecodeError:
return body
def collect_messages(
sqs_client: Any,
queue_url: str,
label: str,
expected: int,
) -> list[Any]:
collected: list[Any] = []
deadline = time.monotonic() + 15
while len(collected) < expected and time.monotonic() < deadline:
response = sqs_client.receive_message(
QueueUrl=queue_url,
MaxNumberOfMessages=min(10, expected - len(collected)),
VisibilityTimeout=900,
WaitTimeSeconds=3,
)
for message in response.get("Messages", []):
collected.append(decode_body(message["Body"]))
print(f"{label}: received {len(collected)} message(s)")
for message in collected:
print(json.dumps(message, indent=2, sort_keys=True))
return collected
def run_sqs(outputs: dict[str, str]) -> None:
sqs_client = boto3.client("sqs", region_name=REGION)
orders = [
{"orderId": "order-1001", "total": 50, "work": "fulfill"},
{"orderId": "order-1002", "total": 800, "work": "fraud-review"},
]
for order in orders:
sqs_client.send_message(
QueueUrl=outputs["DirectQueueUrl"],
MessageBody=json.dumps(order),
)
collect_messages(
sqs_client,
outputs["DirectQueueUrl"],
"One SQS consumer path",
expected=2,
)
print("Observation: SQS buffered both jobs, but it did not choose a consumer by content.")
def run_sns(outputs: dict[str, str]) -> None:
sns_client = boto3.client("sns", region_name=REGION)
sqs_client = boto3.client("sqs", region_name=REGION)
order = {"orderId": "order-2001", "total": 125, "event": "OrderPlaced"}
response = sns_client.publish(
TopicArn=outputs["BroadcastTopicArn"],
Message=json.dumps(order),
)
print(f"SNS message ID: {response['MessageId']}")
collect_messages(
sqs_client,
outputs["SnsAuditQueueUrl"],
"Audit subscription",
expected=1,
)
collect_messages(
sqs_client,
outputs["SnsAnalyticsQueueUrl"],
"Analytics subscription",
expected=1,
)
print("Observation: SNS pushed one publication to both subscriptions.")
def print_report() -> None:
print(
"\n".join(
[
"Observed architecture decision guide",
"SQS | Pull work when consumers need buffering and processing-rate control.",
"SNS | Push one publication to multiple subscribed endpoints.",
]
)
)
def main() -> None:
if len(sys.argv) != 2:
raise SystemExit("Usage: python lab.py [sqs|sns|report]")
mode = sys.argv[1]
if mode == "report":
print_report()
return
outputs = load_outputs()
actions = {"sqs": run_sqs, "sns": run_sns}
if mode not in actions:
raise SystemExit("Usage: python lab.py [sqs|sns|report]")
actions[mode](outputs)
if __name__ == "__main__":
main()
How to use this reference
Compare this reference with your saved experiment file. The output keys must match the names exported by the stack.
That broadcast path now gives each subscriber its own order copy. Next, you will move routing decisions into EventBridge so each consumer receives events that match its needs.
Route Orders with EventBridge Subscribers
Your Amazon SNS experiment proved that one publication can reach two queues. Every subscriber still received the same order.
Fraud review only needs expensive orders. Amazon EventBridge Subscribers move that content decision into the routing layer.
Your AWS CDK stack will send every OrderPlaced event to fulfillment. It will send totals above 500 to fraud review.
In this step, get ready to:
- Create an enhanced EventBridge Custom Event Bus with two Amazon SQS destinations.
- Define independently filtered Subscribers with a delivery role.
- Publish two orders and observe selective delivery by order value.
Create the enhanced event bus and destinations
The enhanced Custom Event Bus retains incoming events for one day. Two queues give each Subscriber an independent buffered destination.
An AWS IAM role gives EventBridge permission to send matching events into those queues. This keeps delivery permissions separate from publisher permissions.
- In messaging_decision_lab_stack.py, replace the import section at the top with this group:
import json
from typing import Any
from aws_cdk import CfnOutput, CfnResource, Duration, Stack
from aws_cdk import aws_eventsv2 as eventsv2
from aws_cdk import aws_iam as iam
from aws_cdk import aws_sns as sns
from aws_cdk import aws_sns_subscriptions as subscriptions
from aws_cdk import aws_sqs as sqs
from constructs import Construct
What do these imports enable?
- json converts each Subscriber filter into the JSON string expected by CloudFormation.
- eventsv2 provides the enhanced Custom Event Bus construct.
- iam provides the delivery role trusted by EventBridge.
- CfnResource lets the stack emit the required Subscriber resource properties.
- In messaging_decision_lab_stack.py, find the second broadcast_topic.add_subscription() call.
- Add the destination queues and event bus below that call by pasting this code:
all_orders_queue = create_queue("AllOrdersQueue")
high_value_queue = create_queue("HighValueQueue")
orders_bus = eventsv2.CfnEventBus(
self,
"OrdersBus",
name="messaging-decision-lab",
description="Order events for the messaging decision lab",
storage_configuration=eventsv2.CfnEventBus.StorageConfigurationProperty(
retention_period_in_days=1,
),
)
What does this code do?
- AllOrdersQueue receives every order that matches the general order-event filter.
- HighValueQueue receives only orders selected by the high-value filter.
- OrdersBus creates an enhanced Custom Event Bus named messaging-decision-lab.
- retention_period_in_days=1 keeps published events available for one day.
- Save messaging_decision_lab_stack.py.
- Synthesize the updated stack by running this command:
cdk synth
What should you see?
The AWS CloudFormation template should include an AWS::EventsV2::EventBus resource. The synthesis should finish without a Python error.
Seeing an import or synthesis error?
Confirm that aws-cdk-lib==2.272.0 remains installed in the active virtual environment. An older library may not provide aws_eventsv2.
Check that the new imports remain above class MessagingDecisionLabStack.
Help me diagnose an AWS CDK synthesis error after adding the EventBridgeV2 imports and OrdersBus resource.
- In messaging_decision_lab_stack.py, place this delivery role immediately below the orders_bus definition:
delivery_role = iam.Role(
self,
"EventBridgeDeliveryRole",
assumed_by=iam.ServicePrincipal("events.amazonaws.com"),
)
all_orders_queue.grant_send_messages(delivery_role)
high_value_queue.grant_send_messages(delivery_role)
How does delivery permission work?
- EventBridgeDeliveryRole creates a role for deliveries from the event bus.
- events.amazonaws.com allows the EventBridge service to assume that role.
- grant_send_messages() grants the role permission to send messages to each target queue.
- Save messaging_decision_lab_stack.py.
- Confirm that the delivery permissions synthesize by running this command:
cdk synth
What should you see now?
The synthesized template should complete without an error. You should also see the delivery role policy near the generated queue resources.
Does the role fail to synthesize?
Check that delivery_role is defined after both target queues. Each grant_send_messages() call needs a queue that already exists.
Help me fix the EventBridge delivery role or SQS grant in my CDK stack.
Define filtered Subscribers
Each Subscriber owns its filter and destination. Producers can publish one event shape without knowing which consumers accept it.
Why use a raw CloudFormation resource?
The typed Subscriber construct in the pinned CDK library does not expose the required delivery configuration. CfnResource emits the official AWS::EventsV2::Subscriber resource with its target and role properties.
- In messaging_decision_lab_stack.py, place the first part of create_subscriber() below the delivery-role grants:
def create_subscriber(
construct_id: str,
name: str,
queue: sqs.Queue,
pattern: dict[str, Any],
state: str = "RUNNING",
resume_position: str | None = None,
) -> CfnResource:
properties: dict[str, Any] = {
"Name": name,
"EventBusArn": orders_bus.attr_event_bus_arn,
"Type": "UNORDERED",
"StartingPosition": "LATEST",
"State": state,
"FilterConfiguration": {
"Filters": [
{
"Scope": "DATA",
"Pattern": json.dumps(pattern, separators=(",", ":")),
}
]
},
"InvokeConfiguration": {
"TargetArn": queue.queue_arn,
"RoleArn": delivery_role.role_arn,
},
"Transformer": {"Type": "RAW"},
}
if resume_position is not None:
properties["ResumePosition"] = resume_position
How are Subscriber properties assembled?
- LATEST starts a new Subscriber from the newest events on the bus.
- DATA applies the pattern to the event envelope containing the order detail.
- InvokeConfiguration connects the selected queue ARN with the delivery-role ARN.
- RAW sends the event payload to the queue without an added metadata envelope.
- Select the four existing CfnOutput lines at the bottom of messaging_decision_lab_stack.py.
- Replace that selection with the helper continuation and Subscriber definitions below:
return CfnResource(
self,
construct_id,
type="AWS::EventsV2::Subscriber",
properties=properties,
)
create_subscriber(
"AllOrdersSubscriber",
"all-orders",
all_orders_queue,
{"detail-type": ["OrderPlaced"]},
)
create_subscriber(
"HighValueSubscriber",
"high-value-orders",
high_value_queue,
{"detail": {"total": [{"numeric": [">", 500]}]}},
)
CfnOutput(self, "DirectQueueUrl", value=direct_queue.queue_url)
CfnOutput(self, "SnsAuditQueueUrl", value=sns_audit_queue.queue_url)
CfnOutput(self, "SnsAnalyticsQueueUrl", value=sns_analytics_queue.queue_url)
CfnOutput(self, "AllOrdersQueueUrl", value=all_orders_queue.queue_url)
CfnOutput(self, "HighValueQueueUrl", value=high_value_queue.queue_url)
CfnOutput(self, "BroadcastTopicArn", value=broadcast_topic.topic_arn)
CfnOutput(self, "OrdersBusArn", value=orders_bus.attr_event_bus_arn)
How do the filters differ?
AllOrdersSubscriber accepts every event whose detail type is OrderPlaced. Its destination is AllOrdersQueue.
HighValueSubscriber checks detail.total with a numeric threshold above 500. Its destination is HighValueQueue.
- Save messaging_decision_lab_stack.py.
- Confirm that both Subscribers synthesize by running this command:
cdk synth
What does synthesis prove?
The template should contain two AWS::EventsV2::Subscriber resources. A successful synthesis confirms that the raw property mappings form valid CloudFormation resources.
Seeing a Subscriber synthesis error?
Check the indentation of return CfnResource(). It belongs inside create_subscriber().
Check that each filter remains a Python dictionary. json.dumps() converts it inside the helper.
Help me troubleshoot the two AWS::EventsV2::Subscriber resources in my CDK stack.
- Compare your completed messaging_decision_lab_stack.py with the full reference below.
✔️ Awesome, I've got everything!
Your stack now defines the event bus, delivery role, filtered Subscribers, target queues, and outputs.
ⓧ I'd like to double check the full code
import json
from typing import Any
from aws_cdk import CfnOutput, CfnResource, Duration, Stack
from aws_cdk import aws_eventsv2 as eventsv2
from aws_cdk import aws_iam as iam
from aws_cdk import aws_sns as sns
from aws_cdk import aws_sns_subscriptions as subscriptions
from aws_cdk import aws_sqs as sqs
from constructs import Construct
class MessagingDecisionLabStack(Stack):
def __init__(
self,
scope: Construct,
construct_id: str,
**kwargs: Any,
) -> None:
super().__init__(scope, construct_id, **kwargs)
def create_queue(construct_id: str) -> sqs.Queue:
return sqs.Queue(
self,
construct_id,
retention_period=Duration.days(1),
visibility_timeout=Duration.minutes(15),
)
direct_queue = create_queue("DirectQueue")
sns_audit_queue = create_queue("SnsAuditQueue")
sns_analytics_queue = create_queue("SnsAnalyticsQueue")
broadcast_topic = sns.Topic(self, "BroadcastTopic")
broadcast_topic.add_subscription(
subscriptions.SqsSubscription(
sns_audit_queue,
raw_message_delivery=True,
)
)
broadcast_topic.add_subscription(
subscriptions.SqsSubscription(
sns_analytics_queue,
raw_message_delivery=True,
)
)
all_orders_queue = create_queue("AllOrdersQueue")
high_value_queue = create_queue("HighValueQueue")
orders_bus = eventsv2.CfnEventBus(
self,
"OrdersBus",
name="messaging-decision-lab",
description="Order events for the messaging decision lab",
storage_configuration=eventsv2.CfnEventBus.StorageConfigurationProperty(
retention_period_in_days=1,
),
)
delivery_role = iam.Role(
self,
"EventBridgeDeliveryRole",
assumed_by=iam.ServicePrincipal("events.amazonaws.com"),
)
all_orders_queue.grant_send_messages(delivery_role)
high_value_queue.grant_send_messages(delivery_role)
def create_subscriber(
construct_id: str,
name: str,
queue: sqs.Queue,
pattern: dict[str, Any],
state: str = "RUNNING",
resume_position: str | None = None,
) -> CfnResource:
properties: dict[str, Any] = {
"Name": name,
"EventBusArn": orders_bus.attr_event_bus_arn,
"Type": "UNORDERED",
"StartingPosition": "LATEST",
"State": state,
"FilterConfiguration": {
"Filters": [
{
"Scope": "DATA",
"Pattern": json.dumps(pattern, separators=(",", ":")),
}
]
},
"InvokeConfiguration": {
"TargetArn": queue.queue_arn,
"RoleArn": delivery_role.role_arn,
},
"Transformer": {"Type": "RAW"},
}
if resume_position is not None:
properties["ResumePosition"] = resume_position
return CfnResource(
self,
construct_id,
type="AWS::EventsV2::Subscriber",
properties=properties,
)
create_subscriber(
"AllOrdersSubscriber",
"all-orders",
all_orders_queue,
{"detail-type": ["OrderPlaced"]},
)
create_subscriber(
"HighValueSubscriber",
"high-value-orders",
high_value_queue,
{"detail": {"total": [{"numeric": [">", 500]}]}},
)
CfnOutput(self, "DirectQueueUrl", value=direct_queue.queue_url)
CfnOutput(self, "SnsAuditQueueUrl", value=sns_audit_queue.queue_url)
CfnOutput(self, "SnsAnalyticsQueueUrl", value=sns_analytics_queue.queue_url)
CfnOutput(self, "AllOrdersQueueUrl", value=all_orders_queue.queue_url)
CfnOutput(self, "HighValueQueueUrl", value=high_value_queue.queue_url)
CfnOutput(self, "BroadcastTopicArn", value=broadcast_topic.topic_arn)
CfnOutput(self, "OrdersBusArn", value=orders_bus.attr_event_bus_arn)
Publish orders and verify selective delivery
The infrastructure now owns the routing decision. Boto3 will publish two orders through the enhanced EventBridge client.
One order has a total of 50. The other has a total of 800.
- In lab.py, add the publishing portion of run_eventbridge() immediately before print_report():
def run_eventbridge(outputs: dict[str, str]) -> None:
eventbridge_client = boto3.client("eventbridgev2", region_name=REGION)
sqs_client = boto3.client("sqs", region_name=REGION)
orders = [
{"orderId": "order-3001", "total": 50},
{"orderId": "order-3002", "total": 800},
]
response = eventbridge_client.put_events(
EventBusArn=outputs["OrdersBusArn"],
Entries=[
{
"Source": "com.example.orders",
"DetailType": "OrderPlaced",
"Detail": json.dumps(order),
}
for order in orders
],
)
print(json.dumps(response, indent=2, default=str))
if response["FailedEntryCount"]:
raise SystemExit("At least one EventBridge entry failed to publish.")
How are the events published?
- eventbridgev2 selects the enhanced EventBridge API client.
- Entries turns both orders into one publish request.
- DetailType carries OrderPlaced for the all-orders filter.
- Detail carries each order as a JSON payload for the numeric filter.
- FailedEntryCount stops the experiment if any event fails to publish.
- In the same function, place the queue checks immediately below the FailedEntryCount condition:
collect_messages(
sqs_client,
outputs["AllOrdersQueueUrl"],
"All-orders Subscriber",
expected=2,
)
collect_messages(
sqs_client,
outputs["HighValueQueueUrl"],
"High-value Subscriber",
expected=1,
)
print("Observation: EventBridge selected targets by event content before delivery.")
What will the queue checks prove?
The all-orders queue waits for two messages because both events have the required detail type. The high-value queue waits for one message because only order-3002 exceeds the numeric threshold.
- Save lab.py.
Before you deploy
This deployment creates additional resources in your AWS account. The experiment publishes only a handful of small messages and is expected to remain under $1 for its intended workload.
CloudFormation can take several minutes to create the bus, queues, role, and Subscribers. A quiet terminal during that period does not mean the deployment has stopped.
- Deploy the updated stack and refresh cdk-outputs.json by running this command:
cdk deploy --outputs-file cdk-outputs.json
What does this deployment change?
CloudFormation creates the enhanced event bus, two target queues, the delivery role, and both Subscribers. The refreshed outputs file gains AllOrdersQueueUrl, HighValueQueueUrl, and OrdersBusArn.
Did the deployment fail?
Confirm that your configured AWS account can create EventBridge, IAM, SQS, and CloudFormation resources in us-east-1.
If CloudFormation reports a Subscriber property problem, compare the capitalization of EventBusArn, InvokeConfiguration, TargetArn, and RoleArn.
Help me diagnose my CDK deployment failure for the enhanced EventBridge bus and Subscribers.
- In lab.py, replace the existing print_report() function with this version:
def print_report() -> None:
print(
"\n".join(
[
"Observed architecture decision guide",
"SQS | Pull work when consumers need buffering and processing-rate control.",
"SNS | Push one publication to multiple subscribed endpoints.",
"EventBridge | Route domain events to independently filtered Subscribers.",
]
)
)
Why update the report now?
The report now records the routing behavior implemented in this step. Its EventBridge rule connects domain events with independently filtered Subscribers.
- In lab.py, replace the existing main() function with this version:
def main() -> None:
if len(sys.argv) != 2:
raise SystemExit("Usage: python lab.py [sqs|sns|eventbridge|report]")
mode = sys.argv[1]
if mode == "report":
print_report()
return
outputs = load_outputs()
actions = {
"sqs": run_sqs,
"sns": run_sns,
"eventbridge": run_eventbridge,
}
if mode not in actions:
raise SystemExit("Usage: python lab.py [sqs|sns|eventbridge|report]")
actions[mode](outputs)
How is the new experiment selected?
The actions dictionary maps the eventbridge mode to run_eventbridge(). The usage message now lists that mode as a valid choice.
- Compare your completed lab.py with the full reference below.
✔️ Awesome, I've got everything!
Your experiment harness can now publish order events and inspect both filtered destinations.
ⓧ I'd like to double check the full code
#!/usr/bin/env python3
import json
import sys
import time
from pathlib import Path
from typing import Any
import boto3
REGION = "us-east-1"
STACK_NAME = "MessagingDecisionLabStack"
OUTPUTS_FILE = Path("cdk-outputs.json")
def load_outputs() -> dict[str, str]:
if not OUTPUTS_FILE.exists():
raise SystemExit(
"cdk-outputs.json is missing. Run cdk deploy --outputs-file cdk-outputs.json first."
)
document = json.loads(OUTPUTS_FILE.read_text(encoding="utf-8"))
return document[STACK_NAME]
def decode_body(body: str) -> Any:
try:
return json.loads(body)
except json.JSONDecodeError:
return body
def collect_messages(
sqs_client: Any,
queue_url: str,
label: str,
expected: int,
) -> list[Any]:
collected: list[Any] = []
deadline = time.monotonic() + 15
while len(collected) < expected and time.monotonic() < deadline:
response = sqs_client.receive_message(
QueueUrl=queue_url,
MaxNumberOfMessages=min(10, expected - len(collected)),
VisibilityTimeout=900,
WaitTimeSeconds=3,
)
for message in response.get("Messages", []):
collected.append(decode_body(message["Body"]))
print(f"{label}: received {len(collected)} message(s)")
for message in collected:
print(json.dumps(message, indent=2, sort_keys=True))
return collected
def run_sqs(outputs: dict[str, str]) -> None:
sqs_client = boto3.client("sqs", region_name=REGION)
orders = [
{"orderId": "order-1001", "total": 50, "work": "fulfill"},
{"orderId": "order-1002", "total": 800, "work": "fraud-review"},
]
for order in orders:
sqs_client.send_message(
QueueUrl=outputs["DirectQueueUrl"],
MessageBody=json.dumps(order),
)
collect_messages(
sqs_client,
outputs["DirectQueueUrl"],
"One SQS consumer path",
expected=2,
)
print("Observation: SQS buffered both jobs, but it did not choose a consumer by content.")
def run_sns(outputs: dict[str, str]) -> None:
sns_client = boto3.client("sns", region_name=REGION)
sqs_client = boto3.client("sqs", region_name=REGION)
order = {"orderId": "order-2001", "total": 125, "event": "OrderPlaced"}
response = sns_client.publish(
TopicArn=outputs["BroadcastTopicArn"],
Message=json.dumps(order),
)
print(f"SNS message ID: {response['MessageId']}")
collect_messages(
sqs_client,
outputs["SnsAuditQueueUrl"],
"Audit subscription",
expected=1,
)
collect_messages(
sqs_client,
outputs["SnsAnalyticsQueueUrl"],
"Analytics subscription",
expected=1,
)
print("Observation: SNS pushed one publication to both subscriptions.")
def run_eventbridge(outputs: dict[str, str]) -> None:
eventbridge_client = boto3.client("eventbridgev2", region_name=REGION)
sqs_client = boto3.client("sqs", region_name=REGION)
orders = [
{"orderId": "order-3001", "total": 50},
{"orderId": "order-3002", "total": 800},
]
response = eventbridge_client.put_events(
EventBusArn=outputs["OrdersBusArn"],
Entries=[
{
"Source": "com.example.orders",
"DetailType": "OrderPlaced",
"Detail": json.dumps(order),
}
for order in orders
],
)
print(json.dumps(response, indent=2, default=str))
if response["FailedEntryCount"]:
raise SystemExit("At least one EventBridge entry failed to publish.")
collect_messages(
sqs_client,
outputs["AllOrdersQueueUrl"],
"All-orders Subscriber",
expected=2,
)
collect_messages(
sqs_client,
outputs["HighValueQueueUrl"],
"High-value Subscriber",
expected=1,
)
print("Observation: EventBridge selected targets by event content before delivery.")
def print_report() -> None:
print(
"\n".join(
[
"Observed architecture decision guide",
"SQS | Pull work when consumers need buffering and processing-rate control.",
"SNS | Push one publication to multiple subscribed endpoints.",
"EventBridge | Route domain events to independently filtered Subscribers.",
]
)
)
def main() -> None:
if len(sys.argv) != 2:
raise SystemExit("Usage: python lab.py [sqs|sns|eventbridge|report]")
mode = sys.argv[1]
if mode == "report":
print_report()
return
outputs = load_outputs()
actions = {
"sqs": run_sqs,
"sns": run_sns,
"eventbridge": run_eventbridge,
}
if mode not in actions:
raise SystemExit("Usage: python lab.py [sqs|sns|eventbridge|report]")
actions[mode](outputs)
if __name__ == "__main__":
main()
- Save lab.py.
Before you run this ask yourself which Subscriber receives one order.
- Run the EventBridge experiment with this command:
python lab.py eventbridge
What should you see?
The publish response should report no failed entries. All-orders Subscriber should report two messages.
High-value Subscriber should report one message containing order-3002 with a total of 800. The final observation confirms that EventBridge selected targets before delivery.
Are the delivery counts different?
Confirm that cdk-outputs.json came from the latest deployment. An older file does not contain the new bus and queue outputs.
Check that the all-orders filter uses detail-type with OrderPlaced. Check that the high-value filter uses detail.total with the numeric threshold.
Help me troubleshoot unexpected EventBridge Subscriber delivery counts in this messaging lab.
That routing split is working. Your bus now sends every order toward fulfillment while isolating high-value work for fraud review.
You now have three messaging behaviors you can demonstrate from evidence. Next up, you will turn those observations into a reusable architecture decision.
Create the Architecture Decision Report
Your Amazon EventBridge experiment now routes each order to the correct Subscriber based on its content. You have terminal evidence for buffering, fanout, and selective routing.
A working demo becomes engineering judgment when you can turn observed behavior into a repeatable decision rule. This step converts your Amazon SQS, Amazon SNS, and EventBridge evidence into an architecture decision report.
In this step, get ready to:
- Map your AWS CDK constructs to synthesized AWS CloudFormation resources.
- Complete the report mode in lab.py.
- Apply the report to three messaging scenarios.
Synthesize the final template
AWS CDK turns the constructs in your Python stack into an AWS CloudFormation template. Inspecting that template connects each abstraction to the infrastructure AWS deploys.
- Return to the terminal in your messaging-decision-lab folder from earlier.
- Synthesize the final application into template.yaml by running this command:
cdk synth > template.yaml
What does this command do?
- The AWS CDK reads the application configured in cdk.json.
- The shell writes the synthesized AWS CloudFormation template into template.yaml.
- The deployed resources remain unchanged because synthesis only generates a local template.
Does synthesis stop with an error?
- Confirm the .venv environment from earlier is active.
- Confirm the packages listed in requirements.txt are installed in that environment.
- Check that your terminal is inside the messaging-decision-lab folder.
Help me diagnose why my AWS CDK application does not synthesize successfully.
Before you search, predict which messaging resource families the synthesized template should reveal.
- Search template.yaml for the messaging resource types by running this command:
grep -E 'AWS::(SQS::Queue|SNS::Topic|EventsV2::EventBus|EventsV2::Subscriber)' template.yaml
What should I see?
- Lines containing AWS::SQS::Queue connect the queue constructs to their generated resources.
- A line containing AWS::SNS::Topic identifies the broadcast topic.
- A line containing AWS::EventsV2::EventBus identifies the enhanced Custom Event Bus.
- Lines containing AWS::EventsV2::Subscriber identify the two filtered Subscribers.
Does the search return no matches?
- Confirm that synthesis created template.yaml inside the messaging-decision-lab folder.
- Confirm that the search command uses the uppercase resource names shown above.
- Run the synthesis command again after saving your project files.
Help me find the expected messaging resource types in my synthesized AWS CloudFormation template.
You have exposed the infrastructure beneath your Python constructs. The template now provides concrete evidence of what the AWS CDK manages for you.
Complete the decision harness
The finished harness needs a focused queue check for a retained high-value event. Its command interface also needs to expose that check alongside the three experiments and the report.
- In lab.py, find the run_eventbridge() function.
- Add the focused queue helper directly below run_eventbridge() by copying this code:
def check_high_value(outputs: dict[str, str]) -> None:
sqs_client = boto3.client("sqs", region_name=REGION)
collect_messages(
sqs_client,
outputs["HighValueQueueUrl"],
"Recovered high-value event",
expected=1,
)
What does this helper do?
- The helper creates an SQS client in the lab region.
- It reads from the queue identified by HighValueQueueUrl.
- It reuses collect_messages() to retrieve one high-value event.
- In lab.py, find the main() function.
- Replace the complete main() function plus its final entry point with this version:
def main() -> None:
if len(sys.argv) != 2:
raise SystemExit("Usage: python lab.py [sqs|sns|eventbridge|check-high|report]")
mode = sys.argv[1]
if mode == "report":
print_report()
return
outputs = load_outputs()
actions = {
"sqs": run_sqs,
"sns": run_sns,
"eventbridge": run_eventbridge,
"check-high": check_high_value,
}
if mode not in actions:
raise SystemExit("Usage: python lab.py [sqs|sns|eventbridge|check-high|report]")
actions[mode](outputs)
if __name__ == "__main__":
main()
How does the command routing work?
- The usage message lists every mode supported by the finished harness.
- The report mode prints the decision guide without loading stack outputs.
- The actions dictionary maps each experiment mode to its Python function.
- The check-high mode maps to check_high_value().
Before you run this check, predict whether the existing report mode still starts after the command routing change.
- Save lab.py.
- Verify that the updated harness still parses and starts by running:
python lab.py report
What should I see?
You should see the existing decision guide with rows for SQS, SNS, and EventBridge. This confirms that Python parsed the new helper plus the updated command routing.
Does the report fail to start?
- Check that check_high_value() appears above main().
- Check that the actions dictionary maps check-high to check_high_value.
- Check the indentation inside main() against the reference below.
Help me diagnose why the report mode stopped working after I added the focused queue check.
Each service row captures one behavior you observed. The report also needs a composition rule because production systems often route an event before buffering or broadcasting its delivery.
- In lab.py, find the print_report() function.
- Replace the complete function with this final version:
def print_report() -> None:
print(
"\n".join(
[
"Observed architecture decision guide",
"SQS | Pull work when consumers need buffering and processing-rate control.",
"SNS | Push one publication to multiple subscribed endpoints.",
"EventBridge | Route domain events to independently filtered Subscribers.",
"Composition | Route with EventBridge, buffer worker-bound delivery with SQS, and use SNS for broadcast endpoints.",
]
)
)
What does the report capture?
- The SQS row captures pull-based buffering plus consumer processing-rate control.
- The SNS row captures push-based fanout to multiple subscribed endpoints.
- The EventBridge row captures routing through independent Subscriber filters.
- The Composition row shows how routing, buffering, and broadcast can work together.
✔️ Awesome, I've got everything!
Your command modes now match the finished harness. Save lab.py before the final check.
ⓧ I'd like to double check the full code
#!/usr/bin/env python3
import json
import sys
import time
from pathlib import Path
from typing import Any
import boto3
REGION = "us-east-1"
STACK_NAME = "MessagingDecisionLabStack"
OUTPUTS_FILE = Path("cdk-outputs.json")
def load_outputs() -> dict[str, str]:
if not OUTPUTS_FILE.exists():
raise SystemExit(
"cdk-outputs.json is missing. Run cdk deploy --outputs-file cdk-outputs.json first."
)
document = json.loads(OUTPUTS_FILE.read_text(encoding="utf-8"))
return document[STACK_NAME]
def decode_body(body: str) -> Any:
try:
return json.loads(body)
except json.JSONDecodeError:
return body
def collect_messages(
sqs_client: Any,
queue_url: str,
label: str,
expected: int,
) -> list[Any]:
collected: list[Any] = []
deadline = time.monotonic() + 15
while len(collected) < expected and time.monotonic() < deadline:
response = sqs_client.receive_message(
QueueUrl=queue_url,
MaxNumberOfMessages=min(10, expected - len(collected)),
VisibilityTimeout=900,
WaitTimeSeconds=3,
)
for message in response.get("Messages", []):
collected.append(decode_body(message["Body"]))
print(f"{label}: received {len(collected)} message(s)")
for message in collected:
print(json.dumps(message, indent=2, sort_keys=True))
return collected
def run_sqs(outputs: dict[str, str]) -> None:
sqs_client = boto3.client("sqs", region_name=REGION)
orders = [
{"orderId": "order-1001", "total": 50, "work": "fulfill"},
{"orderId": "order-1002", "total": 800, "work": "fraud-review"},
]
for order in orders:
sqs_client.send_message(
QueueUrl=outputs["DirectQueueUrl"],
MessageBody=json.dumps(order),
)
collect_messages(
sqs_client,
outputs["DirectQueueUrl"],
"One SQS consumer path",
expected=2,
)
print("Observation: SQS buffered both jobs, but it did not choose a consumer by content.")
def run_sns(outputs: dict[str, str]) -> None:
sns_client = boto3.client("sns", region_name=REGION)
sqs_client = boto3.client("sqs", region_name=REGION)
order = {"orderId": "order-2001", "total": 125, "event": "OrderPlaced"}
response = sns_client.publish(
TopicArn=outputs["BroadcastTopicArn"],
Message=json.dumps(order),
)
print(f"SNS message ID: {response['MessageId']}")
collect_messages(
sqs_client,
outputs["SnsAuditQueueUrl"],
"Audit subscription",
expected=1,
)
collect_messages(
sqs_client,
outputs["SnsAnalyticsQueueUrl"],
"Analytics subscription",
expected=1,
)
print("Observation: SNS pushed one publication to both subscriptions.")
def run_eventbridge(outputs: dict[str, str]) -> None:
eventbridge_client = boto3.client("eventbridgev2", region_name=REGION)
sqs_client = boto3.client("sqs", region_name=REGION)
orders = [
{"orderId": "order-3001", "total": 50},
{"orderId": "order-3002", "total": 800},
]
response = eventbridge_client.put_events(
EventBusArn=outputs["OrdersBusArn"],
Entries=[
{
"Source": "com.example.orders",
"DetailType": "OrderPlaced",
"Detail": json.dumps(order),
}
for order in orders
],
)
print(json.dumps(response, indent=2, default=str))
if response["FailedEntryCount"]:
raise SystemExit("At least one EventBridge entry failed to publish.")
collect_messages(
sqs_client,
outputs["AllOrdersQueueUrl"],
"All-orders Subscriber",
expected=2,
)
collect_messages(
sqs_client,
outputs["HighValueQueueUrl"],
"High-value Subscriber",
expected=1,
)
print("Observation: EventBridge selected targets by event content before delivery.")
def check_high_value(outputs: dict[str, str]) -> None:
sqs_client = boto3.client("sqs", region_name=REGION)
collect_messages(
sqs_client,
outputs["HighValueQueueUrl"],
"Recovered high-value event",
expected=1,
)
def print_report() -> None:
print(
"\n".join(
[
"Observed architecture decision guide",
"SQS | Pull work when consumers need buffering and processing-rate control.",
"SNS | Push one publication to multiple subscribed endpoints.",
"EventBridge | Route domain events to independently filtered Subscribers.",
"Composition | Route with EventBridge, buffer worker-bound delivery with SQS, and use SNS for broadcast endpoints.",
]
)
)
def main() -> None:
if len(sys.argv) != 2:
raise SystemExit("Usage: python lab.py [sqs|sns|eventbridge|check-high|report]")
mode = sys.argv[1]
if mode == "report":
print_report()
return
outputs = load_outputs()
actions = {
"sqs": run_sqs,
"sns": run_sns,
"eventbridge": run_eventbridge,
"check-high": check_high_value,
}
if mode not in actions:
raise SystemExit("Usage: python lab.py [sqs|sns|eventbridge|check-high|report]")
actions[mode](outputs)
if __name__ == "__main__":
main()
How to use this reference
This reference shows the exact finished lab.py file. It includes every experiment mode plus the final decision report.
Test the architecture decision
The final report turns each experiment into a selection rule. You can now choose a messaging service from the behavior the workload requires.
Before you run the final check, predict which report row covers a content-routed event that still needs worker buffering.
- Save lab.py.
- Print the completed architecture decision report by running:
python lab.py report
What should I see?
- You should see Observed architecture decision guide at the top.
- You should see separate rows for SQS, SNS, and EventBridge.
- You should see a final Composition row that combines routing, buffering, and broadcast.
Is the Composition row missing?
- Confirm that you saved lab.py after replacing print_report().
- Confirm that the Composition string sits inside the list passed to join().
- Confirm that your terminal is running the lab.py file inside messaging-decision-lab.
Help me find why my architecture decision report does not show the Composition row.
Strong work. Your terminal now explains when each messaging model belongs in an architecture instead of leaving the choice to guesswork.
- Choose the service for one job pulled by one worker pool.
- Choose the service for one notification delivered to many endpoints.
- Choose the service for one domain event routed according to its content.
Check your decisions
- SQS fits one worker pool because consumers pull buffered work at their own processing rate.
- SNS fits one notification sent to many endpoints because each subscription receives the publication.
- EventBridge fits a domain event routed by content because Subscriber filters select the delivery path.
- Composition fits filtered events that still need SQS buffering or SNS broadcast after routing.
Secret mission
Pause and Resume a Subscriber
Test what happens when a Subscriber pauses while orders keep arriving. You will stop high-value delivery, publish two events, resume from the last processed position, and recover the retained high-value event.
Clean Up Your Resources
Clean Up Your Resources
Choose whether to keep your lab active, pause message delivery, or delete the deployment. AWS usage-based charges should stay under $1 for the intended workload.
Cost warning
The lab sends only a handful of small messages. Additional publishing or unrelated activity in the same account can increase your charges.
Choose Pause or Delete if you have finished experimenting. Pause limits new delivery activity. Delete removes the project resources.
Resources you used:
- The MessagingDecisionLabStack application stack in AWS CloudFormation.
- The messaging-decision-lab Custom Event Bus in Amazon EventBridge.
- The AllOrdersSubscriber and HighValueSubscriber EventBridge Subscribers.
- The EventBridge delivery role in AWS Identity and Access Management (IAM).
- Five Amazon SQS queues for direct work plus subscriber delivery.
- One Amazon SNS topic with two SQS subscriptions.
- The local messaging-decision-lab folder containing your AWS CDK source files plus generated artifacts.
Keep everything running
No action is needed. Choose this if you want to keep demonstrating the messaging experiments or extend the architecture.
- The MessagingDecisionLabStack remains deployed in us-east-1.
- Both EventBridge Subscribers remain in the RUNNING state.
- The OrdersBus continues to retain events for one day.
- Your local source files plus the synthesized template remain available for future experiments.
- Run the experiments only when you need fresh evidence.
Keep Test Traffic Intentional
The intended project workload stays small. Repeated publishing creates additional usage-based activity across the messaging services.
Pause - I'll come back to this later
A paused lab keeps the cloud resources while stopping new Subscriber deliveries. The OrdersBus continues to retain events for one day.
The helper's default state controls AllOrdersSubscriber. The explicit state in the second call controls HighValueSubscriber.
- Return to messaging_decision_lab_stack.py in VS Code.
- Change state: str = "RUNNING" to state: str = "STOPPED" in create_subscriber().
- Change state="RUNNING" to state="STOPPED" in the HighValueSubscriber call.
- Leave resume_position="LAST_PROCESSED" unchanged.
- Save messaging_decision_lab_stack.py.
- Deploy the paused Subscriber state by running:
cdk deploy --outputs-file cdk-outputs.json
What Does This Deployment Change?
- The default state change pauses delivery through AllOrdersSubscriber.
- The explicit state change pauses delivery through HighValueSubscriber.
- The LAST_PROCESSED resume position preserves the recovery behavior you tested.
The terminal reports a successful stack update. Your lab is now parked without deleting its architecture.
Deployment Did Not Finish?
- Confirm the terminal is inside the messaging-decision-lab folder.
- Confirm your configured AWS credentials still point to the account that owns the stack.
Help me diagnose why my CDK deployment did not update both EventBridge Subscribers.
When you return, restore both Subscribers before publishing more events.
- Change the default state in create_subscriber() back to RUNNING.
- Change the explicit HighValueSubscriber state back to RUNNING.
- Leave resume_position="LAST_PROCESSED" unchanged.
- Redeploy the restored state with the same deployment command shown above.
Remember the Retention Window
Only events still inside the one-day retention window are available when delivery resumes. Older retained events expire from the bus.
Delete - I don't want to use this again
Deleting the stack permanently removes the deployed lab resources. Your local source stays untouched until the separate folder cleanup.
CDK asks for confirmation before removing the stack. This gives you one final checkpoint before cloud deletion begins.
- Return to the terminal in the messaging-decision-lab folder.
- Destroy the application stack by running:
cdk destroy MessagingDecisionLabStack
What Does This Command Delete?
The command deletes MessagingDecisionLabStack. CloudFormation removes the EventBridge bus, both Subscribers, the delivery role, all five queues, the SNS topic, and both SNS subscriptions with it.
The shared CDKToolkit stack remains intact. Keep it when other CDK projects may depend on its S3 or ECR assets.
- Confirm the deletion when CDK prompts you.
The terminal reports successful stack deletion when the cloud cleanup finishes. That closes the deployed side of your lab cleanly.
Stack Deletion Failed?
- Confirm your configured AWS credentials still point to the account that owns MessagingDecisionLabStack.
- Review the failed resource named in the terminal before retrying the destroy command.
Help me diagnose why MessagingDecisionLabStack did not delete.
Local Cleanup Is Permanent
The next cleanup removes your source files, virtual environment, synthesized template, and deployment outputs from the local project folder.
- Copy any terminal evidence or source files you want to keep before continuing.
- Return to the terminal inside messaging-decision-lab.
- Remove the local project folder by running these commands:
cd ..
rm -rf messaging-decision-lab
What Do These Commands Remove?
- The cd .. command moves the terminal into the parent directory.
- The rm -rf messaging-decision-lab command permanently deletes the project folder plus everything inside it.
- Check the parent directory in your Linux file manager for the messaging-decision-lab folder.
You should no longer see the project folder. Your lab is now removed from AWS and your machine.
Nice Work!
Nice Work!
You made it! Your Python AWS CDK messaging lab now provides real evidence for choosing between three AWS messaging models.
You've learned how to:
- Build an infrastructure as code application that synthesizes into an AWS CloudFormation template. You can connect Python constructs to the AWS resources they generate.
- Demonstrate Amazon SQS buffering by sending different jobs through one pull-based queue. You observed that the consumer still owns the content-based processing decision.
- Compare Amazon SNS fanout with Amazon EventBridge routing. You proved that SNS broadcasts one publication to independent subscribers while EventBridge selects Subscribers using event content.
- Secret Mission: Pause the high-value Subscriber. Publish while delivery is stopped. Resume from LAST_PROCESSED to recover the event retained by the enhanced Custom Event Bus.
Ready to quiz yourself?