How to Implement Breakpoint Resume for Large-Scale Crawling Jobs in MediaCrawler
MediaCrawler supports breakpoint resume through the existing --start CLI argument combined with a JSON checkpoint persistence layer that saves and restores the last processed page across interruptions.
This guide walks through building a robust checkpoint system on top of the NanmiCoder/MediaCrawler framework. By leveraging the existing config.start_page infrastructure, you can recover from crashes and network failures without re-scraping already-processed data.
Understanding the Built-in --start Mechanism
The foundation for breakpoint resume already exists in api/services/crawler_manager.py. The _build_command method (lines 22-27) constructs subprocess commands that include a --start <num> argument:
# api/services/crawler_manager.py (excerpt)
def _build_command(self, platform: str, start_page: int) -> List[str]:
return [
"python", "main.py",
"--platform", platform,
"--start", str(start_page), # <-- existing entry point
# ...
]
This --start value flows into main.py, where cmd_arg.parse_cmd() populates config.start_page. Platform crawlers then read this value to begin pagination from the specified page.
Architecture for Robust Checkpoint Resume
To transform this static start option into an automatic resume system, implement six coordinated components:
1. Create a Checkpoint Persistence Module
Add tools/checkpoint.py to handle JSON-based state storage. This module exposes three core operations: save, load, and clear.
# tools/checkpoint.py
import json
import pathlib
def checkpoint_path(platform: str) -> pathlib.Path:
"""Return platform-specific checkpoint file path."""
return pathlib.Path("checkpoints") / f"{platform}.json"
def save(platform: str, page: int) -> None:
"""Persist last successfully processed page."""
p = checkpoint_path(platform)
p.parent.mkdir(parents=True, exist_ok=True)
p.write_text(
json.dumps({"last_page": page, "platform": platform}),
encoding="utf-8"
)
def load(platform: str) -> int:
"""Restore last page from checkpoint, defaulting to 1."""
p = checkpoint_path(platform)
if not p.exists():
return 1
data = json.loads(p.read_text(encoding="utf-8"))
return data.get("last_page", 1)
def clear(platform: str) -> None:
"""Remove checkpoint after successful completion."""
p = checkpoint_path(platform)
if p.exists():
p.unlink()
2. Inject Checkpoints at Startup
Modify main.py to load the checkpoint before crawler initialization. This ensures config.start_page reflects the resume point:
# main.py (modified excerpt)
from tools.checkpoint import load as load_checkpoint, clear as clear_checkpoint
async def main() -> None:
# Parse CLI arguments first
await cmd_arg.parse_cmd()
# Override start_page with checkpoint if no explicit --start provided
if config.start_page == 1: # Only auto-resume when not manually specified
resume_page = load_checkpoint(config.PLATFORM)
config.start_page = resume_page
if resume_page > 1:
logger.info(f"Resuming {config.PLATFORM} from page {resume_page}")
crawler = CrawlerFactory.create_crawler(platform=config.PLATFORM)
await crawler.start()
# Clean checkpoint on successful completion
clear_checkpoint(config.PLATFORM)
3. Propagate Through CrawlerManager
No changes needed to CrawlerManager._build_command. The existing method already forwards --start to subprocesses. Simply ensure the API layer passes the checkpoint-loaded value.
4. Save Progress in Platform Crawlers
Each platform crawler must call save_checkpoint() after successfully storing a page's data. Here's the pattern for media_platform/xhs/core.py:
# media_platform/xhs/core.py (pagination loop excerpt)
from tools.checkpoint import save as save_checkpoint
async def fetch(self) -> None:
page = config.start_page
max_page = self._calculate_max_pages()
while page <= max_page:
items = await self._fetch_page(page)
if not items:
break
await self.store.save_items(items)
# Persist checkpoint after each successful page
save_checkpoint(config.PLATFORM, page)
page += 1
await asyncio.sleep(config.CRAWL_INTERVAL)
Apply this same pattern to all platform crawlers: media_platform/douyin/core.py, media_platform/weibo/core.py, etc.
5. Cleanup on Normal Exit
The async_cleanup function in main.py (lines 23-33) already handles final resource disposal. Extend it to ensure checkpoint cleanup occurs after storage flush:
# main.py (in async_cleanup)
async def async_cleanup():
tasks = [
SaveDataFactory.close(),
# ... other cleanup tasks
]
await asyncio.gather(*tasks, return_exceptions=True)
# Checkpoint already cleared in main() on success,
# but defensive cleanup here handles edge cases
from tools.checkpoint import clear as clear_checkpoint
clear_checkpoint(config.PLATFORM)
6. Add Manual Reset CLI Flag
For fresh starts, add --reset-checkpoint to cmd_arg/arg.py:
# cmd_arg/arg.py (addition)
parser.add_argument(
"--reset-checkpoint",
action="store_true",
help="Delete existing checkpoint and start from page 1"
)
Then handle it in main.py:
# main.py (startup logic)
if args.reset_checkpoint:
clear_checkpoint(config.PLATFORM)
config.start_page = 1
elif config.start_page == 1:
config.start_page = load_checkpoint(config.PLATFORM)
Platform-Specific Considerations
Different platforms in MediaCrawler implement pagination differently. Check your target platform's core module:
| Platform | File Path | Pagination Pattern |
|---|---|---|
| XiaoHongShu (XHS) | media_platform/xhs/core.py |
Cursor-based with page param |
| Douyin | media_platform/douyin/core.py |
Offset/limit or cursor |
media_platform/weibo/core.py |
Page number in URL | |
| Bilibili | media_platform/bilibili/core.py |
pn parameter |
All should respect config.start_page and call save_checkpoint() after batch persistence.
Handling Edge Cases
- Partial page failures: Save checkpoint only after
store.save_items()succeeds - Duplicate items: Store processed IDs in checkpoint JSON to enable item-level deduplication
- Multiple platforms: Each platform maintains independent checkpoint files
- Corrupted checkpoints: Validate JSON structure; fall back to page 1 on parse errors
Summary
- MediaCrawler's
--startargument provides the foundation for breakpoint resume - Add
tools/checkpoint.pyfor JSON persistence oflast_pageper platform - Load checkpoints in
main.pybefore crawler initialization - Call
save_checkpoint()after each successful page in platform crawlers - Clear checkpoints on successful completion via
async_cleanup - Add
--reset-checkpointCLI flag for manual fresh starts
Frequently Asked Questions
How does the checkpoint system handle concurrent crawls for different platforms?
Each platform uses an isolated checkpoint file named checkpoints/{platform}.json. The checkpoint_path() function in tools/checkpoint.py ensures no collision between simultaneous XHS, Douyin, or Weibo crawling jobs.
What happens if a page fetch succeeds but storage fails?
Do not call save_checkpoint() until after store.save_items() completes successfully. This guarantees that resuming will re-attempt the failed page, preventing data gaps. Wrap the persistence call in try/except and only save checkpoint on confirmed success.
Can I resume from an item-level position instead of just page numbers?
Yes. Extend the checkpoint schema to include processed_ids: List[str] alongside last_page. Modify platform crawlers to skip items whose IDs appear in this list. This adds overhead but prevents re-processing partial pages when item-level idempotency matters.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →