How batch_runner.py Implements Parallel Trajectory Generation in Hermes Agent

The batch_runner.py module in Hermes Agent orchestrates large-scale research runs by splitting JSONL datasets into batches, processing them across multiple worker processes using Python's multiprocessing.Pool, and merging the results into ShareGPT-compatible trajectory files with incremental checkpointing for fault tolerance.

The Hermes Agent repository by NousResearch provides a framework for generating high-quality conversational trajectories through tool-augmented agents. At the heart of its research infrastructure lies batch_runner.py, which enables parallel trajectory generation across distributed datasets. This module transforms raw prompt files into structured training data by leveraging multiprocessing to maximize throughput while maintaining data integrity through robust checkpointing mechanisms.

Batch Creation and Dataset Preparation

Before parallel execution begins, batch_runner.py loads the input dataset and segments it into manageable chunks. The BatchRunner class initializes by calling self.batches = self._create_batches() (lines 5053-5064), which groups the loaded dataset into lists of (index, entry) tuples sized according to the batch_size parameter.

This segmentation ensures that each worker process receives a balanced workload, preventing memory bottlenecks when processing large JSONL files containing thousands of research prompts.

Parallel Execution Architecture

The Multiprocessing Pool

The core parallelization mechanism utilizes Python's multiprocessing.Pool to distribute batch processing across CPU cores. The implementation creates a process pool of size self.num_workers (default 4) at line 7776:

with Pool(processes=self.num_workers) as pool:
    tasks = [
        (batch_num, batch_data, str(self.output_dir),
         completed_prompts_set, config)
        for batch_num, batch_data in enumerate(self.batches)
    ]
    for result in pool.imap_unordered(_process_batch_worker, tasks):
        # Progress tracking and checkpointing

The code constructs task tuples containing batch metadata, output directory paths, and configuration objects. Using imap_unordered allows the main process to consume results asynchronously as workers complete, enabling real-time progress monitoring through a Rich Progress bar (lines 9998-10010).

Worker Process Implementation

Each worker executes _process_batch_worker() (lines 8585-8989), which handles the actual trajectory generation for its assigned batch. The worker opens a dedicated output file batch_{batch_num}.jsonl and iterates through each prompt entry.

For every prompt, the worker calls _process_single_prompt() (lines 8629-8639), which:

  1. Samples toolset distributions using sample_toolsets_from_distribution (lines 8041-8043) from toolset_distributions.py
  2. Instantiates an AIAgent from run_agent.py with the sampled tools
  3. Executes agent.run_conversation() to generate the full message history
  4. Extracts tool statistics via _extract_tool_stats and reasoning coverage via _extract_reasoning_stats

The worker filters trajectories with zero reasoning steps (lines 9039-9045), ensuring only high-quality samples enter the final dataset. Valid trajectories are written as JSON lines (lines 9069-9072).

Incremental Checkpointing and Fault Tolerance

To support long-running research jobs, batch_runner.py implements robust checkpointing. After each batch completion, the main process updates a JSON checkpoint file:

checkpoint_data['completed_prompts'] = sorted(completed_prompts_set)
self._save_checkpoint(checkpoint_data, lock=checkpoint_lock)

The checkpoint lock ensures thread-safe writes from the main process (lines 9023-9028). The completed_prompts_set tracks which prompts have finished processing, allowing precise resume capabilities without duplicate work.

Merging and Post-Processing

When all workers complete, the run() method aggregates individual batch files into a unified output. It scans every batch_*.jsonl file, validates tool names against ALL_POSSIBLE_TOOLS from model_tools.py (line 542), filters corrupted entries, and writes the final trajectories.jsonl (lines 9892-10032).

The validation ensures that only tools defined in model_tools.TOOL_TO_TOOLSET_MAP appear in the final dataset, maintaining schema consistency for downstream training pipelines.

Resume Support for Interrupted Runs

The module provides seamless resume functionality through the --resume flag. When resuming, _scan_completed_prompts_by_content() (lines 7111-7173) reads existing batch_*.jsonl files to reconstruct the set of completed prompts. The runner then rebuilds the batch list, excluding already-processed entries (lines 8008-8032), ensuring zero duplicate generation when restarting crashed or preempted experiments.

Summary

  • Batch Processing: batch_runner.py splits large JSONL datasets into configurable batches using _create_batches() to balance memory usage and parallel efficiency.
  • Multiprocessing Architecture: Uses multiprocessing.Pool with imap_unordered to distribute trajectory generation across workers while maintaining real-time progress tracking via Rich.
  • Worker Implementation: Each worker instantiates AIAgent from run_agent.py, samples toolsets from distributions, and generates ShareGPT-style trajectories with extracted statistics.
  • Fault Tolerance: Implements JSON checkpointing with file locking to track completed prompts, supporting safe interruption and resumption of research runs.
  • Data Validation: Merges batch outputs while validating tool names against model_tools.TOOL_TO_TOOLSET_MAP to ensure schema-compliant training data.

Frequently Asked Questions

How does batch_runner.py handle failures in individual worker processes?

The multiprocessing.Pool implementation in batch_runner.py uses imap_unordered to yield results as workers complete. If a worker process crashes or raises an exception, the main process captures the error and continues processing remaining batches. The checkpointing mechanism tracks completed prompts via completed_prompts_set (lines 9023-9028), ensuring that successfully processed items are not re-attempted while failed items can be identified and reprocessed upon resume.

What determines how many trajectories are generated per prompt?

Each prompt generates exactly one trajectory through _process_single_prompt() (lines 8629-8639), but the complexity of that trajectory depends on the toolset distribution sampled via sample_toolsets_from_distribution (lines 8041-8043). The distribution configuration (passed as the distribution parameter) determines which subsets of tools from model_tools.py are available to the agent, directly influencing the conversation length and tool usage patterns captured in the final trajectory.

Can batch_runner.py resume a run if the server crashes?

Yes, batch_runner.py provides robust resume support through the --resume flag. When resuming, the runner calls _scan_completed_prompts_by_content() (lines 7111-7173) to read existing batch_*.jsonl files and reconstruct the set of completed prompts. It then rebuilds the batch list excluding already-processed entries (lines 8008-8032) and loads the checkpoint file to restore the exact state, allowing research runs to safely resume after hardware failures or preemption without duplicate generation.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →