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

Python中如何定期合并双线程生成的音视频文件?

问题描述

我想要同时并定期将音视频保存为固定时长的MP4文件。根据Stack Overflow上的相关帖子,我知道需要启动两个独立线程并行录制视频和音频。我已修改代码如下,目前可将音视频同步保存为固定时长的分离文件。我的问题是:如何在两个线程生成音视频文件后即时合并它们?我可以等线程结束后再合并,但这不是我想要的。

现有代码
#!/usr/bin/env python
# -*- coding: utf-8 -*-
# VideoRecorder.py

from __future__ import print_function, division
import numpy as np
import cv2
import pyaudio
import wave
import threading
import time
from datetime import datetime
import subprocess
import os
import ffmpeg

LOOP_DURATION = 5


class VideoRecorder():
    """基于openCV的视频录制类"""
    def __init__(self, name="temp_video.avi", fourcc="MJPG", sizex=640, sizey=480, camindex=0, fps=10):
        self.open = True
        self.device_index = camindex
        self.fps = fps                  # fps应为相机能稳定捕获画面的最低恒定帧率(需测试验证)
        self.fourcc = fourcc            # 视频编码格式,取决于所用相机
        self.frameSize = (sizex, sizey) # 视频分辨率,同样取决于相机
        self.video_filename = name
        self.video_cap = cv2.VideoCapture(self.device_index)
        self.video_cap.set(cv2.CAP_PROP_FPS, self.fps)
        self.video_writer = cv2.VideoWriter_fourcc(*self.fourcc)
        self.video_out = cv2.VideoWriter(self.video_filename, self.video_writer, self.fps, self.frameSize)
        self.frame_counts = 1
        self.start_time = time.time()

    def record(self):
        """开始录制视频"""
        timer_start = time.time()
        timer_current = 0
        while self.open:
            start = time.time()
            now = datetime.now()
            if not os.path.isdir("output"):
                os.mkdir("output")
            self.video_filename = os.path.join("./output", now.strftime("%y-%m-%d-%H_%M_%S") + '.avi')
            self.video_out = cv2.VideoWriter(self.video_filename, self.video_writer, self.fps, self.frameSize)
            while time.time() - start < LOOP_DURATION:
                ret, video_frame = self.video_cap.read()
                self.video_out.write(video_frame)

    def stop(self):
        """停止视频录制并释放线程资源"""
        if self.open:
            self.open=False
            self.video_out.release()
            self.video_cap.release()
            cv2.destroyAllWindows()

    def start(self):
        """启动视频录制线程"""
        video_thread = threading.Thread(target=self.record)
        video_thread.start()

class AudioRecorder():
    """基于pyAudio和Wave的音频录制类"""
    def __init__(self, filename="temp_audio.wav", rate=44100, fpb=1024, channels=2):
        self.open = True
        self.rate = rate
        self.frames_per_buffer = fpb
        self.channels = channels
        self.format = pyaudio.paInt16
        self.audio_filename = filename
        self.audio = pyaudio.PyAudio()
        self.stream = self.audio.open(format=self.format,
                                      channels=self.channels,
                                      rate=self.rate,
                                      input=True,
                                      frames_per_buffer = self.frames_per_buffer)
        self.audio_frames = []

    def record(self):
        """开始录制音频"""
        self.stream.start_stream()
        while self.open:
            start = time.time()
            now = datetime.now()
            if not os.path.isdir("output"):
                os.mkdir("output")
            self.audio_filename = os.path.join("./output", now.strftime("%y-%m-%d-%H_%M_%S") + '.wav')
            while time.time() - start < LOOP_DURATION:
                data = self.stream.read(self.frames_per_buffer)
                self.audio_frames.append(data)
            waveFile = wave.open(self.audio_filename, 'wb')
            waveFile.setnchannels(self.channels)
            waveFile.setsampwidth(self.audio.get_sample_size(self.format))
            waveFile.setframerate(self.rate)
            waveFile.writeframes(b''.join(self.audio_frames))
            waveFile.close()
            self.audio_frames = []
            if not self.open:
                break

    def stop(self):
        """停止音频录制并释放线程资源"""
        if self.open:
            self.open = False
            self.stream.stop_stream()
            self.stream.close()
            self.audio.terminate()
            waveFile = wave.open(self.audio_filename, 'wb')
            waveFile.setnchannels(self.channels)
            waveFile.setsampwidth(self.audio.get_sample_size(self.format))
            waveFile.setframerate(self.rate)
            waveFile.writeframes(b''.join(self.audio_frames))
            waveFile.close()

    def start(self):
        """启动音频录制线程"""
        audio_thread = threading.Thread(target=self.record)
        audio_thread.start()

def start_AVrecording(filename="test"):
    global video_thread
    global audio_thread
    video_thread = VideoRecorder()
    audio_thread = AudioRecorder()
    audio_thread.start()
    video_thread.start()
    return filename


