feat(dedup): 动态抽帧 + 滑动窗口时序匹配 (#1659)
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 1s
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 2s
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 22s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 9s
CI/CD Pipeline / Unit Tests (pull_request) Successful in 1m22s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 1m25s
AI Code Review / AI Code Review (pull_request) Failing after 1m50s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m39s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 2m44s
CI/CD Pipeline / Validate - Style (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Security (pull_request) Has been cancelled
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been cancelled
CI/CD Pipeline / Build Production API Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Web Image (pull_request) Has been cancelled
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been cancelled
CI/CD Pipeline / Deploy Production (pull_request) Has been cancelled
CI/CD Pipeline / Production Browser E2E (pull_request) Has been cancelled
CI/CD Pipeline / Canary Release to Production (pull_request) Has been cancelled
CI/CD Pipeline / CI Gate (pull_request) Has been cancelled
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been cancelled
CI/CD Pipeline / ACR Image Cleanup (pull_request) Failing after 304h23m6s
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Failing after 304h23m16s
CI/CD Pipeline / Staging API Integration Tests (pull_request) Failing after 304h23m6s
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Failing after 304h23m16s
CI/CD Pipeline / Staging E2E Tests (pull_request) Failing after 304h23m7s
CI/CD Pipeline / Build Staging Web Image (pull_request) Failing after 304h23m21s
CI/CD Pipeline / Build Staging API Image (pull_request) Failing after 304h23m21s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Failing after 304h23m16s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Failing after 304h23m43s
CI/CD Pipeline / Frontend Lint (pull_request) Failing after 304h23m43s
CI/CD Pipeline / Check push changed paths (pull_request) Failing after 304h23m53s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Failing after 304h57m36s
CI/CD Pipeline / Build Staging Worker Image (pull_request) Failing after 304h57m45s
CI/CD Pipeline / PR Build Web Image (pull_request) Failing after 304h57m45s

1. 动态抽帧策略 — detect_keyframe_timestamps()
   - 降采样到 320x240 逐帧灰度差异检测场景切换
   - 最小间隔过滤(保留差异最大的候选帧)
   - 数量裁剪到 [MIN_KEYFRAMES=5, MAX_KEYFRAMES=30]
   - 长视频(>3分钟)每 30 秒分段保底

2. 滑动窗口时序匹配 — find_duplicate_segments()
   - 逐帧最佳匹配 → 连续 run 检测(允许 MAX_GAP=2 间隙)
   - 最少 MIN_CONSECUTIVE_MATCHES=5 帧才报告
   - 返回 DuplicateSegment(query/target 时间范围 + 平均距离)

3. 查重算法升级
   - 均值距离 → 中位数距离(抵抗异常值)
   - 新增帧匹配比例条件(match_ratio >= 0.7)
   - Bhattacharyya 系数替代余弦相似度
   - pHash + 直方图加权融合(0.7/0.3)
   - 判定重复后附加 duplicate_segments 字段

4. 删除旧代码
   - 移除 SHORT_VIDEO_CHUNK_SEC/LONG_VIDEO_CHUNK_SEC 固定间隔
   - 移除 compute_chunk_interval()
   - 移除 _average_histogram_similarity()

5. 测试
   - 新增 test_dedup_v2.py: 34 个测试
   - 更新 test_dedup_engine.py/test_duplicate_rate.py/test_dedup_pure.py
   - 清理 test_fingerprint_chunks.py 中旧常量测试
This commit is contained in:
xiaoxia
2026-09-03 22:52:47 +08:00
parent 159a62f9a5
commit fac80b1f77
6 changed files with 1037 additions and 196 deletions
+458 -103
View File
@@ -1,8 +1,12 @@
"""Video deduplication module - compute fingerprints and detect duplicates."""
"""Video deduplication module - compute fingerprints and detect duplicates.
Dynamic keyframe detection + sliding window temporal matching (Issue #1659).
"""
import hashlib
import logging
import os
import statistics
import tempfile
from dataclasses import dataclass, field
from typing import Optional
@@ -21,11 +25,28 @@ from packages.shared.storage import get_storage_service
logger = logging.getLogger(__name__)
# 分片策略常量
SHORT_VIDEO_CHUNK_SEC = 2 # ≤60秒视频,每 2 秒一个分片
LONG_VIDEO_CHUNK_SEC = 5 # >60秒视频,每 5 秒一个分片
SHORT_VIDEO_THRESHOLD_SEC = 60
# ── 关键帧检测常量 ──────────────────────────────────────────────
SCENE_CHANGE_THRESHOLD = 30 # 灰度差异阈值
MIN_KEYFRAME_INTERVAL_SEC = 1.0 # 最小关键帧间隔(秒)
MAX_KEYFRAMES = 30 # 最大关键帧数
MIN_KEYFRAMES = 5 # 最小关键帧数
LONG_VIDEO_SEGMENT_SEC = 30 # 长视频每段秒数
LONG_VIDEO_DURATION_THRESHOLD_SEC = 180 # 3 分钟阈值
MIN_FRAMES_PER_SEGMENT = 2 # 长视频每段最少帧数
# ── 滑动窗口匹配常量 ────────────────────────────────────────────
SEGMENT_MATCH_THRESHOLD = 8 # 帧匹配汉明距离阈值
MIN_CONSECUTIVE_MATCHES = 5 # 最少连续匹配帧数
MAX_GAP = 2 # 允许的最大间隙帧数
# ── 融合判定常量 ────────────────────────────────────────────────
PHASH_WEIGHT = 0.7 # pHash 权重
HISTOGRAM_WEIGHT = 0.3 # 直方图权重
MATCH_RATIO_THRESHOLD = 0.7 # 至少 70% 帧匹配
DUPLICATE_THRESHOLD = 0.70 # 融合后相似度阈值
# ── 感知哈希 & 颜色直方图工具函数 ────────────────────────────────
def compute_phash(image: np.ndarray, hash_size: int = 8) -> str:
"""计算图像的感知哈希(pHash),基于 DCT(离散余弦变换)。
@@ -87,16 +108,106 @@ def compute_color_histogram(image: np.ndarray, bins: int = 32) -> list[float]:
return hist
def compute_chunk_interval(duration: float) -> float:
"""根据视频时长返回分片间隔(秒)。
# ── 关键帧检测 ──────────────────────────────────────────────────
短视频(≤60秒):每 2 秒一个分片
长视频(>60秒):每 5 秒一个分片
def detect_keyframe_timestamps(
video_path: str,
*,
min_interval_sec: float = MIN_KEYFRAME_INTERVAL_SEC,
max_frames: int = MAX_KEYFRAMES,
min_frames: int = MIN_KEYFRAMES,
) -> list[float]:
"""检测视频中的场景切换点,返回关键帧时间戳列表(秒)。
算法:
1. 降采样到 320x240,逐帧转灰度
2. 计算相邻帧灰度差异(像素均值差)
3. 差异 > SCENE_CHANGE_THRESHOLD(30) 标记为候选关键帧
4. 相邻关键帧间隔 < min_interval_sec 的,保留差异更大的那个
5. 数量裁剪到 [min_frames, max_frames]
对于长视频(>3分钟):
- 每 30 秒一个分段
- 每个分段至少选 2 个关键帧(如果分段内无场景切换,均匀取 2 帧)
"""
if duration <= SHORT_VIDEO_THRESHOLD_SEC:
return SHORT_VIDEO_CHUNK_SEC
return LONG_VIDEO_CHUNK_SEC
cap = cv2.VideoCapture(video_path)
if not cap.isOpened():
raise RuntimeError(f"Cannot open video: {video_path}")
fps = cap.get(cv2.CAP_PROP_FPS)
frame_count = int(cap.get(cv2.CAP_PROP_FRAME_COUNT))
duration = frame_count / fps if fps > 0 else 0
if duration <= 0:
cap.release()
return []
# 逐帧检测场景切换
candidates: list[tuple[float, float]] = [] # (timestamp_sec, diff_score)
prev_gray = None
while True:
ret, frame = cap.read()
if not ret:
break
# 降采样 + 灰度
small = cv2.resize(frame, (320, 240))
gray = cv2.cvtColor(small, cv2.COLOR_BGR2GRAY).astype(np.float32)
if prev_gray is not None:
diff = float(np.mean(np.abs(gray - prev_gray)))
if diff > SCENE_CHANGE_THRESHOLD:
pos_ms = cap.get(cv2.CAP_PROP_POS_MSEC)
candidates.append((pos_ms / 1000.0, diff))
prev_gray = gray
cap.release()
# 按最小间隔过滤(保留差异更大的)
filtered: list[tuple[float, float]] = []
for ts, diff in sorted(candidates):
if filtered and (ts - filtered[-1][0]) < min_interval_sec:
if diff > filtered[-1][1]:
filtered[-1] = (ts, diff)
else:
filtered.append((ts, diff))
keyframe_times = [ts for ts, _ in filtered]
# 数量不足 min_frames 时,在时间轴上均匀补充
if len(keyframe_times) < min_frames:
uniform = [duration * (i + 0.5) / min_frames for i in range(min_frames)]
keyframe_times = sorted(set(uniform) | set(keyframe_times))
# 如果合并后还不足 min_frames,直接用均匀分布
if len(keyframe_times) < min_frames:
keyframe_times = uniform
# 数量超过 max_frames 时,均匀采样
if len(keyframe_times) > max_frames:
step = len(keyframe_times) / max_frames
keyframe_times = [keyframe_times[int(i * step)] for i in range(max_frames)]
# 长视频分段保底(>3分钟)
if duration > LONG_VIDEO_DURATION_THRESHOLD_SEC:
segment_count = int(duration / LONG_VIDEO_SEGMENT_SEC)
for seg_idx in range(segment_count):
seg_start = seg_idx * LONG_VIDEO_SEGMENT_SEC
seg_end = min((seg_idx + 1) * LONG_VIDEO_SEGMENT_SEC, duration)
seg_frames = [t for t in keyframe_times if seg_start <= t < seg_end]
if len(seg_frames) < MIN_FRAMES_PER_SEGMENT:
# 均匀补齐
for i in range(MIN_FRAMES_PER_SEGMENT):
t = seg_start + LONG_VIDEO_SEGMENT_SEC * (i + 0.5) / MIN_FRAMES_PER_SEGMENT
if t not in keyframe_times and seg_start <= t < seg_end:
keyframe_times.append(t)
keyframe_times.sort()
return keyframe_times
# ── 数据类 ──────────────────────────────────────────────────────
@dataclass
class FingerprintChunk:
@@ -109,6 +220,17 @@ class FingerprintChunk:
frame_count: int = 1
@dataclass
class DuplicateSegment:
"""一段重复片段的描述。"""
query_start_ms: int
query_end_ms: int
target_start_ms: int
target_end_ms: int
avg_distance: float # 该段内帧的平均汉明距离
@dataclass
class VideoFingerprint:
"""Video fingerprint containing multiple similarity metrics."""
@@ -163,6 +285,135 @@ class VideoFingerprint:
return models
# ── 滑动窗口时序匹配 ────────────────────────────────────────────
def find_duplicate_segments(
query_chunks: list,
target_chunks: list,
*,
match_threshold: int = SEGMENT_MATCH_THRESHOLD,
min_consecutive: int = MIN_CONSECUTIVE_MATCHES,
max_gap: int = MAX_GAP,
) -> list[DuplicateSegment]:
"""滑动窗口时序匹配:找出两组分片之间的重复片段。
算法:
1. 对每个 query chunk,找到 target 中汉明距离最小的 chunk
2. 距离 <= match_threshold 视为匹配
3. 找连续匹配的 run(允许 max_gap 帧间隙)
4. 连续匹配数 >= min_consecutive 的 run 报告为重复片段
Args:
query_chunks: 查询视频的分片列表(FingerprintChunk 或 dict)
target_chunks: 目标视频的分片列表
match_threshold: 汉明距离匹配阈值
min_consecutive: 最少连续匹配帧数
max_gap: 允许的最大间隙帧数
Returns:
DuplicateSegment 列表
"""
if not query_chunks or not target_chunks:
return []
def _get_phash(chunk) -> str:
if isinstance(chunk, dict):
return chunk["phash_binary"]
return chunk.phash_binary
def _get_start(chunk) -> int:
if isinstance(chunk, dict):
return chunk["start_time_ms"]
return chunk.start_time_ms
def _get_end(chunk) -> int:
if isinstance(chunk, dict):
return chunk["end_time_ms"]
return chunk.end_time_ms
# Step 1: 逐帧匹配
frame_matches: list[tuple[bool, int, int]] = [] # (is_match, min_dist, best_target_idx)
for qc in query_chunks:
qc_phash = _get_phash(qc)
best_dist = 64
best_idx = 0
for j, tc in enumerate(target_chunks):
d = hamming_distance(qc_phash, _get_phash(tc))
if d < best_dist:
best_dist = d
best_idx = j
frame_matches.append((best_dist <= match_threshold, best_dist, best_idx))
# Step 2: 找连续匹配的 runs
runs: list[tuple[int, int]] = [] # list of (start_idx, end_idx)
run_start = None
gap_count = 0
for i, (is_match, dist, idx) in enumerate(frame_matches):
if is_match:
if run_start is None:
run_start = i
gap_count = 0 # 重置间隙
else:
if run_start is not None:
gap_count += 1
if gap_count > max_gap:
# 中断当前 run
run_end = i - gap_count # 最后一个匹配帧的索引
# 计算 run 内的实际匹配帧数(总跨度 - 间隙数)
total_gaps = sum(1 for k in range(run_start, run_end + 1) if not frame_matches[k][0])
matching_count = (run_end - run_start + 1) - total_gaps
if matching_count >= min_consecutive:
runs.append((run_start, run_end))
run_start = None
gap_count = 0
# 处理末尾 run
if run_start is not None:
last_idx = len(frame_matches) - 1
# 回退找到最后一个匹配帧的位置(跳过尾部非匹配帧)
while last_idx >= run_start and not frame_matches[last_idx][0]:
last_idx -= 1
if last_idx >= run_start:
# 计算 run 内的总间隙数
total_gaps = sum(1 for k in range(run_start, last_idx + 1) if not frame_matches[k][0])
matching_count = (last_idx - run_start + 1) - total_gaps
if matching_count >= min_consecutive:
runs.append((run_start, last_idx))
# Step 3: 构建 DuplicateSegment
segments: list[DuplicateSegment] = []
for start, end in runs:
query_start = _get_start(query_chunks[start])
query_end = _get_end(query_chunks[end])
# 取目标范围(按最佳匹配的目标 chunk 时间范围)
target_indices = [frame_matches[k][2] for k in range(start, end + 1) if frame_matches[k][0]]
if target_indices:
t_min = min(target_indices)
t_max = max(target_indices)
target_start = _get_start(target_chunks[t_min])
target_end = _get_end(target_chunks[t_max])
else:
target_start = _get_start(target_chunks[0])
target_end = _get_end(target_chunks[-1])
avg_dist = sum(frame_matches[k][1] for k in range(start, end + 1)) / (end - start + 1)
segments.append(
DuplicateSegment(
query_start_ms=query_start,
query_end_ms=query_end,
target_start_ms=target_start,
target_end_ms=target_end,
avg_distance=avg_dist,
)
)
return segments
# ── VideoDeduplicator ───────────────────────────────────────────
class VideoDeduplicator:
"""Video deduplication using multiple fingerprint methods."""
@@ -170,11 +421,11 @@ class VideoDeduplicator:
HISTOGRAM_THRESHOLD = 0.85
def compute_fingerprint(self, video_path: str) -> VideoFingerprint:
"""Compute video fingerprint using MD5, pHash, and color histogram.
"""Compute video fingerprint using dynamic keyframe detection.
按时间分片抽帧:短视频(≤60s)每 2s 一片,长视频每 5s 一片。
每片取 1 帧计算 pHash + color_histogram。
同时保留 keyframe_phashes/color_histograms 聚合字段(向后兼容)。
使用 detect_keyframe_timestamps() 检测内容感知关键帧,
在每个关键帧处取帧计算 pHash + color_histogram。
同时保留 MD5 计算和分片数据结构。
"""
cap = cv2.VideoCapture(video_path)
if not cap.isOpened():
@@ -186,41 +437,55 @@ class VideoDeduplicator:
width = int(cap.get(cv2.CAP_PROP_FRAME_WIDTH))
height = int(cap.get(cv2.CAP_PROP_FRAME_HEIGHT))
cap.release()
# 1. 检测关键帧时间戳
keyframe_times = detect_keyframe_timestamps(video_path)
if not keyframe_times:
return VideoFingerprint(
md5="",
keyframe_phashes=[],
color_histograms=[],
duration=duration,
resolution=(width, height),
chunks=[],
)
# 2. 打开视频,逐个关键帧取帧
cap = cv2.VideoCapture(video_path)
md5_hash = hashlib.md5(usedforsecurity=False)
chunks: list[FingerprintChunk] = []
# 分片间隔(秒)
chunk_interval_sec = compute_chunk_interval(duration)
chunk_interval_ms = int(chunk_interval_sec * 1000)
duration_ms = int(duration * 1000)
# 遍历每个分片时间窗口,取 1 帧
start_ms = 0
while start_ms < duration_ms:
end_ms = min(start_ms + chunk_interval_ms, duration_ms)
# 定位到分片中点
seek_ms = (start_ms + end_ms) / 2
for i, t_sec in enumerate(keyframe_times):
seek_ms = t_sec * 1000
cap.set(cv2.CAP_PROP_POS_MSEC, seek_ms)
ret, frame = cap.read()
if ret:
# MD5 计算
_, buffer = cv2.imencode(".jpg", frame)
md5_hash.update(buffer)
if not ret:
continue
phash = compute_phash(frame)
hist = compute_color_histogram(frame)
# MD5 计算
_, buffer = cv2.imencode(".jpg", frame)
md5_hash.update(buffer)
chunks.append(
FingerprintChunk(
start_time_ms=start_ms,
end_time_ms=end_ms,
phash_binary=phash,
color_histogram=hist,
frame_count=1,
)
phash = compute_phash(frame)
hist = compute_color_histogram(frame)
# 计算分片时间范围(从前一个关键帧到下一个关键帧的中点)
prev_boundary = keyframe_times[i - 1] * 1000 if i > 0 else 0
next_boundary = keyframe_times[i + 1] * 1000 if i < len(keyframe_times) - 1 else duration * 1000
start_ms = int((prev_boundary + seek_ms) / 2)
end_ms = int((seek_ms + next_boundary) / 2)
chunks.append(
FingerprintChunk(
start_time_ms=start_ms,
end_time_ms=end_ms,
phash_binary=phash,
color_histogram=hist,
frame_count=1,
)
start_ms = end_ms
)
cap.release()
@@ -255,14 +520,39 @@ class VideoDeduplicator:
for r in rows
]
@staticmethod
def _bhattacharyya_coefficient(hist_a: list[float], hist_b: list[float]) -> float:
"""Bhattacharyya 系数:Σ √(a[i] * b[i]),范围 [0, 1],1=完全相同。"""
min_len = min(len(hist_a), len(hist_b))
a = hist_a[:min_len]
b = hist_b[:min_len]
return float(sum(np.sqrt(ai * bi) for ai, bi in zip(a, b)))
@staticmethod
def _compute_histogram_similarity(
histograms_a: list[list[float]],
histograms_b: list[list[float]],
) -> float:
"""对每组直方图,找到最佳匹配的 Bhattacharyya 系数,取平均。"""
if not histograms_a or not histograms_b:
return 0.0
similarities = []
for ha in histograms_a:
best = 0.0
for hb in histograms_b:
bc = VideoDeduplicator._bhattacharyya_coefficient(ha, hb)
best = max(best, bc)
similarities.append(best)
return sum(similarities) / len(similarities) if similarities else 0.0
def check_duplicate(self, fingerprint: VideoFingerprint, project_id: str, session: Session) -> Optional[dict]:
"""检查视频是否与项目中已有视频重复。
查重逻辑:
1. MD5 精确匹配 → similarity=1.0
2. pHash 相似度(优先从分片表读取,回退到 JSON 字段)
2. pHash 中位数距离 + 帧匹配比例 + 直方图融合判定
判定阈值:avg_distance < PHASH_THRESHOLD(10)
判定为重复后,调用 find_duplicate_segments() 获取具体重复片段。
Args:
fingerprint: 待检测视频的指纹
@@ -270,7 +560,7 @@ class VideoDeduplicator:
session: 数据库会话
Returns:
重复信息字典(含 duplicate, duplicate_of, reason, similarity),
重复信息字典(含 duplicate, duplicate_of, reason, similarity, duplicate_segments),
或 None 表示未找到重复。
"""
video_repo = SQLAlchemyGeneratedVideoRepository(session)
@@ -298,23 +588,64 @@ class VideoDeduplicator:
if not existing_phashes:
continue
# 计算每个新关键帧到已有关键帧的最小汉明距离,取平均
# 计算每个新关键帧到已有关键帧的最小汉明距离
min_distances = []
for phash in fingerprint.keyframe_phashes:
distances = [hamming_distance(phash, ep) for ep in existing_phashes]
min_distances.append(min(distances))
avg_distance = sum(min_distances) / len(min_distances) if min_distances else 100
if avg_distance >= self.PHASH_THRESHOLD:
# 帧匹配比例检查
matching_frames = sum(1 for d in min_distances if d < self.PHASH_THRESHOLD)
match_ratio = matching_frames / len(min_distances) if min_distances else 0
if match_ratio < 0.7:
continue
phash_similarity = 1.0 - (avg_distance / 64)
# 中位数距离
median_distance = statistics.median(min_distances) if min_distances else 64
if median_distance >= self.PHASH_THRESHOLD:
continue
# 直方图融合
existing_histograms = []
if chunk_data:
existing_histograms = [c["color_histogram"] for c in chunk_data if c.get("color_histogram")]
else:
existing_histograms = ef.get("color_histograms", [])
phash_similarity = 1.0 - (median_distance / 64)
hist_similarity = (
self._compute_histogram_similarity(fingerprint.color_histograms, existing_histograms)
if existing_histograms
else 0.5
)
combined_score = 0.7 * phash_similarity + 0.3 * hist_similarity
# DUPLICATE_THRESHOLD from module level
if combined_score < DUPLICATE_THRESHOLD:
continue
# 滑动窗口时序匹配:获取具体重复片段
existing_chunk_objects = chunk_data if chunk_data else [
{"phash_binary": p, "start_time_ms": 0, "end_time_ms": 0}
for p in existing_phashes
]
segments = find_duplicate_segments(fingerprint.chunks, existing_chunk_objects)
return {
"duplicate": True,
"duplicate_of": existing.id,
"reason": "phash_similar",
"similarity": phash_similarity,
"reason": "phash_histogram_fusion",
"similarity": combined_score,
"duplicate_segments": [
{
"query_start_ms": s.query_start_ms,
"query_end_ms": s.query_end_ms,
"target_start_ms": s.target_start_ms,
"target_end_ms": s.target_end_ms,
"avg_distance": round(s.avg_distance, 2),
}
for s in segments
],
}
return None
@@ -328,7 +659,8 @@ class VideoDeduplicator:
) -> Optional[dict]:
"""检查视频是否与同批次内其他视频重复。
逻辑与 check_duplicate 一致(MD5 + pHash),但搜索范围限定为同 batch_id 的视频。
逻辑与 check_duplicate 一致(MD5 + pHash + 直方图融合 + 时序匹配),
但搜索范围限定为同 batch_id 的视频。
Args:
fingerprint: 待检测视频的指纹
@@ -373,59 +705,62 @@ class VideoDeduplicator:
for phash in fingerprint.keyframe_phashes:
distances = [hamming_distance(phash, ep) for ep in existing_phashes]
min_distances.append(min(distances))
avg_distance = sum(min_distances) / len(min_distances) if min_distances else 100
if avg_distance >= self.PHASH_THRESHOLD:
# 帧匹配比例检查
matching_frames = sum(1 for d in min_distances if d < self.PHASH_THRESHOLD)
match_ratio = matching_frames / len(min_distances) if min_distances else 0
if match_ratio < 0.7:
continue
phash_similarity = 1.0 - (avg_distance / 64)
median_distance = statistics.median(min_distances) if min_distances else 64
if median_distance >= self.PHASH_THRESHOLD:
continue
# 直方图融合
existing_histograms = []
if chunk_data:
existing_histograms = [c["color_histogram"] for c in chunk_data if c.get("color_histogram")]
else:
existing_histograms = ef.get("color_histograms", [])
phash_similarity = 1.0 - (median_distance / 64)
hist_similarity = (
self._compute_histogram_similarity(fingerprint.color_histograms, existing_histograms)
if existing_histograms
else 0.5
)
combined_score = 0.7 * phash_similarity + 0.3 * hist_similarity
# DUPLICATE_THRESHOLD from module level
if combined_score < DUPLICATE_THRESHOLD:
continue
# 滑动窗口时序匹配
existing_chunk_objects = chunk_data if chunk_data else [
{"phash_binary": p, "start_time_ms": 0, "end_time_ms": 0}
for p in existing_phashes
]
segments = find_duplicate_segments(fingerprint.chunks, existing_chunk_objects)
return {
"duplicate": True,
"duplicate_of": existing.id,
"reason": "batch_phash_similar",
"similarity": phash_similarity,
"reason": "batch_phash_histogram_fusion",
"similarity": combined_score,
"duplicate_segments": [
{
"query_start_ms": s.query_start_ms,
"query_end_ms": s.query_end_ms,
"target_start_ms": s.target_start_ms,
"target_end_ms": s.target_end_ms,
"avg_distance": round(s.avg_distance, 2),
}
for s in segments
],
}
return None
@staticmethod
def _average_histogram_similarity(histograms_a: list[list[float]], histograms_b: list[list[float]]) -> float:
"""
计算两组颜色直方图之间的平均余弦相似度。
对每组直方图对取最小长度对齐,计算余弦相似度后取平均。
Args:
histograms_a: 第一组直方图(每帧一个 list)
histograms_b: 第二组直方图
Returns:
平均余弦相似度,范围 [0, 1]
"""
if not histograms_a or not histograms_b:
return 0.0
similarities = []
for ha in histograms_a:
best = 0.0
vec_a = np.array(ha, dtype=np.float64)
norm_a = np.linalg.norm(vec_a)
if norm_a == 0:
continue
for hb in histograms_b:
vec_b = np.array(hb, dtype=np.float64)
# 对齐长度
min_len = min(len(vec_a), len(vec_b))
va, vb = vec_a[:min_len], vec_b[:min_len]
norm_b = np.linalg.norm(vb)
if norm_b == 0:
continue
sim = float(np.dot(va, vb) / (norm_a * norm_b))
best = max(best, sim)
similarities.append(best)
return sum(similarities) / len(similarities) if similarities else 0.0
def compute_duplicate_rate(
self,
fingerprint: VideoFingerprint,
@@ -438,9 +773,9 @@ class VideoDeduplicator:
"""计算当前视频与用户库内已有视频的最高相似度百分比。
优先按 user_id 全局比较(跨项目),user_id 为空时回退到项目级比较。
遍历最近 200 个其他有指纹的视频,对每个计算相似度:
遍历最近 200 个其他有指纹的视频,对每个计算融合相似度:
- MD5 精确匹配 → 100%
- pHash 相似度 → (1.0 - avg_distance / 64) * 100
- pHash + 直方图融合 → 0.7 * phash_sim + 0.3 * hist_sim
取最高值作为 duplicate_rate(0~100)。
如果没有其他视频可比较,返回 0.0。
@@ -454,7 +789,6 @@ class VideoDeduplicator:
Returns:
duplicate_rate: 0~100 的浮点数
"""
# 限制查询最近 200 个视频,避免大库内存溢出
from packages.adapters.sqlalchemy_impl.models import GeneratedVideoModel
# 优先按 user_id 全局比较(跨项目),否则回退到项目级
@@ -469,7 +803,7 @@ class VideoDeduplicator:
)
logger.debug("compute_duplicate_rate: project-level fallback project_id=%s", project_id)
# 排除当前视频自身(记录可能已写入 DB,必须在查询层排除)
# 排除当前视频自身
if current_video_id:
query = query.filter(GeneratedVideoModel.id != current_video_id)
@@ -505,9 +839,30 @@ class VideoDeduplicator:
for phash in fingerprint.keyframe_phashes:
distances = [hamming_distance(phash, ep) for ep in existing_phashes]
min_distances.append(min(distances))
avg_distance = sum(min_distances) / len(min_distances) if min_distances else 64
similarity = (1.0 - avg_distance / 64) * 100
max_similarity = max(max_similarity, similarity)
# 帧匹配比例检查
matching_frames = sum(1 for d in min_distances if d < self.PHASH_THRESHOLD)
match_ratio = matching_frames / len(min_distances) if min_distances else 0
if match_ratio < 0.7:
continue
median_distance = statistics.median(min_distances) if min_distances else 64
# 直方图融合
existing_histograms = []
if chunk_data:
existing_histograms = [c["color_histogram"] for c in chunk_data if c.get("color_histogram")]
else:
existing_histograms = ef.get("color_histograms", [])
phash_similarity = (1.0 - median_distance / 64) * 100
hist_similarity = (
self._compute_histogram_similarity(fingerprint.color_histograms, existing_histograms) * 100
if existing_histograms
else 50.0
)
combined_score = 0.7 * phash_similarity + 0.3 * hist_similarity
max_similarity = max(max_similarity, combined_score)
return round(max(max_similarity, 0.0), 2)