Commit e60df3d0 authored by williamzhangNU's avatar williamzhangNU
Browse files

update dataset creation

parent 6ed4a33a
Loading
Loading
Loading
Loading
+10 −13
Original line number Diff line number Diff line
@@ -5,6 +5,7 @@ import os
import pandas as pd
import argparse
from pathlib import Path
from typing import Union, List

class DatasetCreator:

@@ -16,7 +17,7 @@ class DatasetCreator:
        
        

    def create_dataset(self, start_seed, train_size, test_size, force_gen=False):
    def create_dataset(self, seed: Union[int, List[int]], train_size, test_size, force_gen=False):
        train_file_path = os.path.join(self.data_dir, 'train.parquet')
        test_file_path = os.path.join(self.data_dir, 'test.parquet')
        
@@ -29,16 +30,12 @@ class DatasetCreator:
        # Ensure data directory exists
        os.makedirs(self.data_dir, exist_ok=True)
        
        seeds = range(start_seed, start_seed + train_size + test_size)
        instructions = []
        for seed in seeds:
            # Generate instruction based on environment
            # This is a placeholder - actual implementation would depend on the environment
            instruction = f"Instruction for seed {seed}"
            instructions.append(instruction)
        if isinstance(seed, int):
            seeds = range(seed, seed + train_size + test_size)
        else:
            seeds = seed
            
        def _create_instance(seed_idx, instruction):
            split = "train" if seed_idx < start_seed + train_size else "test"
        def _create_instance(seed_idx, split: str = 'train'):
            env_settings = {
                'env_name': self.env_name,
                'env_config': self.env_config,
@@ -48,12 +45,12 @@ class DatasetCreator:
            # TODO: no reward model defined here for the reward will be generated while rollout
            return {
                "data_source": self.env_name,
                "prompt": [{"role": "user", "content": instruction}],
                "prompt": [{"role": "user", "content": ''}],
                "extra_info": {"split": split, **env_settings}
            }

        train_instances = [_create_instance(start_seed + i, '') for i in range(train_size)]
        test_instances = [_create_instance(start_seed + train_size + i, '') for i in range(test_size)]
        train_instances = [_create_instance(seeds[i], split='train') for i in range(train_size)]
        test_instances = [_create_instance(seeds[train_size + i], split='test') for i in range(test_size)]
        
        train_dataset = Dataset.from_list(train_instances)
        test_dataset = Dataset.from_list(test_instances)
+123 −12
Original line number Diff line number Diff line
@@ -12,21 +12,109 @@ from tqdm import tqdm
from verl.utils.hdfs_io import copy, makedirs
import argparse
import datasets
import multiprocessing as mp
from functools import partial
import numpy as np

from vagen.env.create_dataset import DatasetCreator

from vagen.env.sokoban.env import SokobanInterface
from vagen.env.sokoban.room_utils import get_shortest_action_path, plot_animation
class SokobanDatasetCreator(DatasetCreator):

    def _process_seed(self, seed: int, max_action_length: int = 5):
        env_interface = SokobanInterface(**self.env_config)
        env_interface.reset(seed=seed)
        gt_action_sequence = get_shortest_action_path(
            env_interface.env.room_fixed, 
            env_interface.env.room_state, 
            MAX_DEPTH=max_action_length,
        )
        if len(gt_action_sequence) > max_action_length:
            return seed, []
        
        # images = []
        # obs = env_interface.env._render('rgb_array')
        # images.append(obs)
        # for action in gt_action_sequence:
        #     env_interface.env._step(action)
        #     obs = env_interface.env._render('rgb_array')
        #     images.append(obs)
        # animation = plot_animation(images)
        # animation.save(f'animation_{seed}.gif')

        return seed, gt_action_sequence
    
    def create_filtered_dataset(
        self,
        start_seed: int,
        train_size: int,
        test_size: int,
        max_steps: int = 5
        seed: int = 0,
        train_ratio: float = 0.8,
        max_action_length: int = 5,
        n_candidate: int = 20000,
        force_gen: bool = False,
    ):
        # TODO
        """
        Create a filtered dataset with given seeds
        """
        train_file_path = os.path.join(self.data_dir, 'train.parquet')
        test_file_path = os.path.join(self.data_dir, 'test.parquet')
        
        # Check if files already exist and force_gen is False
        if not force_gen and os.path.exists(train_file_path) and os.path.exists(test_file_path):
            print(f"Dataset files already exist at {self.data_dir}. Skipping generation.")
            print(f"Use --force-gen to override and regenerate the dataset.")
            return

        num_processes = mp.cpu_count()
        print(f"Using {num_processes} processes for seed processing")
        pool = mp.Pool(processes=num_processes)
        process_seed_partial = partial(self._process_seed, max_action_length=max_action_length)
        seeds = range(seed, seed + n_candidate)
        results = list(tqdm(pool.imap(process_seed_partial, seeds), total=len(seeds), desc="Processing seeds"))
        pool.close()
        pool.join()

        valid_seeds = [seed for seed, gt_action_sequence in results if gt_action_sequence and len(gt_action_sequence) <= max_action_length]
        train_size = int(len(valid_seeds) * train_ratio)
        test_size = len(valid_seeds) - train_size
        print(f"Train size: {train_size}, Test size: {test_size}")
        # Analyze statistics of action sequences
        action_lengths = [len(gt_action_sequence) for _, gt_action_sequence in results if gt_action_sequence]
        


        # Calculate basic statistics
        avg_length = np.mean(action_lengths) if action_lengths else 0
        median_length = np.median(action_lengths) if action_lengths else 0
        min_length = min(action_lengths) if action_lengths else 0
        max_length = max(action_lengths) if action_lengths else 0
        
        # Count frequency of each action length
        length_counts = {}
        for length in action_lengths:
            length_counts[length] = length_counts.get(length, 0) + 1
        
        # Calculate percentage of valid seeds
        valid_percentage = (len(valid_seeds) / n_candidate) * 100
        
        # Print statistics
        print("\nAction Sequence Statistics:")
        print(f"Total candidates processed: {n_candidate}")
        print(f"Valid seeds found: {len(valid_seeds)} ({valid_percentage:.2f}%)")
        print(f"Average action length: {avg_length:.2f}")
        print(f"Median action length: {median_length}")
        print(f"Min action length: {min_length}")
        print(f"Max action length: {max_length}")
        print("\nAction length distribution:")
        for length in sorted(length_counts.keys()):
            count = length_counts[length]
            percentage = (count / len(action_lengths)) * 100
            print(f"  Length {length}: {count} instances ({percentage:.2f}%)")



        self.create_dataset(valid_seeds, train_size, test_size, force_gen=force_gen)





@@ -35,11 +123,13 @@ class SokobanDatasetCreator(DatasetCreator):
if __name__ == "__main__":
    parser = argparse.ArgumentParser()
    parser.add_argument('--start_seed', type=int, default=0)
    parser.add_argument('--train_size', type=int, default=100)
    parser.add_argument('--test_size', type=int, default=100)
    parser.add_argument('--train_ratio', type=float, default=0.8)
    parser.add_argument('--n_candidate', type=int, default=20000)
    parser.add_argument('--max_action_length', type=int, default=None)
    parser.add_argument('--force-gen', action='store_true', 
                        help='Force dataset generation even if files already exist')
    # Added arguments based on the YAML config
    parser.add_argument('--data_dir', type=str, default='data/sokoban',)

    parser.add_argument('--dim_room', type=int, nargs=2, default=[6, 6],
                        help='Dimensions of the room [height, width]')
    parser.add_argument('--num_boxes', type=int, default=1,
@@ -50,7 +140,14 @@ if __name__ == "__main__":
                        help='Search depth that affects the starting position of the player')
    parser.add_argument('--visual_env', action='store_true',
                        help='Whether to use visual environment')
    parser.add_argument('--data_dir', type=str, default='data/sokoban',)
    
    import os
    if 'PYTHONHASHSEED' not in os.environ:
        os.environ['PYTHONHASHSEED'] = '0'
        print("Set PYTHONHASHSEED to 0 for reproducibility")
    else:
        print(f"PYTHONHASHSEED already set to {os.environ['PYTHONHASHSEED']}")
    

    args = parser.parse_args()
    args.name = 'sokoban'
@@ -62,5 +159,19 @@ if __name__ == "__main__":
        'visual_env': args.visual_env
    }
    creator = SokobanDatasetCreator(config=vars(args))
    #creator.create_filtered_dataset(start_seed=args.start_seed, train_size=args.train_size, test_size=args.test_size)
    creator.create_dataset(start_seed=args.start_seed, train_size=args.train_size, test_size=args.test_size)
    if args.max_action_length:
        creator.create_filtered_dataset(
            seed=args.start_seed,
            train_ratio=args.train_ratio,
            max_action_length=args.max_action_length,
            n_candidate=args.n_candidate,
            force_gen=args.force_gen
        )
    else:
        train_size = int(args.train_ratio * args.n_candidate)
        test_size = args.n_candidate - train_size
        creator.create_dataset(
            seed=args.start_seed,
            train_size=train_size,
            test_size=test_size,
            force_gen=args.force_gen)
