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

如何用GStreamer监控目录并实时推送新MP4至Kinesis视频流

实现持续监控目录并推送视频片段到Kinesis Video Stream

Absolutely! This is totally feasible with GStreamer, and we can set up a clean, automated pipeline that handles new files as they arrive—no need to manually configure a pipeline for each individual MP4. Let’s walk through the best approaches and how to integrate with Kinesis Video Streams (KVS).

核心思路

Your setup involves a directory where 60-second MP4 clips are continuously written, with each file named by the start timestamp of the segment. The goal is to:

  • Monitor /tmp/media for new, fully-written MP4 files
  • Automatically decode each clip
  • Stream the decoded content sequentially to KVS, maintaining the continuous video flow
  • Clean up files older than an hour to save space

方案1:Python Watchdog + GStreamer Pipeline Queue

This approach uses a lightweight Python script to monitor the directory, track when files are fully written, and feed them into a persistent GStreamer pipeline. It’s flexible and easy to customize.

Step 1: Install Dependencies

First, install the watchdog library for directory monitoring:

pip install watchdog

Ensure you have GStreamer and the KVS plugin installed (you’ll need the AWS GStreamer SDK for KVS support).

Step 2: Directory Monitoring Script

This script watches for new files, waits until they’re fully written (by checking if the file size stabilizes), then adds them to a processing queue. The GStreamer pipeline will consume this queue sequentially:

import time
import os
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
import gi
gi.require_version('Gst', '1.0')
from gi.repository import Gst, GObject

Gst.init(None)
GObject.threads_init()

class MediaFileHandler(FileSystemEventHandler):
    def __init__(self, pipeline_queue):
        self.pipeline_queue = pipeline_queue
        self.file_sizes = {}

    def on_created(self, event):
        if not event.is_directory and event.src_path.endswith('.mp4'):
            self.file_sizes[event.src_path] = os.path.getsize(event.src_path)
            # Wait for the file to finish writing (check size every 2s for 10s)
            for _ in range(5):
                time.sleep(2)
                new_size = os.path.getsize(event.src_path)
                if new_size == self.file_sizes[event.src_path]:
                    self.pipeline_queue.put(event.src_path)
                    break
                self.file_sizes[event.src_path] = new_size

def run_pipeline(queue):
    # Base pipeline template: decode MP4, convert to KVS-compatible format, push to KVS
    pipeline_template = """
    filesrc location={file_path} ! qtdemux ! h264parse ! avdec_h264 ! 
    videoconvert ! video/x-raw,format=I420 ! x264enc bitrate=500000 ! 
    h264parse ! kvssink stream-name=your-kvs-stream-name 
    access-key=your-aws-access-key secret-key=your-aws-secret-key
    """
    while True:
        file_path = queue.get()
        print(f"Processing file: {file_path}")
        pipeline_str = pipeline_template.format(file_path=file_path)
        pipeline = Gst.parse_launch(pipeline_str)
        bus = pipeline.get_bus()
        pipeline.set_state(Gst.State.PLAYING)
        # Wait for the pipeline to finish processing the file
        bus.timed_pop_filtered(Gst.CLOCK_TIME_NONE, Gst.MessageType.EOS | Gst.MessageType.ERROR)
        pipeline.set_state(Gst.State.NULL)
        # Optional: Delete the file after processing (or keep for cleanup later)
        # os.remove(file_path)

if __name__ == "__main__":
    import queue
    file_queue = queue.Queue()
    event_handler = MediaFileHandler(file_queue)
    observer = Observer()
    observer.schedule(event_handler, path='/tmp/media', recursive=False)
    observer.start()
    # Start the pipeline processing thread
    import threading
    pipeline_thread = threading.Thread(target=run_pipeline, args=(file_queue,))
    pipeline_thread.start()
    try:
        while True:
            time.sleep(1)
    except KeyboardInterrupt:
        observer.stop()
    observer.join()
    pipeline_thread.join()

方案2:原生GStreamer multifilesrc (更简洁)

If you prefer a pure GStreamer solution without Python, you can use multifilesrc to automatically detect and load new files matching your timestamp pattern. This is more streamlined but requires precise file naming matching.

Pipeline Example

gst-launch-1.0 multifilesrc location=/tmp/media/video-%Y-%m-%dT%H:%M:00.mp4 \
    index=0 next-file=forever poll-interval=1000 timeout=1000000000 \
    ! qtdemux ! h264parse ! avdec_h264 \
    ! videoconvert ! video/x-raw,format=I420 ! x264enc bitrate=500000 \
    ! h264parse ! kvssink stream-name=your-kvs-stream-name \
    access-key=your-aws-access-key secret-key=your-aws-secret-key

Key Parameters Explained:

  • location: Uses strftime-style formatting to match your file naming pattern (video-%Y-%m-%dT%H:%M:00.mp4)
  • next-file=forever: Tells multifilesrc to keep checking for new files indefinitely
  • poll-interval=1000: Checks for new files every 1 second (in milliseconds)
  • timeout=1000000000: Waits up to 1 second for the next file before continuing (prevents pipeline stalls)

Note: This assumes files are written with exact timestamp alignment (e.g., every minute at :00). If there’s any drift, you might need to adjust the pattern or add logic to detect the latest file.


集成Kinesis Video Streams Notes

  • Ensure the kvssink plugin is installed (part of the AWS Kinesis Video Streams GStreamer SDK)
  • You can avoid hardcoding credentials by setting AWS environment variables (AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY) or using an IAM role if running on EC2/EKS
  • Adjust the encoding parameters (x264enc bitrate) to match your desired quality and bandwidth

旧文件清理

To maintain only 1 hour of history, add a cron job that runs every minute to delete old files:

* * * * * find /tmp/media -name "video-*.mp4" -mmin +60 -delete

This will delete any MP4 files modified more than 60 minutes ago.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 09:24:38