如何从AWS API Gateway获取客户端WebSocket的响应结果?
架构与问题背景
我的应用架构如下:
- Web应用:运行在EC2上,提供照片上传功能;
- 若干人脸识别摄像头:作为WebSocket客户端连接至API Gateway;
- Lambda #1:作为API Gateway的
$connect路由实现,将WebSocket连接ID存储至DynamoDB; - Lambda #2:独立Lambda函数,当Web应用上传照片时触发,从DynamoDB获取连接ID,调用摄像头的WebSocket API上传照片。
遇到的问题:Lambda #2调用post_to_connection后,返回的ResponseMetadata中HTTPStatusCode为200,但content-length为0,无法获取摄像头的响应。已确认摄像头确实收到了消息,且用wscat作为WebSocket服务器时能正常接收摄像头响应。
原因分析
API Gateway WebSocket的post_to_connection接口仅负责单向推送消息到客户端,它返回的200响应只是API Gateway确认消息已成功投递,并不会转发客户端的回复给调用该接口的Lambda。客户端(摄像头)的响应会进入API Gateway的路由处理流程,需要额外配置才能传递到后端服务。
解决方案:基于API Gateway路由的回调机制
步骤1:配置API Gateway自定义路由
在WebSocket API中添加一个自定义路由(比如camera-response),将其集成目标设置为新的Lambda函数(Lambda #3),用于接收摄像头的响应消息。
步骤2:Lambda #2发送消息时携带唯一请求ID
在Lambda #2中发送消息给摄像头时,生成一个唯一requestId,并将该ID与业务上下文(如照片ID、连接ID)存储到DynamoDB,用于后续关联响应:
import boto3 import json import uuid apigw_client = boto3.client('apigatewaymanagementapi', endpoint_url="https://{api-id}.execute-api.{region}.amazonaws.com/{stage}") ddb_resource = boto3.resource('dynamodb') requests_table = ddb_resource.Table('WebSocketRequests') def lambda_handler(event, context): # 从DynamoDB获取摄像头连接ID connection_id = "your-camera-connection-id" # 生成唯一请求ID request_id = str(uuid.uuid4()) photo_id = event['photoId'] photo_url = event['photoUrl'] # 存储请求上下文 requests_table.put_item( Item={ 'requestId': request_id, 'connectionId': connection_id, 'photoId': photo_id, 'status': 'pending' } ) # 推送消息到摄像头 apigw_client.post_to_connection( ConnectionId=connection_id, Data=json.dumps({ 'action': 'process-photo', 'requestId': request_id, 'photoUrl': photo_url }) ) return {'statusCode': 200}
步骤3:摄像头携带请求ID回复消息
摄像头收到消息后,处理完成后将结果携带requestId,通过WebSocket发送到camera-response路由,示例消息格式:
{ "action": "response", "requestId": "your-unique-request-id", "result": "success", "faceMatches": [{"id": "123", "confidence": 0.95}] }
步骤4:Lambda #3处理响应并关联业务逻辑
Lambda #3接收API Gateway传递的响应消息,根据requestId从DynamoDB获取业务上下文,再将响应传递给Lambda #2处理(或直接执行业务逻辑):
import boto3 import json ddb_resource = boto3.resource('dynamodb') requests_table = ddb_resource.Table('WebSocketRequests') lambda_client = boto3.client('lambda') def lambda_handler(event, context): # 解析客户端响应 message = json.loads(event['body']) request_id = message.get('requestId') camera_response = message.get('result') # 获取请求上下文 request_item = requests_table.get_item(Key={'requestId': request_id})['Item'] photo_id = request_item['photoId'] # 触发Lambda #2处理响应(异步调用) lambda_client.invoke( FunctionName='Lambda-2-Function-Name', InvocationType='Event', Payload=json.dumps({ 'action': 'handle-camera-response', 'photoId': photo_id, 'requestId': request_id, 'response': camera_response }) ) # 更新请求状态 requests_table.update_item( Key={'requestId': request_id}, UpdateExpression='SET #status = :completed', ExpressionAttributeNames={'#status': 'status'}, ExpressionAttributeValues={':completed': 'completed'} ) return {'statusCode': 200}
替代方案:使用AppSync WebSocket(可选)
如果希望更便捷地处理双向通信,可以考虑改用AWS AppSync的WebSocket功能,它内置了订阅/发布机制,支持客户端与服务端的双向响应,无需手动配置路由和消息中转逻辑。
内容的提问来源于stack exchange,提问作者Cyron

