如何用GStreamer监控目录并实时推送新MP4至Kinesis视频流
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/mediafor 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: Tellsmultifilesrcto keep checking for new files indefinitelypoll-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
kvssinkplugin is installed (part of the AWS Kinesis Video Streams GStreamer SDK) - You can avoid hardcoding credentials by setting AWS environment variables (
AWS_ACCESS_KEY_IDandAWS_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