def stop_AVrecording(filename="test"):
    audio_thread.stop()
    frame_counts = video_thread.frame_counts
    elapsed_time = time.time() - video_thread.start_time
    recorded_fps = frame_counts / elapsed_time
    print("总帧数 " + str(frame_counts))
    print("录制时长 " + str(elapsed_time))
    print("实际录制帧率 " + str(recorded_fps))
    video_thread.stop()

    # 等待所有录制线程结束
    while threading.active_count() > 1:
        print(f"活跃线程数: {threading.active_count()}")
        time.sleep(1)

    video_stream = ffmpeg.input(video_thread.video_filename)
    audio_stream = ffmpeg.input(audio_thread.audio_filename)
    now = datetime.now()
    output_file = os.path.join("./output", now.strftime("%y-%m-%d-%H_%M_%S") + '.mp4')
    ffmpeg.output(audio_stream, video_stream, output_file).run(overwrite_output=True)

def file_manager(filename="test"):
    """清理临时文件"""
    local_path = os.getcwd()
    if os.path.exists(str(local_path) + "/temp_audio.wav"):
        os.remove(str(local_path) + "/temp_audio.wav")
    if os.path.exists(str(local_path) + "/temp_video.avi"):
        os.remove(str(local_path) + "/temp_video.avi")
    if os.path.exists(str(local_path) + "/temp_video2.avi"):
        os.remove(str(local_path) + "/temp_video2.avi")

if __name__ == '__main__':
    start_AVrecording()
    time.sleep(30)
    stop_AVrecording()
解决方案

要实现每段固定时长的音视频录制完成后即时合并,核心是通过线程通信机制,让音视频线程在生成文件后通知合并线程处理,同时继续下一段录制。具体实现如下:

1. 核心思路

  • 用queue.Queue传递待合并的音视频文件标识(统一时间戳)
  • 新增独立合并线程,持续监听队列,收到完整文件对后立即执行合并
  • 确保音视频文件用同一时间戳命名,保证匹配准确性

2. 修改后的完整代码

#!/usr/bin/env python
# -*- coding: utf-8 -*-
# VideoRecorder.py

from __future__ import print_function, division
import numpy as np
import cv2
import pyaudio
import wave
import threading
import time
from datetime import datetime
import subprocess
import os
import ffmpeg
import queue

LOOP_DURATION = 5
# 创建队列用于传递待合并的文件时间戳
merge_queue = queue.Queue()

class VideoRecorder():
    """基于openCV的视频录制类"""
    def __init__(self, name="temp_video.avi", fourcc="MJPG", sizex=640, sizey=480, camindex=0, fps=10):
        self.open = True
        self.device_index = camindex
        self.fps = fps                  
        self.fourcc = fourcc            
        self.frameSize = (sizex, sizey) 
        self.video_filename = name
        self.video_cap = cv2.VideoCapture(self.device_index)
        self.video_cap.set(cv2.CAP_PROP_FPS, self.fps)
        self.video_writer = cv2.VideoWriter_fourcc(*self.fourcc)
        self.video_out = cv2.VideoWriter(self.video_filename, self.video_writer, self.fps, self.frameSize)
        self.frame_counts = 1
        self.start_time = time.time()

    def record(self):
        """开始录制视频"""
        while self.open:
            start = time.time()
            timestamp = datetime.now().strftime("%y-%m-%d-%H_%M_%S")
            if not os.path.isdir("output"):
                os.mkdir("output")
            self.video_filename = os.path.join("./output", f"{timestamp}.avi")
            self.video_out = cv2.VideoWriter(self.video_filename, self.video_writer, self.fps, self.frameSize)
            
            # 录制固定时长视频
            while time.time() - start < LOOP_DURATION and self.open:
                ret, video_frame = self.video_cap.read()
                if ret:
                    self.video_out.write(video_frame)
            
            # 释放当前段视频写入器
            self.video_out.release()
            # 向队列发送视频文件的时间戳
            if self.open:
                merge_queue.put(('video', timestamp))

    def stop(self):
        """停止视频录制并释放资源"""
        if self.open:
            self.open=False
            self.video_out.release()
            self.video_cap.release()
            cv2.destroyAllWindows()

    def start(self):
        """启动视频录制线程"""
        video_thread = threading.Thread(target=self.record)
        video_thread.start()

