You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()

关键修复点说明

  1. 传递必要参数:将client和iot_topic作为参数传入upload_image函数,解决NameError问题。
  2. 线程锁保护全局列表:使用Lock()确保同一时间只有一个线程操作liste_of_frame,避免竞争条件导致的异常。
  3. 优化线程模型:创建一个持续运行的上传守护线程,而非每次循环创建新线程,真正实现捕获与上传的并行。
  4. 避免索引错误:上传时使用pop(0)而非按索引pop,彻底避免列表长度变化导致的索引越界问题。

不需要更换线程库,threading完全可以满足你的需求,问题主要出在代码逻辑的错误上。

内容的提问来源于stack exchange,提问作者Matthieu Rousseau

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 06:42:30