分布式进程崩溃后,如何基于MySQL与ZooKeeper恢复任务信息?
Great question—this is a super common pain point in distributed systems, and your stack (MySQL + ZooKeeper) gives you all the tools you need to build a solid recovery flow. Let’s break this down into actionable steps that play to each tool’s strengths.
Core Approach
The key is to:
- Track which processes are alive (ZooKeeper’s bread and butter)
- Maintain clear, transactional state for tasks in MySQL
- Safely claim and resume crashed processes’ tasks without duplicates or data loss
Step-by-Step Implementation
1. Use ZooKeeper for Process Liveness Tracking
Every time a process starts up, have it create a ephemeral node in ZooKeeper (e.g., /workers/worker-<unique-id>). Ephemeral nodes are automatically deleted when the process disconnects (crash, network split, etc.).
- Store a small payload in the node: the worker’s unique ID (could be a UUID or sequential ID from ZK) and maybe its current load.
- Set up a Watcher on the
/workerspath (either in a dedicated monitoring service or other worker processes) to detect when a node disappears. This is your first trigger for recovery.
2. Enhance MySQL Task Table for State Tracking
Update your task table to include these critical fields (on top of your existing task data):
worker_id: The unique ID of the process currently handling the tasktask_status: Enum with values likePENDING,RUNNING,RECOVERING,COMPLETED,FAILEDheartbeat_ts: Timestamp updated by the worker every N seconds (e.g., 10s) to prove it’s aliveprogress_details: Optional JSON field to store checkpoint data (e.g.,"last_processed_record_id": 1234) for resumable tasks
Example schema snippet:
ALTER TABLE tasks ADD COLUMN worker_id VARCHAR(64), ADD COLUMN task_status ENUM('PENDING', 'RUNNING', 'RECOVERING', 'COMPLETED', 'FAILED') DEFAULT 'PENDING', ADD COLUMN heartbeat_ts DATETIME, ADD COLUMN progress_details JSON;
3. Trigger Recovery When a Process Crashes
You have two complementary ways to detect a crashed worker:
- ZooKeeper Watcher Trigger: When a worker’s ephemeral node is deleted, the watcher immediately marks all tasks assigned to that worker as
RECOVERINGin MySQL (use a transaction to ensure atomicity). - MySQL Heartbeat Check: Run a periodic job (e.g., every 30s) that scans for
RUNNINGtasks whereheartbeat_tsis older than your threshold (e.g., 20s). Mark these asRECOVERING—this acts as a fallback if ZooKeeper has a temporary blip.
4. New Worker Claims and Resumes Tasks
When a new process starts up:
- Register itself in ZooKeeper (create an ephemeral node, get its unique
worker_id). - Query MySQL for tasks marked
RECOVERINGthat were assigned to the crashed worker. Use an optimistic lock to claim tasks safely:
This ensures only one worker can claim each task (MySQL’s row-level locking prevents race conditions).UPDATE tasks SET worker_id = '<new-worker-id>', task_status = 'RUNNING', heartbeat_ts = NOW() WHERE worker_id = '<crashed-worker-id>' AND task_status = 'RECOVERING'; - For each claimed task:
- If it’s a resumable task, load the
progress_detailsand pick up from the last checkpoint. - If it’s a non-resumable task (e.g., a one-time API call), verify it wasn’t already completed (check logs or an
is_executedflag) before re-running.
- If it’s a resumable task, load the
- Start sending regular heartbeats to update
heartbeat_tsfor all active tasks.
5. Guard Against Edge Cases
- Avoid Duplicate Execution: Always use the optimistic lock update when claiming tasks—never just select and update separately (that’s a race condition waiting to happen).
- Handle Partial Failures: If a worker crashes mid-task, make sure your task logic is idempotent (can be run multiple times without side effects) or uses checkpointing to resume from where it left off.
- ZooKeeper Split Brain: Since ZooKeeper uses quorum consensus, it’s resistant to split brain, but your recovery logic should rely on the ephemeral node deletion as the source of truth for worker liveness.
Bonus Optimizations
- Task Sharding: Assign tasks to workers using a consistent hash on
task_idandworker_id—this makes recovery faster because new workers only need to scan a subset of tasks. - Recovery Logging: Add a
task_recovery_logtable to track every recovery attempt (worker ID, timestamp, outcome) for debugging. - Auto-Scaling Integration: Tie the recovery trigger to your auto-scaler—if a worker crashes, spin up a new one automatically and kick off the task claim process.
内容的提问来源于stack exchange,提问作者hehe

