dedup/dedup.py
Greg Pomerantz 434330be1d Initial release: parallel duplicate finder with checkpointing
- SHA-256 hashing with multiprocessing (configurable workers)
- Resumable via checkpoints (saves every N files, auto-resumes)
- Reports exact matches, hash-only matches, same-name/different-content, unique
- --move (quarantine to __duplicates__/) or --delete
- Python 3.8+, stdlib only
2026-09-30 17:06:56 -04:00

494 lines
18 KiB
Python

#!/usr/bin/env python3
"""
dedup - Find and remove duplicate files between two directories.
Uses SHA-256 hashing with multiprocessing for speed.
Supports resumable runs via checkpoints.
Python 3.8+, standard library only.
"""
import argparse
import hashlib
import json
import multiprocessing as mp
import os
import shutil
import sys
import time
from collections import defaultdict
from pathlib import Path
__version__ = "1.0.0"
DEFAULT_EXTENSIONS = {
".jpg", ".jpeg", ".png", ".gif", ".bmp", ".webp",
".mov", ".mp4", ".3gp", ".m4v", ".tif", ".tiff",
".cr2", ".cr3", ".nef", ".arw", ".raw", ".orf", ".sr2",
".avi", ".mkv", ".flv", ".wmv", ".webm",
".psd", ".ai", ".svg", ".eps",
}
SKIP_FILES = {".stfolder", ".DS_Store", "Thumbs.db", "desktop.ini"}
CHECKPOINT_DIR = ".dedup_checkpoint"
# ---------------------------------------------------------------------------
# Hashing
# ---------------------------------------------------------------------------
def sha256(filepath):
"""Return the SHA-256 hex digest of a file."""
h = hashlib.sha256()
with open(filepath, "rb") as f:
for chunk in iter(lambda: f.read(1 << 20), b""):
h.update(chunk)
return h.hexdigest()
def hash_batch(file_paths):
"""Worker: hash a batch of files. Returns (results_dict, errors_list)."""
results = {}
errors = []
for fp in file_paths:
try:
results[fp] = sha256(fp)
except Exception as e:
errors.append((fp, str(e)))
return results, errors
def hash_files_parallel(file_list, num_workers=None, worker_batch=50):
"""Hash files in parallel. Returns dict mapping filepath -> hash."""
if not file_list:
return {}
if num_workers is None:
num_workers = mp.cpu_count()
batches = [file_list[i:i + worker_batch] for i in range(0, len(file_list), worker_batch)]
with mp.Pool(processes=num_workers) as pool:
raw_results = pool.map(hash_batch, batches)
results = {}
for batch_results, batch_errors in raw_results:
results.update(batch_results)
for fp, err in batch_errors:
print(f" WARNING: could not hash {fp}: {err}", file=sys.stderr)
return results
# ---------------------------------------------------------------------------
# File walking
# ---------------------------------------------------------------------------
def walk_files(directory, recurse=True, extensions=None):
"""Yield file paths matching the given extensions (case-insensitive)."""
directory = Path(directory)
if not directory.is_dir():
print(f"ERROR: {directory} is not a directory", file=sys.stderr)
sys.exit(1)
if extensions is None:
extensions = DEFAULT_EXTENSIONS
if recurse:
for root, dirs, files in os.walk(directory):
dirs[:] = [d for d in dirs if d != CHECKPOINT_DIR]
for fname in files:
fpath = Path(root) / fname
if fname in SKIP_FILES:
continue
if fpath.suffix.lower() in extensions:
try:
fpath.stat()
except OSError:
continue
yield str(fpath)
else:
for item in directory.iterdir():
if item.is_file() and item.name not in SKIP_FILES:
if item.suffix.lower() in extensions:
try:
item.stat()
except OSError:
continue
yield str(item)
# ---------------------------------------------------------------------------
# Checkpointing
# ---------------------------------------------------------------------------
def get_checkpoint_path(dir_a):
return Path(dir_a) / CHECKPOINT_DIR
def save_checkpoint(cp_path, phase, hash_map, name_map, file_list, processed_set, batch_num, total_files):
"""Save checkpoint to disk."""
cp_path.mkdir(parents=True, exist_ok=True)
meta = {
"phase": phase,
"batch_num": batch_num,
"total_files": total_files,
"files_processed": len(processed_set),
"timestamp": time.time(),
}
(cp_path / "meta.json").write_text(json.dumps(meta, indent=2))
(cp_path / "processed.json").write_text(json.dumps(sorted(processed_set)))
(cp_path / f"{phase}_hash.json").write_text(json.dumps(dict(hash_map)))
(cp_path / f"{phase}_name.json").write_text(json.dumps(dict(name_map)))
(cp_path / "file_list.json").write_text(json.dumps(file_list))
def load_checkpoint(cp_path):
"""Load checkpoint. Returns (meta, hash_map, name_map, file_list, processed_set) or None."""
if not cp_path.exists() or not (cp_path / "meta.json").exists():
return None
try:
meta = json.loads((cp_path / "meta.json").read_text())
processed_set = set(json.loads((cp_path / "processed.json").read_text()))
hash_map = defaultdict(list, json.loads((cp_path / f"{meta['phase']}_hash.json").read_text()))
name_map = defaultdict(list, json.loads((cp_path / f"{meta['phase']}_name.json").read_text()))
file_list = json.loads((cp_path / "file_list.json").read_text())
return meta, hash_map, name_map, file_list, processed_set
except (json.JSONDecodeError, FileNotFoundError, KeyError):
return None
def clean_checkpoint(dir_a):
"""Remove all checkpoint data."""
cp_path = get_checkpoint_path(dir_a)
if cp_path.exists():
shutil.rmtree(cp_path)
print(f" Cleaned checkpoint: {cp_path}")
# ---------------------------------------------------------------------------
# Indexing
# ---------------------------------------------------------------------------
def index_dir(directory, recurse, extensions, checkpoint_every, phase, cp_path,
num_workers=None, file_list=None, resume_from=None,
existing_hash_map=None, existing_name_map=None):
"""
Index a directory with parallel hashing and periodic checkpointing.
Returns (hash_map, name_map, count).
"""
hash_map = defaultdict(list, existing_hash_map or {})
name_map = defaultdict(list, existing_name_map or {})
if resume_from is None:
resume_from = set()
if file_list is None:
all_files = list(walk_files(directory, recurse, extensions))
print(f" Found {len(all_files)} files to hash")
else:
all_files = [fp for fp in file_list if Path(fp).exists()]
print(f" Loaded {len(all_files)} files from checkpoint list")
remaining = [fp for fp in all_files if fp not in resume_from]
already_done = len(all_files) - len(remaining)
if already_done > 0:
print(f" Resuming: {already_done} hashed, {len(remaining)} remaining")
if not remaining:
print(" All files already processed.")
return hash_map, name_map, len(all_files)
total_all = len(all_files)
processed_count = already_done
batch_num = 0
files_since_checkpoint = 0
# Hash in sub-batches of 200, checkpoint every `checkpoint_every` files
for i in range(0, len(remaining), 200):
batch = remaining[i:i + 200]
batch_results = hash_files_parallel(batch, num_workers=num_workers)
for fp, h in batch_results.items():
hash_map[h].append(fp)
name_map[Path(fp).name].append(fp)
processed_count += len(batch_results)
total_done = already_done + len(batch_results)
batch_num += 1
files_since_checkpoint += len(batch_results)
if total_done % 5000 == 0 or total_done == total_all:
print(f" Hashed {total_done}/{total_all} files...")
if checkpoint_every > 0 and files_since_checkpoint >= checkpoint_every:
processed = resume_from | set(remaining[:i + len(batch)])
save_checkpoint(cp_path, phase, dict(hash_map), dict(name_map),
all_files, processed, batch_num, total_all)
files_since_checkpoint = 0
# Final checkpoint
if checkpoint_every > 0:
processed = resume_from | set(remaining)
save_checkpoint(cp_path, phase, dict(hash_map), dict(name_map),
all_files, processed, batch_num, total_all)
return hash_map, name_map, processed_count
# ---------------------------------------------------------------------------
# Comparison
# ---------------------------------------------------------------------------
def compare_directories(a_hash_map, a_name_map, b_hash_map, b_name_map):
"""
Compare dirA files against dirB.
Returns (exact_matches, hash_matches, same_name_diff, unique_in_a).
"""
exact_matches = []
hash_matches = []
same_name_diff = []
unique_in_a = []
# Build path -> hash for dirA
dir_a_files = {}
for h, paths in a_hash_map.items():
for p in paths:
dir_a_files[p] = h
for a_path in sorted(dir_a_files.keys()):
a_hash = dir_a_files[a_path]
a_name = Path(a_path).name
b_same_name = b_name_map.get(a_name, [])
b_same_hash = b_hash_map.get(a_hash, [])
if b_same_name:
name_match_found = False
for b_path in b_same_name:
if b_path in b_same_hash:
exact_matches.append((a_path, b_path, a_hash))
name_match_found = True
break
if not name_match_found:
for b_path in b_same_name:
same_name_diff.append((a_path, a_hash, b_path))
break
elif b_same_hash:
hash_matches.append((a_path, b_same_hash[0], a_hash))
else:
unique_in_a.append((a_path, a_hash))
return exact_matches, hash_matches, same_name_diff, unique_in_a
# ---------------------------------------------------------------------------
# Reporting
# ---------------------------------------------------------------------------
def write_report(report_path, exact, hash_only, same_name_diff, unique, dir_a, dir_b):
with open(report_path, "w") as f:
f.write("=" * 80 + "\n")
f.write("DUPLICATE FILES REPORT\n")
f.write(f"Dir A: {dir_a}\nDir B: {dir_b}\n")
f.write("=" * 80 + "\n\n")
f.write(f"--- EXACT MATCHES ({len(exact)}) ---\n\n")
for a, b, h in exact:
f.write(f" A: {a}\n B: {b}\n Hash: {h}\n\n")
f.write(f"--- HASH MATCHES / DIFFERENT NAMES ({len(hash_only)}) ---\n\n")
for a, b, h in hash_only:
f.write(f" A: {a}\n B: {b}\n Hash: {h}\n\n")
f.write(f"--- SAME NAME, DIFFERENT CONTENT ({len(same_name_diff)}) ---\n\n")
for a, ah, b in same_name_diff:
f.write(f" A: {a}\n B: {b}\n A hash: {ah}\n\n")
f.write(f"--- UNIQUE TO DIR A ({len(unique)}) ---\n\n")
for a, h in unique:
f.write(f" {a} ({h[:16]}...)\n")
# ---------------------------------------------------------------------------
# Actions
# ---------------------------------------------------------------------------
def move_duplicates(duplicates):
"""Move duplicates into a __duplicates__/ quarantine folder."""
if not duplicates:
return
quarantine = Path(duplicates[0][0]).parent / "__duplicates__"
quarantine.mkdir(exist_ok=True)
moved = errs = 0
for i, (a_path, _, _) in enumerate(duplicates, 1):
try:
shutil.move(a_path, str(quarantine / Path(a_path).name))
moved += 1
except Exception as e:
print(f" ERROR: {a_path}: {e}", file=sys.stderr)
errs += 1
if i % 500 == 0 or i == len(duplicates):
print(f" Moved {i}/{len(duplicates)}...")
print(f" Moved {moved} files to {quarantine}/")
if errs:
print(f" {errs} errors.", file=sys.stderr)
def delete_duplicates(duplicates):
"""Permanently delete duplicate files."""
deleted = errs = 0
for i, (a_path, _, _) in enumerate(duplicates, 1):
try:
os.unlink(a_path)
deleted += 1
except Exception as e:
print(f" ERROR: {a_path}: {e}", file=sys.stderr)
errs += 1
if i % 500 == 0 or i == len(duplicates):
print(f" Deleted {i}/{len(duplicates)}...")
print(f" Deleted {deleted} files")
if errs:
print(f" {errs} errors.", file=sys.stderr)
# ---------------------------------------------------------------------------
# Main
# ---------------------------------------------------------------------------
def main():
parser = argparse.ArgumentParser(
prog="dedup",
description="Find and optionally remove duplicate files between two directories.",
formatter_class=argparse.RawDescriptionHelpFormatter,
epilog=f"""
version: {__version__}
python: {sys.version.split()[0]} (stdlib only)
Examples:
dedup.py /photos/backup /photos/original
dedup.py --move /photos/backup /photos/original
dedup.py --delete /photos/backup /photos/original
dedup.py --workers 4 --checkpoint-every 500 /large/backup /original
dedup.py --clean /photos/backup /photos/original
Checkpoints are saved to dirA/.dedup_checkpoint/ every --checkpoint-every files.
Re-run the same command to resume. Use --clean to discard checkpoints.
""")
parser.add_argument("dir_a", help="Source directory (duplicates found here)")
parser.add_argument("dir_b", help="Reference directory (kept)")
parser.add_argument("--move", action="store_true",
help="Move duplicates to __duplicates__/ quarantine folder")
parser.add_argument("--delete", action="store_true",
help="Permanently delete duplicates (irreversible!)")
parser.add_argument("--no-subdirs", action="store_true",
help="Only compare top-level files (no recursion)")
parser.add_argument("--ext", default=None,
help="Comma-separated extensions to include (default: all media)")
parser.add_argument("--report", default=None,
help="Report output path (default: dirA/duplicates_report.txt)")
parser.add_argument("--workers", type=int, default=None,
help="Number of parallel hash workers (default: CPU count)")
parser.add_argument("--checkpoint-every", type=int, default=1000,
help="Files between checkpoints, 0 to disable (default: 1000)")
parser.add_argument("--clean", action="store_true",
help="Remove checkpoint data and exit")
parser.add_argument("--version", action="version", version=f"%(prog)s {__version__}")
args = parser.parse_args()
# Validate
if args.move and args.delete:
print("ERROR: specify only one of --move or --delete", file=sys.stderr)
sys.exit(1)
if not Path(args.dir_a).is_dir():
print(f"ERROR: {args.dir_a} not a directory", file=sys.stderr)
sys.exit(1)
if not Path(args.dir_b).is_dir():
print(f"ERROR: {args.dir_b} not a directory", file=sys.stderr)
sys.exit(1)
extensions = set(args.ext.split(",")) if args.ext else DEFAULT_EXTENSIONS
recurse = not args.no_subdirs
cp_path = get_checkpoint_path(args.dir_a)
if args.clean:
clean_checkpoint(args.dir_a)
sys.exit(0)
# ==========================================================
# Phase 1: Index dirB (reference)
# ==========================================================
print(f"=== Phase 1: Indexing reference ({args.dir_b}) ===")
b_cp = load_checkpoint(cp_path / "b")
if b_cp:
meta, b_hash, b_name, b_file_list, b_processed = b_cp
print(f" Resuming: dirB already indexed ({len(b_hash)} unique hashes)")
else:
b_hash, b_name, b_count = index_dir(
args.dir_b, recurse, extensions, args.checkpoint_every, "b",
cp_path / "b", num_workers=args.workers)
print(f" Indexed {b_count} files, {len(b_hash)} unique hashes")
# ==========================================================
# Phase 2: Index dirA (source)
# ==========================================================
print(f"\n=== Phase 2: Indexing source ({args.dir_a}) ===")
a_cp = load_checkpoint(cp_path / "a")
resume_from = None
file_list = None
a_hash_existing = None
a_name_existing = None
if a_cp:
meta, a_hash_existing, a_name_existing, file_list, resume_from = a_cp
print(f" Resuming: {meta['files_processed']} files already processed")
a_hash, a_name, a_count = index_dir(
args.dir_a, recurse, extensions, args.checkpoint_every, "a",
cp_path / "a", num_workers=args.workers,
file_list=file_list, resume_from=resume_from,
existing_hash_map=a_hash_existing, existing_name_map=a_name_existing)
print(f" Indexed {a_count} files, {len(a_hash)} unique hashes")
# ==========================================================
# Phase 3: Compare
# ==========================================================
print("\n=== Phase 3: Comparing directories ===")
exact, hash_only, same_name_diff, unique = compare_directories(
a_hash, a_name, b_hash, b_name)
total_dups = len(exact) + len(hash_only)
print(f"\n{'=' * 60}")
print(f"SUMMARY")
print(f"{'=' * 60}")
print(f" Exact matches (same name + hash): {len(exact)}")
print(f" Hash matches (different name): {len(hash_only)}")
print(f" Total duplicates: {total_dups}")
print(f" Same name, different content: {len(same_name_diff)}")
print(f" Unique to dirA: {len(unique)}")
report_path = Path(args.report) if args.report else Path(args.dir_a) / "duplicates_report.txt"
write_report(report_path, exact, hash_only, same_name_diff, unique, args.dir_a, args.dir_b)
print(f"\nReport: {report_path}")
clean_checkpoint(args.dir_a)
# ==========================================================
# Phase 4: Action
# ==========================================================
all_dups = exact + hash_only
if args.delete:
print(f"\n=== Phase 4: DELETING {len(all_dups)} duplicates ===")
delete_duplicates(all_dups)
elif args.move:
print(f"\n=== Phase 4: MOVING {len(all_dups)} duplicates ===")
move_duplicates(all_dups)
else:
if all_dups:
print(f"\nNo action taken. Use --move or --delete to remove {len(all_dups)} duplicates.")
if __name__ == "__main__":
main()