AWS Lambda异步任务异常:DeepLens图像捕获与S3上传线程故障
问题诊断与修复方案
首先,你的上传线程无法工作的核心原因有两个:未传递必要变量导致NameError,以及多线程操作全局列表时缺少同步锁,另外当前的线程执行逻辑其实并没有真正实现并行。下面一步步拆解并修复:
1. 直接致命错误:upload_image函数未获取client和iot_topic变量
看你的upload_image函数代码:
def upload_image(): client.publish(topic=iot_topic, payload='upload_image') ...
这里的client和iot_topic是在init_greengrass函数中创建的局部变量,upload_image并没有访问权限。当上传线程启动后,会立即抛出NameError: name 'client' is not defined异常,导致线程直接崩溃——这就是你看到上传线程"无法工作"的直接原因。
2. 多线程共享全局列表的竞争条件
liste_of_frame是全局列表,在捕获线程中执行append,上传线程中执行pop(i),这种无保护的多线程操作会导致列表状态不一致:比如循环遍历列表时长度突然变化,抛出IndexError,或者丢失元素。
3. 当前线程逻辑并未真正并行
在你的主循环中:
t1.start() t2.start() t1.join() t2.join()
join()方法会阻塞当前线程直到目标线程完成,所以实际执行顺序是先等捕获线程完成,再执行上传线程,完全是串行执行,并没有实现你想要的并行效果。
修复后的代码实现
下面是修复后的完整代码,同时优化了线程模型(让上传线程持续运行,而非每次循环创建新线程):
from threading import Thread, Event, Lock, Timer import os import json import numpy as np import awscam import cv2 import greengrasssdk import boto3 from botocore.session import Session import time # 全局变量与线程同步工具 liste_of_frame = [] frame_lock = Lock() # 保护列表的线程锁,避免多线程竞争 stop_event = Event() # 用于优雅停止上传线程 def write_image_to_s3(img): session = Session() s3 = session.create_client('s3') file_name = 'DeepLens/image-'+time.strftime("%Y%m%d-%H%M%S")+'.jpg' encode_param=[int(cv2.IMWRITE_JPEG_QUALITY),100] _, jpg_data = cv2.imencode('.jpg', img, encode_param) response = s3.put_object(ACL='public-read', Body=jpg_data.tostring(),Bucket='deeplens-sagemaker-f2e',Key=file_name) image_url = 'https://s3.amazonaws.com/deeplens-sagemaker-f2e/'+file_name return image_url def upload_image(client, iot_topic): # 持续运行直到收到停止信号 while not stop_event.is_set(): try: client.publish(topic=iot_topic, payload='Checking for images to upload') # 获取锁后安全操作列表 with frame_lock: if len(liste_of_frame) > 0: # 弹出列表第一个元素(避免索引越界问题) img = liste_of_frame.pop(0) write_image_to_s3(img) client.publish(topic=iot_topic, payload='Uploaded one image') # 短暂休眠,避免占用过多系统资源 time.sleep(0.5) except Exception as e: client.publish(topic=iot_topic, payload=f'Upload error: {str(e)}') time.sleep(1) def capture_img(model_type, output_map, client, iot_topic, local_display, model, detection_threshold, input_height, input_width): ret, frame = awscam.getLastFrame() if not ret: raise Exception('Failed to get frame from the stream') frame_resize = cv2.resize(frame, (input_height, input_width)) parsed_inference_results = model.parseResult(model_type, model.doInference(frame_resize)) yscale = float(frame.shape[0]/input_height) xscale = float(frame.shape[1]/input_width) cloud_output = {} for obj in parsed_inference_results[model_type]: if obj['prob'] > detection_threshold: cloud_output[output_map[obj['label']]] = obj['prob'] # 获取锁后安全添加图像到列表 with frame_lock: liste_of_frame.append(frame) local_display.set_frame_data(frame) client.publish(topic=iot_topic, payload=json.dumps(cloud_output)) def init_greengrass(): model_type = 'ssd' output_map = {1: 'face'} client = greengrasssdk.client('iot-data') iot_topic = '$aws/things/{}/infer'.format(os.environ['AWS_IOT_THING_NAME']) local_display = LocalDisplay('480p') local_display.start() model_path = '/opt/awscam/artifacts/mxnet_deploy_ssd_FP16_FUSED.xml' client.publish(topic=iot_topic, payload='Loading face detection model') model = awscam.Model(model_path, {'GPU': 1}) client.publish(topic=iot_topic, payload='Face detection model loaded') detection_threshold = 0.5 input_height = 300 input_width = 300 return model_type, output_map, client, iot_topic, local_display, model, detection_threshold, input_height, input_width class LocalDisplay(Thread): def __init__(self, resolution): super(LocalDisplay, self).__init__() RESOLUTION = {'1080p' : (1920, 1080), '720p' : (1280, 720), '480p' : (858, 480)} if resolution not in RESOLUTION: raise Exception("Invalid resolution") self.resolution = RESOLUTION[resolution] self.frame = cv2.imencode('.jpg', 255*np.ones([640, 480, 3]))[1] self.stop_request = Event() def run(self): result_path = '/tmp/results.mjpeg' if not os.path.exists(result_path): os.mkfifo(result_path) with open(result_path, 'w') as fifo_file: while not self.stop_request.isSet(): try: fifo_file.write(self.frame.tobytes()) except IOError: continue def set_frame_data(self, frame): ret, jpeg = cv2.imencode('.jpg', cv2.resize(frame, self.resolution)) if not ret: raise Exception('Failed to set frame data') self.frame = jpeg def join(self): self.stop_request.set() def greengrass_infinite_infer_run(): try: model_type, output_map, client, iot_topic, local_display, model, detection_threshold, input_height, input_width = init_greengrass() # 启动持续运行的上传线程,传入必要参数 upload_thread = Thread(target=upload_image, args=(client, iot_topic)) upload_thread.daemon = True # 设置为守护线程,主进程退出时自动终止 upload_thread.start() # 持续执行图像捕获逻辑 while True: capture_img(model_type, output_map, client, iot_topic, local_display, model, detection_threshold, input_height, input_width) time.sleep(0.1) # 控制捕获频率,避免资源占用过高 except Exception as ex: client.publish(topic=iot_topic, payload='Error in face detection lambda: {}'.format(ex)) # 发送停止信号给上传线程 stop_event.set() greengrass_infinite_infer_run()
关键修复点说明
- 传递必要参数:将
client和iot_topic作为参数传入upload_image函数,解决NameError问题。 - 线程锁保护全局列表:使用
Lock()确保同一时间只有一个线程操作liste_of_frame,避免竞争条件导致的异常。 - 优化线程模型:创建一个持续运行的上传守护线程,而非每次循环创建新线程,真正实现捕获与上传的并行。
- 避免索引错误:上传时使用
pop(0)而非按索引pop,彻底避免列表长度变化导致的索引越界问题。
不需要更换线程库,threading完全可以满足你的需求,问题主要出在代码逻辑的错误上。
内容的提问来源于stack exchange,提问作者Matthieu Rousseau
相关产品推荐
相关产品推荐

