-
Notifications
You must be signed in to change notification settings - Fork 1
perf: speed up ARIS→MP4 convert (~15× small / ~2.6× large) #23
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
a1d54e1
f32228b
ac7ee8c
f14fcc9
f3bf3b1
8b416f9
93216d8
df9ca9b
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
|
@@ -12,16 +12,39 @@ | |||||||||||||||||
| import array | ||||||||||||||||||
| import datetime | ||||||||||||||||||
| import struct | ||||||||||||||||||
| import subprocess as sp | ||||||||||||||||||
| from fractions import Fraction | ||||||||||||||||||
|
|
||||||||||||||||||
| import av | ||||||||||||||||||
| import cv2 | ||||||||||||||||||
| import numpy as np | ||||||||||||||||||
| import pytz | ||||||||||||||||||
| import tqdm | ||||||||||||||||||
| from PIL import Image, ImageDraw, ImageFont | ||||||||||||||||||
|
|
||||||||||||||||||
| from ..video.processing import ProcessingPipeline | ||||||||||||||||||
| from . import beamLookUp | ||||||||||||||||||
|
|
||||||||||||||||||
| # FF_QP2LAMBDA: ffmpeg -q:v N sets global_quality = N * 118. | ||||||||||||||||||
| _FF_QP2LAMBDA = 118 | ||||||||||||||||||
|
|
||||||||||||||||||
| # mpeg4 / yuv420p does not reliably signal full range. Map 0–255 → 16–235 | ||||||||||||||||||
| # so a limited-range player stretch recovers the remapped gray. | ||||||||||||||||||
| _FULL_TO_LIMITED = np.array( | ||||||||||||||||||
| [16 + int(round(i * 219 / 255)) for i in range(256)], | ||||||||||||||||||
| dtype=np.uint8, | ||||||||||||||||||
| ) | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| def _mpeg4_compatible_rate(fps: float) -> Fraction: | ||||||||||||||||||
| """Return a frame rate Fraction legal for mpeg4's timebase limit. | ||||||||||||||||||
|
|
||||||||||||||||||
| Mpeg4 caps the timebase denominator at 65535. The timebase is ``1/rate``, | ||||||||||||||||||
| so the rate *numerator* must be bounded: limit the reciprocal, then invert. | ||||||||||||||||||
| Exact rates such as 24 and 29.97 are preserved; near-15 ARIS rates (e.g. | ||||||||||||||||||
| ``15.000149726867676``) become ``65521/4368``. | ||||||||||||||||||
| """ | ||||||||||||||||||
| return 1 / (1 / Fraction(fps)).limit_denominator(65535) | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| class ARIS_File: | ||||||||||||||||||
| "This is a class container for the ARIS file headers" | ||||||||||||||||||
|
|
@@ -691,6 +714,8 @@ def DataImport(filename, startFrame=1, frameBuffer=0): | |||||||||||||||||
| output_data.LUT = False | ||||||||||||||||||
| output_data.map_y = None | ||||||||||||||||||
| output_data.map_x = None | ||||||||||||||||||
| output_data.map_1 = None | ||||||||||||||||||
| output_data.map_2 = None | ||||||||||||||||||
|
|
||||||||||||||||||
| # Load the first frame | ||||||||||||||||||
| frame = FrameRead(output_data, startFrame) | ||||||||||||||||||
|
|
@@ -699,17 +724,17 @@ def DataImport(filename, startFrame=1, frameBuffer=0): | |||||||||||||||||
| return output_data, frame | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| def FrameRead(ARIS_data, frameIndex, frameBuffer=None) -> ARIS_Frame: | ||||||||||||||||||
| """The FrameRead function loads in the specified frame data from the raw ARIS data. | ||||||||||||||||||
| The function then calls the remapARIS() function which remaps the raw data into | ||||||||||||||||||
| a 2D real world projection. | ||||||||||||||||||
| def FrameRead(ARIS_data, frameIndex, frameBuffer=None, remap=True) -> ARIS_Frame: | ||||||||||||||||||
| """Load frame data from the raw ARIS file; optionally remap to 2D. | ||||||||||||||||||
|
|
||||||||||||||||||
| Parameters | ||||||||||||||||||
| ----------- | ||||||||||||||||||
| ARIS_data : ARIS data structure returned via pyARIS.DataImport() | ||||||||||||||||||
| frameIndex : frame number | ||||||||||||||||||
| frameBuffer : This parameter add a specified number of pixels around the edges | ||||||||||||||||||
| of the remapped frame. | ||||||||||||||||||
| remap : If True (default), call :func:`remapARIS`. If False, only read raw | ||||||||||||||||||
| beam/bin data (and build the LUT once if needed); caller remaps later. | ||||||||||||||||||
|
|
||||||||||||||||||
| Returns | ||||||||||||||||||
| ------- | ||||||||||||||||||
|
|
@@ -1191,16 +1216,15 @@ def FrameRead(ARIS_data, frameIndex, frameBuffer=None) -> ARIS_Frame: | |||||||||||||||||
| output.frame_data = frame | ||||||||||||||||||
| output.WinStart = output.samplestartdelay * 0.000001 * output.soundspeed / 2 | ||||||||||||||||||
|
|
||||||||||||||||||
| # Close the data file | ||||||||||||||||||
| data.close() | ||||||||||||||||||
|
|
||||||||||||||||||
| # Create the lookup table if not present | ||||||||||||||||||
| if ARIS_data.LUT == False: | ||||||||||||||||||
| LUT(ARIS_data, output) | ||||||||||||||||||
| ARIS_data.LUT = True | ||||||||||||||||||
|
|
||||||||||||||||||
| # Remap the frame | ||||||||||||||||||
| remapARIS(ARIS_data, output, frameBuffer) | ||||||||||||||||||
| if remap: | ||||||||||||||||||
| remapARIS(ARIS_data, output, frameBuffer) | ||||||||||||||||||
|
|
||||||||||||||||||
| return output | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
@@ -1246,17 +1270,23 @@ def getBeam(x, y, beamcount): | |||||||||||||||||
| return beamnum | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| getBeamVec = np.vectorize(getBeam) | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| def getBin(x, y, winstart, sampleperiod, soundspeed): | ||||||||||||||||||
| angle = np.rad2deg(np.tan(x / y)) | ||||||||||||||||||
| hyp = y / np.cos(np.deg2rad(angle)) | ||||||||||||||||||
| binnum2 = int((2 * (hyp - winstart)) / (sampleperiod * 0.000001 * soundspeed)) | ||||||||||||||||||
| return binnum2 | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| getBinVec = np.vectorize(getBin) | ||||||||||||||||||
| def _beam_breakpoints(beamcount): | ||||||||||||||||||
| if beamcount == 128: | ||||||||||||||||||
| return np.asarray(beamLookUp.breakpoints128, dtype=np.float64) | ||||||||||||||||||
| if beamcount == 96: | ||||||||||||||||||
| return np.asarray(beamLookUp.breakpoints96, dtype=np.float64) | ||||||||||||||||||
| if beamcount == 64: | ||||||||||||||||||
| return np.asarray(beamLookUp.breakpoints64, dtype=np.float64) | ||||||||||||||||||
| if beamcount == 48: | ||||||||||||||||||
| return np.asarray(beamLookUp.breakpoints48, dtype=np.float64) | ||||||||||||||||||
| raise ValueError(f"Unsupported BeamCount: {beamcount}") | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| def LUT(ARISFile, frame): | ||||||||||||||||||
|
|
@@ -1295,24 +1325,39 @@ def LUT(ARISFile, frame): | |||||||||||||||||
| # Create a matrix in real-world coordinates with the sonar located at 0,0 | ||||||||||||||||||
| xx, yy = np.meshgrid(x1, y1) | ||||||||||||||||||
|
|
||||||||||||||||||
| # Use the getBeamVec function to return the map_x | ||||||||||||||||||
| beamcount = frame.BeamCount | ||||||||||||||||||
| beams = getBeamVec(xx, yy, beamcount) | ||||||||||||||||||
| # Shared angle (same formula as getBeam / getBin) | ||||||||||||||||||
| angle = np.rad2deg(np.tan(xx / yy)) | ||||||||||||||||||
| lookup_angle = -angle | ||||||||||||||||||
|
|
||||||||||||||||||
| # map_x: beam index via searchsorted (same as bisect_right - 1) | ||||||||||||||||||
| breakpoints = _beam_breakpoints(frame.BeamCount) | ||||||||||||||||||
| beams = np.searchsorted(breakpoints, lookup_angle, side="right") - 1 | ||||||||||||||||||
| valid = (lookup_angle > breakpoints[0]) & (lookup_angle <= breakpoints[-1]) | ||||||||||||||||||
| beams = beams.astype(np.float64, copy=False) | ||||||||||||||||||
| beams[~valid] = 999 | ||||||||||||||||||
|
|
||||||||||||||||||
| # Create a map_y based on the bin locations | ||||||||||||||||||
| # map_y: range bin (truncate toward zero like int(...)) | ||||||||||||||||||
| winstart = frame.WinStart | ||||||||||||||||||
| sampleperiod = frame.sampleperiod | ||||||||||||||||||
| soundspeed = frame.soundspeed | ||||||||||||||||||
| bins = getBinVec(xx, yy, winstart, sampleperiod, soundspeed) | ||||||||||||||||||
| hyp = yy / np.cos(np.deg2rad(angle)) | ||||||||||||||||||
| bins = np.trunc((2 * (hyp - winstart)) / (sampleperiod * 0.000001 * soundspeed)) | ||||||||||||||||||
|
|
||||||||||||||||||
| # Clip the map_y mask | ||||||||||||||||||
| bins = bins.astype(np.float32) | ||||||||||||||||||
| mask = np.isnan(beams) | ||||||||||||||||||
| bins[mask] = beams[mask] | ||||||||||||||||||
|
|
||||||||||||||||||
| # Embed the maps in the ARISFile object | ||||||||||||||||||
| # Embed float maps (kept for debugging / callers) and fixed-point maps for remap. | ||||||||||||||||||
| ARISFile.map_y = bins | ||||||||||||||||||
| ARISFile.map_x = beams.astype(np.float32) | ||||||||||||||||||
| # nninterpolation=True matches INTER_NEAREST (map_2 may be empty). | ||||||||||||||||||
| ARISFile.map_1, ARISFile.map_2 = cv2.convertMaps( | ||||||||||||||||||
| np.ascontiguousarray(ARISFile.map_x), | ||||||||||||||||||
| np.ascontiguousarray(ARISFile.map_y), | ||||||||||||||||||
| cv2.CV_16SC2, | ||||||||||||||||||
| nninterpolation=True, | ||||||||||||||||||
| ) | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| def remapARIS(ARISFile, frame, frameBuffer=None): | ||||||||||||||||||
|
|
@@ -1329,12 +1374,13 @@ def remapARIS(ARISFile, frame, frameBuffer=None): | |||||||||||||||||
| ------- | ||||||||||||||||||
| A remapped frame which is stored in the frames data structure as frame.remap | ||||||||||||||||||
| """ | ||||||||||||||||||
| # Create an empty frame | ||||||||||||||||||
| Remap = np.zeros([ARISFile.ydim, ARISFile.xdim]) | ||||||||||||||||||
|
|
||||||||||||||||||
| # Use OpenCV's remap function to populate the empty frame | ||||||||||||||||||
| Remap = cv2.remap(frame.frame_data, ARISFile.map_x, ARISFile.map_y, cv2.INTER_NEAREST) | ||||||||||||||||||
|
|
||||||||||||||||||
| # uint8 in → uint8 out; fixed-point maps from LUT/convertMaps | ||||||||||||||||||
| Remap = cv2.remap( | ||||||||||||||||||
| frame.frame_data, | ||||||||||||||||||
| ARISFile.map_1, | ||||||||||||||||||
| ARISFile.map_2, | ||||||||||||||||||
| cv2.INTER_NEAREST, | ||||||||||||||||||
| ) | ||||||||||||||||||
| Remap = np.flipud(Remap) | ||||||||||||||||||
|
|
||||||||||||||||||
| # Add buffer is requested | ||||||||||||||||||
|
|
@@ -1343,22 +1389,21 @@ def remapARIS(ARISFile, frame, frameBuffer=None): | |||||||||||||||||
| buffX = int(ARISFile.xdim * frameBuffer) | ||||||||||||||||||
| Remap = np.concatenate( | ||||||||||||||||||
| ( | ||||||||||||||||||
| np.ones([ARISFile.ydim, buffX]), | ||||||||||||||||||
| np.ones([ARISFile.ydim, buffX], dtype=np.uint8), | ||||||||||||||||||
| Remap, | ||||||||||||||||||
| np.ones([ARISFile.ydim, buffX]), | ||||||||||||||||||
| np.ones([ARISFile.ydim, buffX], dtype=np.uint8), | ||||||||||||||||||
| ), | ||||||||||||||||||
| axis=1, | ||||||||||||||||||
| ) | ||||||||||||||||||
| Remap = np.concatenate( | ||||||||||||||||||
| ( | ||||||||||||||||||
| np.ones([buffY, ARISFile.xdim + buffX * 2]), | ||||||||||||||||||
| np.ones([buffY, ARISFile.xdim + buffX * 2], dtype=np.uint8), | ||||||||||||||||||
| Remap, | ||||||||||||||||||
| np.ones([buffY, ARISFile.xdim + buffX * 2]), | ||||||||||||||||||
| np.ones([buffY, ARISFile.xdim + buffX * 2], dtype=np.uint8), | ||||||||||||||||||
| ) | ||||||||||||||||||
| ) | ||||||||||||||||||
|
|
||||||||||||||||||
| # Add to frame data | ||||||||||||||||||
| frame.remap = Remap.astype("uint8") | ||||||||||||||||||
| frame.remap = Remap | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| def VideoExport( | ||||||||||||||||||
|
|
@@ -1371,9 +1416,14 @@ def VideoExport( | |||||||||||||||||
| fontsize=30, | ||||||||||||||||||
| ts_pos=(0, 0), | ||||||||||||||||||
| vbr=10, | ||||||||||||||||||
| num_workers=1, | ||||||||||||||||||
| codec="mpeg4", | ||||||||||||||||||
| show_progress=True, | ||||||||||||||||||
| ): | ||||||||||||||||||
| """Output video using the ffmpeg pipeline. The current implementation | ||||||||||||||||||
| outputs compresses png files and outputs a mp4. | ||||||||||||||||||
| """Output video using in-process PyAV, overlapping I/O with remap/encode. | ||||||||||||||||||
|
|
||||||||||||||||||
| Producer reads raw frames (no remap); consumers remap + timestamp/pad + | ||||||||||||||||||
| pixel prep; ordered writer only encodes/muxes. | ||||||||||||||||||
|
|
||||||||||||||||||
| Parameters | ||||||||||||||||||
| ----------- | ||||||||||||||||||
|
|
@@ -1384,15 +1434,17 @@ def VideoExport( | |||||||||||||||||
| timestamp : (Bool) Add the timestamp from the sonar to the video frames | ||||||||||||||||||
| fontsize : (Int) Size of timestamp font | ||||||||||||||||||
| ts_pos : (Tuple) (x,y) location of the timestamp | ||||||||||||||||||
| vbr : (Int) Variable Bitrate of output video (1-31) 1 being highest quality, 31 being lowest quality | ||||||||||||||||||
| vbr : (Int) mpeg4 quantizer (1–31); 1 is highest quality, 31 lowest. Same as ffmpeg -q:v. | ||||||||||||||||||
| num_workers : (Int) Pass-through consumer threads (default 1; raise when consumers do real work) | ||||||||||||||||||
| codec : (Str) PyAV encoder name (default ``mpeg4``; e.g. ``h264_videotoolbox``, ``h264_nvenc``) | ||||||||||||||||||
| show_progress : (Bool) Show tqdm progress bar | ||||||||||||||||||
|
|
||||||||||||||||||
| Returns | ||||||||||||||||||
| ------- | ||||||||||||||||||
| Returns a video into the current working directory | ||||||||||||||||||
|
|
||||||||||||||||||
| Notes | ||||||||||||||||||
| ------ | ||||||||||||||||||
| ffmpeg must be available in the system PATH. | ||||||||||||||||||
| Must have the '*.mp4' file extension. | ||||||||||||||||||
| Uses the tqdm package to display a status bar. | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
@@ -1402,38 +1454,26 @@ def VideoExport( | |||||||||||||||||
|
|
||||||||||||||||||
| """ | ||||||||||||||||||
|
|
||||||||||||||||||
| # Command to send via the command prompt which specifies the pipe parameters | ||||||||||||||||||
| # Linux compatible: use 'ffmpeg' instead of 'ffmpeg.exe' | ||||||||||||||||||
| command = [ | ||||||||||||||||||
| "ffmpeg", | ||||||||||||||||||
| "-y", # (optional) overwrite output file if it exists | ||||||||||||||||||
| "-f", | ||||||||||||||||||
| "image2pipe", | ||||||||||||||||||
| "-vcodec", | ||||||||||||||||||
| "mjpeg", | ||||||||||||||||||
| "-r", | ||||||||||||||||||
| str(fps), # frames per second | ||||||||||||||||||
| "-i", | ||||||||||||||||||
| "-", # The input comes from a pipe | ||||||||||||||||||
| "-an", # Tells FFMPEG not to expect any audio | ||||||||||||||||||
| "-vcodec", | ||||||||||||||||||
| "mpeg4", | ||||||||||||||||||
| "-q:v", | ||||||||||||||||||
| str(vbr), # Variable bitrate quality | ||||||||||||||||||
| filename, | ||||||||||||||||||
| ] | ||||||||||||||||||
|
|
||||||||||||||||||
| # Open the pipe | ||||||||||||||||||
| pipe = sp.Popen(command, stdin=sp.PIPE, stderr=sp.DEVNULL) | ||||||||||||||||||
|
|
||||||||||||||||||
| if end_frame is None: | ||||||||||||||||||
| end_frame = data.FrameCount | ||||||||||||||||||
|
|
||||||||||||||||||
| # Iterate through the dataframes and push to pipe | ||||||||||||||||||
| for i in tqdm.tqdm(range(start_frame - 1, end_frame)): | ||||||||||||||||||
| frame = FrameRead(data, i) | ||||||||||||||||||
| im = Image.fromarray(frame.remap) | ||||||||||||||||||
| frame_indices = range(start_frame - 1, end_frame) | ||||||||||||||||||
| n_frames = len(frame_indices) | ||||||||||||||||||
| container = None | ||||||||||||||||||
| stream = None | ||||||||||||||||||
| progress = tqdm.tqdm(total=n_frames, disable=not show_progress) | ||||||||||||||||||
|
|
||||||||||||||||||
| def producer(): | ||||||||||||||||||
| for out_i, frame_i in enumerate(frame_indices): | ||||||||||||||||||
| frame = FrameRead(data, frame_i, remap=False) | ||||||||||||||||||
| yield out_i, frame | ||||||||||||||||||
|
|
||||||||||||||||||
| def consumer(item): | ||||||||||||||||||
| out_i, frame = item | ||||||||||||||||||
| remapARIS(data, frame, None) | ||||||||||||||||||
| remap = frame.remap | ||||||||||||||||||
| if timestamp: | ||||||||||||||||||
| im = Image.fromarray(remap) | ||||||||||||||||||
| ts = str( | ||||||||||||||||||
| datetime.datetime.fromtimestamp( | ||||||||||||||||||
| frame.sonartimestamp / 1000000, pytz.timezone("UTC") | ||||||||||||||||||
|
|
@@ -1443,10 +1483,70 @@ def VideoExport( | |||||||||||||||||
| try: | ||||||||||||||||||
| font = ImageFont.truetype("DejaVuSans.ttf", fontsize) | ||||||||||||||||||
| except OSError: | ||||||||||||||||||
| # Fallback to default font if DejaVuSans not available | ||||||||||||||||||
| font = ImageFont.load_default() | ||||||||||||||||||
| draw.text(ts_pos, ts, font=font, fill="white") | ||||||||||||||||||
| im.save(pipe.stdin, "JPEG") | ||||||||||||||||||
| remap = np.asarray(im) | ||||||||||||||||||
| remap = cv2.LUT(np.ascontiguousarray(remap), _FULL_TO_LIMITED) | ||||||||||||||||||
| # yuv420p (I420) needs even W/H; pad odd ARIS remaps. | ||||||||||||||||||
| h, w = remap.shape[:2] | ||||||||||||||||||
| if (h | w) & 1: | ||||||||||||||||||
| remap = np.pad(remap, ((0, h & 1), (0, w & 1)), mode="edge") | ||||||||||||||||||
| h, w = remap.shape[:2] | ||||||||||||||||||
| # Y = limited-range gray, U/V = 128 (I420 = Y + half-height chroma). | ||||||||||||||||||
| yuv = np.vstack([remap, np.full((h // 2, w), 128, dtype=np.uint8)]) | ||||||||||||||||||
| video_frame = av.VideoFrame.from_ndarray(yuv, format="yuv420p") | ||||||||||||||||||
| return out_i, video_frame | ||||||||||||||||||
|
|
||||||||||||||||||
| # ProcessingPipeline does not reorder; buffer like muimg VideoEncodePipeline. | ||||||||||||||||||
| next_index = 0 | ||||||||||||||||||
| frame_buffer = {} | ||||||||||||||||||
|
|
||||||||||||||||||
| def writer(result): | ||||||||||||||||||
| nonlocal container, stream, next_index | ||||||||||||||||||
| index, video_frame = result | ||||||||||||||||||
| frame_buffer[index] = video_frame | ||||||||||||||||||
| while next_index in frame_buffer: | ||||||||||||||||||
| video_frame = frame_buffer.pop(next_index) | ||||||||||||||||||
| if container is None: | ||||||||||||||||||
| container = av.open(filename, mode="w") | ||||||||||||||||||
| # mpeg4 caps timebase den at 65535; bound the reciprocal of fps. | ||||||||||||||||||
| stream = container.add_stream( | ||||||||||||||||||
| codec, rate=_mpeg4_compatible_rate(fps) | ||||||||||||||||||
| ) | ||||||||||||||||||
|
Comment on lines
+1513
to
+1515
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Blocking: hangs on every ARIS file we have
I hit this on every recording I tried, plus the repo's own test fixture. They all report that same framerate, so this is the normal path, not an edge case. Careful with the fix:
Suggested change
Verified: Worth a regression test that exports two frames at |
||||||||||||||||||
| stream.width = video_frame.width | ||||||||||||||||||
| stream.height = video_frame.height | ||||||||||||||||||
| stream.pix_fmt = "yuv420p" | ||||||||||||||||||
| q = max(int(vbr), 1) | ||||||||||||||||||
| ctx = stream.codec_context | ||||||||||||||||||
| ctx.flags |= av.codec.context.Flags.qscale | ||||||||||||||||||
| ctx.global_quality = q * _FF_QP2LAMBDA | ||||||||||||||||||
| ctx.qmin = ctx.qmax = q | ||||||||||||||||||
|
|
||||||||||||||||||
| for packet in stream.encode(video_frame): | ||||||||||||||||||
| container.mux(packet) | ||||||||||||||||||
| progress.update(1) | ||||||||||||||||||
| next_index += 1 | ||||||||||||||||||
|
|
||||||||||||||||||
| try: | ||||||||||||||||||
| if n_frames > 0: | ||||||||||||||||||
| pipeline = ProcessingPipeline( | ||||||||||||||||||
| producer, | ||||||||||||||||||
| consumer, | ||||||||||||||||||
| writer, | ||||||||||||||||||
| num_workers=num_workers, | ||||||||||||||||||
| total_items=n_frames, | ||||||||||||||||||
| task_name="VideoExport", | ||||||||||||||||||
| ) | ||||||||||||||||||
| pipeline.run() | ||||||||||||||||||
| assert not frame_buffer, ( | ||||||||||||||||||
| f"VideoExport reorder buffer not drained: {sorted(frame_buffer)}" | ||||||||||||||||||
| ) | ||||||||||||||||||
| progress.close() | ||||||||||||||||||
| if stream is not None: | ||||||||||||||||||
| for packet in stream.encode(): | ||||||||||||||||||
| container.mux(packet) | ||||||||||||||||||
| finally: | ||||||||||||||||||
| progress.close() | ||||||||||||||||||
| if container is not None: | ||||||||||||||||||
| container.close() | ||||||||||||||||||
|
|
||||||||||||||||||
| pipe.stdin.close() | ||||||||||||||||||
| pipe.wait() | ||||||||||||||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Reorder buffer never recovers from a dropped frame
ProcessingPipelineswallows consumer exceptions (processing.py:171), so ifconsumerfails for frame k, that index never arrives.next_indexstalls at k permanently: every later frame accumulates inframe_bufferas a decodedav.VideoFrameand none are ever encoded. The video is silently truncated at k, memory grows unbounded, and if k is 0 thencontainerstaysNoneand no file is written at all.process_aris_filepathstill returns True and logs success.Measured on a 40 frame export, forcing
remapARISto raise on exactly one frame:VideoExportreturnedBoth cases look like success to the caller.
Once comment 2 is in place, the cleanest fix is to let consumer failures cancel the run and surface the error, rather than continuing with a hole in the sequence. A cheap guard in the meantime: assert
frame_bufferis empty afterpipeline.run()returns.