+22 −11
Original line number Diff line number Diff line
@@ -2,40 +2,51 @@ set -x

export VLLM_ATTENTION_BACKEND=XFORMERS

python -m vagen.env.sokoban.create_dataset --data_dir data/sokoban-text
export PYTHONHASHSEED=0

python -m vagen.env.sokoban.create_dataset \
    --data_dir data/sokoban-text-1-step \
    --max_action_length 1 \
    --dim_room 6 6 \
    --num_boxes 1 \
    --max_steps 100 \
    --search_depth 30 \
    --start_seed 0 \
    --train_ratio 0.8 \
    --n_candidate 20000

# max_trajectory_length = max_prompt_length + max_response_length
#Set use_remove_padding to false, if true, causing batch size must be postive error in vllm

python3 -m vagen.trainer.main_ppo \
    algorithm.adv_estimator=grpo \
    data.train_files=data/sokoban-text/train.parquet \
    data.val_files=data/sokoban-text/test.parquet \
    data.train_batch_size=32 \
    data.train_files=data/sokoban-text-1-step/train.parquet \
    data.val_files=data/sokoban-text-1-step/test.parquet \
    data.train_batch_size=128 \
    data.max_prompt_length=1024 \
    data.max_response_length=128 \
    data.max_trajectory_length=2048 \
    data.max_trajectory_length=1024 \
    data.image_key=images \
    actor_rollout_ref.model.path=Qwen/Qwen2.5-0.5B-Instruct \
    actor_rollout_ref.actor.optim.lr=1e-6 \
    #Set use_remove_padding to false, if true, causing batch size must be postive error in vllm
    actor_rollout_ref.model.use_remove_padding=False \
    actor_rollout_ref.actor.ppo_mini_batch_size=32 \
    actor_rollout_ref.actor.ppo_micro_batch_size_per_gpu=1 \
    actor_rollout_ref.actor.ppo_micro_batch_size_per_gpu=8 \
    actor_rollout_ref.actor.use_kl_loss=True \
    actor_rollout_ref.actor.kl_loss_coef=0.001 \
    actor_rollout_ref.actor.kl_loss_type=low_var_kl \
    actor_rollout_ref.model.enable_gradient_checkpointing=True \
    actor_rollout_ref.actor.fsdp_config.param_offload=False \
    actor_rollout_ref.actor.fsdp_config.optimizer_offload=False \
    actor_rollout_ref.rollout.log_prob_micro_batch_size_per_gpu=1 \
    actor_rollout_ref.rollout.log_prob_micro_batch_size_per_gpu=8 \
    actor_rollout_ref.rollout.tensor_model_parallel_size=1 \
    actor_rollout_ref.rollout.name=vllm \
    actor_rollout_ref.rollout.gpu_memory_utilization=0.6 \
    actor_rollout_ref.rollout.gpu_memory_utilization=0.4 \
    actor_rollout_ref.rollout.enable_chunked_prefill=False \
    actor_rollout_ref.rollout.enforce_eager=False \
    actor_rollout_ref.rollout.free_cache_engine=False \
    actor_rollout_ref.rollout.n=1 \
    actor_rollout_ref.ref.log_prob_micro_batch_size_per_gpu=1 \
    actor_rollout_ref.ref.log_prob_micro_batch_size_per_gpu=8 \
    actor_rollout_ref.ref.fsdp_config.param_offload=True \
    algorithm.kl_ctrl.kl_coef=0.001 \
    trainer.critic_warmup=0 \
@@ -47,7 +58,7 @@ python3 -m vagen.trainer.main_ppo \
    trainer.save_freq=50 \
    trainer.test_freq=2 \
    trainer.total_epochs=15 \
    rollout_manger.max_turns=2 \
    rollout_manger.max_turns=1 \
    rollout_manger.window_size=5 \
    trainer.val_before_train=True \
    trainer.val_generations_to_log_to_wandb=5 \
+25 −12
Original line number Diff line number Diff line
@@ -2,39 +2,50 @@ set -x

export VLLM_ATTENTION_BACKEND=XFORMERS

python -m vagen.env.sokoban.create_dataset --data_dir data/sokoban-text
export PYTHONHASHSEED=0

python -m vagen.env.sokoban.create_dataset \
    --data_dir data/sokoban-text-1-step \
    --max_action_length 1 \
    --dim_room 6 6 \
    --num_boxes 1 \
    --max_steps 100 \
    --search_depth 30 \
    --start_seed 0 \
    --train_ratio 0.8 \
    --n_candidate 20000

# max_trajectory_length = max_prompt_length + max_response_length

python3 -m vagen.trainer.main_ppo \
    algorithm.adv_estimator=grpo \
    data.train_files=data/sokoban-text/train.parquet \
    data.val_files=data/sokoban-text/test.parquet \
    data.train_batch_size=32 \
    data.max_prompt_length=1024 \
    data.train_files=data/sokoban-text-1-step/train.parquet \
    data.val_files=data/sokoban-text-1-step/test.parquet \
    data.train_batch_size=256 \
    data.max_prompt_length=768 \
    data.max_response_length=128 \
    data.max_trajectory_length=2048 \
    data.max_trajectory_length=1024 \
    data.image_key=images \
    actor_rollout_ref.model.path=Qwen/Qwen2.5-0.5B-Instruct \
    actor_rollout_ref.actor.optim.lr=1e-6 \
    actor_rollout_ref.model.use_remove_padding=False \
    actor_rollout_ref.actor.ppo_mini_batch_size=32 \
    actor_rollout_ref.actor.ppo_micro_batch_size_per_gpu=2 \
    actor_rollout_ref.actor.ppo_mini_batch_size=64 \
    actor_rollout_ref.actor.ppo_micro_batch_size_per_gpu=8 \
    actor_rollout_ref.actor.use_kl_loss=True \
    actor_rollout_ref.actor.kl_loss_coef=0.001 \
    actor_rollout_ref.actor.kl_loss_type=low_var_kl \
    actor_rollout_ref.model.enable_gradient_checkpointing=True \
    actor_rollout_ref.actor.fsdp_config.param_offload=False \
    actor_rollout_ref.actor.fsdp_config.optimizer_offload=False \
    actor_rollout_ref.rollout.log_prob_micro_batch_size_per_gpu=2 \
    actor_rollout_ref.rollout.log_prob_micro_batch_size_per_gpu=8 \
    actor_rollout_ref.rollout.tensor_model_parallel_size=2 \
    actor_rollout_ref.rollout.name=vllm \
    actor_rollout_ref.rollout.gpu_memory_utilization=0.7 \
    actor_rollout_ref.rollout.gpu_memory_utilization=0.4 \
    actor_rollout_ref.rollout.enable_chunked_prefill=False \
    actor_rollout_ref.rollout.enforce_eager=False \
    actor_rollout_ref.rollout.free_cache_engine=False \
    actor_rollout_ref.rollout.n=1 \
    actor_rollout_ref.ref.log_prob_micro_batch_size_per_gpu=2 \
    actor_rollout_ref.ref.log_prob_micro_batch_size_per_gpu=8 \
    actor_rollout_ref.ref.fsdp_config.param_offload=True \
    algorithm.kl_ctrl.kl_coef=0.001 \
    trainer.critic_warmup=0 \
