如何通过触发式Lambda函数处理MQTT消息并发送至另一AWS IoT主题?
Got it, let's walk through exactly how to build this Lambda function to generate a custom message and publish it to your target topic. Since the built-in "Republish" action only forwards the original message, using Lambda gives you full control to transform or create entirely new content.
Step 1: Set Up Lambda IAM Permissions
First, your Lambda function needs permission to publish messages to AWS IoT Core. You'll need to add an inline policy to your Lambda's execution role:
- Go to the IAM console, find your Lambda's execution role.
- Add an inline policy with the following JSON (replace
REGIONandACCOUNT_IDwith your actual values):
{ "Version": "2012-10-17", "Statement": [ { "Effect": "Allow", "Action": "iot:Publish", "Resource": "arn:aws:iot:REGION:ACCOUNT_ID:topic/terminal2/test" } ] }
This policy explicitly grants the Lambda permission to send messages to your terminal2/test topic.
Step 2: Write the Lambda Function Code
Below is a Python example (the most common choice for AWS Lambda) that handles the incoming IoT message, generates a custom new message, and publishes it to your target topic. I've included comments to explain each part:
import boto3 import json import base64 # Initialize the AWS IoT Data Plane client iot_data_client = boto3.client('iot-data') def lambda_handler(event, context): # Parse the incoming message from the IoT rule # Note: IoT rules often encode the payload as base64 by default try: # Decode base64 payload and convert to JSON raw_payload = base64.b64decode(event['payload']).decode('utf-8') original_message = json.loads(raw_payload) except KeyError: # Fallback if payload is already decoded (adjust based on your rule config) original_message = event.get('payload', {}) except json.JSONDecodeError: # Handle cases where original message isn't valid JSON original_message = {"raw_message": raw_payload} # -------------------------- # Customize this section to generate your new message # -------------------------- new_message = { "original_content": original_message, "processed_at": context.get_remaining_time_in_millis(), "message_type": "transformed", "custom_metadata": { "source_topic": "terminal1/", "status": "success" } } # Publish the new message to terminal2/test try: publish_response = iot_data_client.publish( topic='terminal2/test', qos=1, # Use 0 for fire-and-forget, 1 for guaranteed delivery payload=json.dumps(new_message) ) print(f"Message published successfully: {publish_response}") return { "statusCode": 200, "body": json.dumps("Custom message published to terminal2/test") } except Exception as e: error_msg = f"Failed to publish message: {str(e)}" print(error_msg) return { "statusCode": 500, "body": json.dumps(error_msg) }
Key Notes About the Code:
- Payload Decoding: AWS IoT rules typically pass the message payload as a base64-encoded string, so we handle that decoding first. If your rule is configured to pass raw JSON, you can skip the base64 step.
- Custom Message Logic: The
new_messagedictionary is where you'll define your custom content—you can modify this to include any data derived from the original message, static values, or external data. - QoS Setting: Adjust the
qosparameter based on your delivery needs:0for best-effort,1for at-least-once delivery.
Step 3: Update Your AWS IoT Rule
Make sure your existing IoT rule is configured to trigger this Lambda function instead of the republish action:
- Go to the AWS IoT Core console, open your rule.
- Under "Actions", remove the existing republish action (if present).
- Add a new action: Invoke a Lambda function.
- Select your newly created Lambda function from the dropdown list.
- Save the rule.
Step 4: Test the Workflow
- Use an MQTT client (like the AWS IoT Core Test Console, or a tool like Mosquitto) to publish a message to
terminal1/. - Subscribe to the
terminal2/testtopic in the same client. - You should see your custom generated message appear in the subscription feed.
内容的提问来源于stack exchange,提问作者Senura Dissanayake