class AudioRecorder():
    """基于pyAudio和Wave的音频录制类"""
    def __init__(self, filename="temp_audio.wav", rate=44100, fpb=1024, channels=2):
        self.open = True
        self.rate = rate
        self.frames_per_buffer = fpb
        self.channels = channels
        self.format = pyaudio.paInt16
        self.audio_filename = filename
        self.audio = pyaudio.PyAudio()
        self.stream = self.audio.open(format=self.format,
                                      channels=self.channels,
                                      rate=self.rate,
                                      input=True,
                                      frames_per_buffer = self.frames_per_buffer)
        self.audio_frames = []

    def record(self):
        """开始录制音频"""
        self.stream.start_stream()
        while self.open:
            start = time.time()
            timestamp = datetime.now().strftime("%y-%m-%d-%H_%M_%S")
            if not os.path.isdir("output"):
                os.mkdir("output")
            self.audio_filename = os.path.join("./output", f"{timestamp}.wav")
            self.audio_frames = []
            
            # 录制固定时长音频
            while time.time() - start < LOOP_DURATION and self.open:
                data = self.stream.read(self.frames_per_buffer)
                self.audio_frames.append(data)
            
            # 写入音频文件
            if self.open:
                waveFile = wave.open(self.audio_filename, 'wb')
                waveFile.setnchannels(self.channels)
                waveFile.setsampwidth(self.audio.get_sample_size(self.format))
                waveFile.setframerate(self.rate)
                waveFile.writeframes(b''.join(self.audio_frames))
                waveFile.close()
                # 向队列发送音频文件的时间戳
                merge_queue.put(('audio', timestamp))

    def stop(self):
        """停止音频录制并释放资源"""
        if self.open:
            self.open = False
            self.stream.stop_stream()
            self.stream.close()
            self.audio.terminate()
            # 写入最后一段音频
            waveFile = wave.open(self.audio_filename, 'wb')
            waveFile.setnchannels(self.channels)
            waveFile.setsampwidth(self.audio.get_sample_size(self.format))
            waveFile.setframerate(self.rate)
            waveFile.writeframes(b''.join(self.audio_frames))
            waveFile.close()

    def start(self):
        """启动音频录制线程"""
        audio_thread = threading.Thread(target=self.record)
        audio_thread.start()

def merge_worker():
    """合并线程:监听队列,处理音视频合并"""
    pending_files = {}
    while True:
        try:
            # 阻塞等待队列消息,超时1秒检查录制状态
            file_type, timestamp = merge_queue.get(timeout=1)
            pending_files[file_type] = timestamp
            
            # 同一时间戳的音视频都到齐时,执行合并
            if 'video' in pending_files and 'audio' in pending_files and pending_files['video'] == pending_files['audio']:
                current_ts = pending_files['video']
                video_path = os.path.join("./output", f"{current_ts}.avi")
                audio_path = os.path.join("./output", f"{current_ts}.wav")
                output_path = os.path.join("./output", f"{current_ts}.mp4")
                
                # 调用ffmpeg合并,copy模式避免重新编码提升速度
                try:
                    ffmpeg.output(ffmpeg.input(audio_path), ffmpeg.input(video_path), 
                                  output_path, vcodec='copy', acodec='aac').run(overwrite_output=True, quiet=True)
                    print(f"已完成合并:{output_path}")
                    
                    # 可选:删除原分离文件
                    os.remove(video_path)
                    os.remove(audio_path)
                except Exception as e:
                    print(f"合并失败 [{current_ts}]: {str(e)}")
                
                # 清空待处理文件记录
                pending_files.clear()
            
            merge_queue.task_done()
        except queue.Empty:
            # 所有录制线程停止后,退出合并线程
            if not (video_thread.open or audio_thread.open):
                break
            continue

def start_AVrecording(filename="test"):
    global video_thread
    global audio_thread
    global merge_thread
    video_thread = VideoRecorder()
    audio_thread = AudioRecorder()
    # 启动合并线程(守护线程,主进程结束自动退出)
    merge_thread = threading.Thread(target=merge_worker, daemon=True)
    merge_thread.start()
    
    audio_thread.start()
    video_thread.start()
    return filename

def stop_AVrecording(filename="test"):
    audio_thread.stop()
    video_thread.stop()

    # 等待队列中所有任务处理完成
    merge_queue.join()
    merge_thread.join()

    print("所有录制和合并任务已完成")

def file_manager(filename="test"):
    """清理临时文件"""
    local_path = os.getcwd()
    if os.path.exists(str(local_path) + "/temp_audio.wav"):
        os.remove(str(local_path) + "/temp_audio.wav")
    if os.path.exists(str(local_path) + "/temp_video.avi"):
        os.remove(str(local_path) + "/temp_video.avi")
    if os.path.exists(str(local_path) + "/temp_video2.avi"):
        os.remove(str(local_path) + "/temp_video2.avi")

if __name__ == '__main__':
    start_AVrecording()
    time.sleep(30)
    stop_AVrecording()

关键改动说明

  • 队列通信:通过queue.Queue传递音视频文件的时间戳,确保同一时段的文件能准确匹配
  • 独立合并线程:merge_worker线程持续监听队列,无需等待录制线程结束即可处理合并
  • 文件名同步:音视频录制时使用完全相同的时间戳命名,避免匹配错误
  • 资源优化:每段录制完成后立即释放VideoWriter和音频文件句柄,减少资源占用

内容的提问来源于stack exchange,提问作者Albert G Lieu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 21:54:24