@@ -46,8 +57,10 @@ python3 -m vagen.trainer.main_ppo \
    trainer.save_freq=50 \
    trainer.test_freq=2 \
    trainer.total_epochs=15 \
    rollout_manger.max_turns=2 \
    rollout_manger.max_turns=1 \
    rollout_manger.window_size=5 \
    trainer.val_before_train=True \
    trainer.val_generations_to_log_to_wandb=5 \
    2>&1 | tee debug_qwen0_5_4_gpu_grpo.log

# NOTE change gpu_memory_utilization to a smaller value (0.4) to avoid oom error
 No newline at end of file
+19 −11
Original line number Diff line number Diff line
set -x

export VLLM_ATTENTION_BACKEND=XFORMERS
export PYTHONHASHSEED=0

python -m vagen.env.sokoban.create_dataset --visual_env --data_dir data/sokoban
python -m vagen.env.sokoban.create_dataset \
    --visual_env \
    --data_dir data/sokoban-vision-1-step \
    --max_action_length 1 \
    --dim_room 6 6 \
    --num_boxes 1 \
    --max_steps 100 \
    --search_depth 30 \
    --start_seed 0 \
    --train_ratio 0.8 \
    --n_candidate 20000

python3 -m vagen.trainer.main_ppo \
    algorithm.adv_estimator=grpo \
    data.train_files=data/sokoban/train.parquet \
    data.val_files=data/sokoban/test.parquet \
    data.train_batch_size=32 \
    data.max_prompt_length=1536 \
    data.train_files=data/sokoban-vision-1-step/train.parquet \
    data.val_files=data/sokoban-vision-1-step/test.parquet \
    data.train_batch_size=16 \
    data.max_prompt_length=768 \
    data.max_response_length=128 \
    data.max_trajectory_length=3072 \
    data.max_trajectory_length=1024 \
    data.image_key=images \
    actor_rollout_ref.model.path=Qwen/Qwen2.5-VL-3B-Instruct \
    actor_rollout_ref.actor.optim.lr=1e-6 \
    actor_rollout_ref.model.use_remove_padding=True \
    actor_rollout_ref.actor.ppo_mini_batch_size=32 \
    actor_rollout_ref.actor.ppo_mini_batch_size=4 \
    actor_rollout_ref.actor.ppo_micro_batch_size_per_gpu=1 \
    actor_rollout_ref.actor.use_kl_loss=True \
    actor_rollout_ref.actor.kl_loss_coef=0.001 \
@@ -25,7 +36,7 @@ python3 -m vagen.trainer.main_ppo \
    actor_rollout_ref.actor.fsdp_config.param_offload=False \
    actor_rollout_ref.actor.fsdp_config.optimizer_offload=False \
    actor_rollout_ref.rollout.log_prob_micro_batch_size_per_gpu=1 \
    actor_rollout_ref.rollout.tensor_model_parallel_size=2 \
    actor_rollout_ref.rollout.tensor_model_parallel_size=4 \
    actor_rollout_ref.rollout.name=vllm \
    actor_rollout_ref.rollout.gpu_memory_utilization=0.6 \
    actor_rollout_ref.rollout.enable_chunked_prefill=False \
@@ -49,6 +60,3 @@ python3 -m vagen.trainer.main_ppo \
    trainer.val_before_train=True \
    trainer.val_generations_to_log_to_wandb=5 \
    2>&1 | tee debug_qwen2_5_vl_4gpu_grpo.log
 No newline at end of file


# Tested and get oom error
Loading