diff --git a/agent/async_agent.py b/agent/async_agent.py index de1eba8..c0fdac4 100644 --- a/agent/async_agent.py +++ b/agent/async_agent.py @@ -75,14 +75,19 @@ class AsyncAgent: target_network.share_memory() target_network.load_state_dict(learning_network.state_dict()) extra = target_network - elif config.worker == ContinuousAdvantageActorCritic or config.worker == ProximalPolicyOptimization: + elif config.worker == ContinuousAdvantageActorCritic \ + or config.worker == ProximalPolicyOptimization\ + or config.worker == DeterministicPolicyGradient: state_normalizer = StaticNormalizer(task.state_dim) reward_normalizer = StaticNormalizer(1) extra = [state_normalizer, reward_normalizer] + if config.worker == DeterministicPolicyGradient: + extra.append(config.replay_fn()) else: extra = None args = [(i, config, learning_network, extra) for i in range(config.num_workers)] args.append((config, task, learning_network, extra)) + # procs = [] procs = [mp.Process(target=evaluate, args=args[-1])] procs.extend([mp.Process(target=train, args=args[i]) for i in range(config.num_workers)]) for p in procs: p.start() diff --git a/async_worker/__init__.py b/async_worker/__init__.py index ee5e06e..9a2160d 100644 --- a/async_worker/__init__.py +++ b/async_worker/__init__.py @@ -3,4 +3,5 @@ from .continuous_actor_critic import * from .n_step_q import * from .one_step_sarsa import * from .one_step_q import * -from .ppo import * \ No newline at end of file +from .ppo import * +from .dpg import * \ No newline at end of file diff --git a/async_worker/dpg.py b/async_worker/dpg.py new file mode 100644 index 0000000..bb74296 --- /dev/null +++ b/async_worker/dpg.py @@ -0,0 +1,120 @@ +####################################################################### +# Copyright (C) 2017 Shangtong Zhang(zhangshangtong.cpp@gmail.com) # +# Permission given to modify the code as long as you keep this # +# declaration at the top # +####################################################################### + +import numpy as np +import torch.multiprocessing as mp +from network import * +from utils import * +from component import * +from async_worker import * +import pickle +import os +import time + +class DeterministicPolicyGradient: + def __init__(self, config, shared_network, extra): + self.config = config + self.task = config.task_fn() + + self.shared_network = shared_network + self.worker_network = config.network_fn() + self.worker_network.load_state_dict(self.shared_network.state_dict()) + self.target_network = config.network_fn() + self.target_network.load_state_dict(self.worker_network.state_dict()) + self.actor_opt = config.actor_optimizer_fn(self.shared_network.actor.parameters()) + self.critic_opt = config.critic_optimizer_fn(self.shared_network.critic.parameters()) + + self.random_process = config.random_process_fn() + self.criterion = nn.MSELoss() + + # self.state_normalizer = Normalizer(self.task.state_dim) + self.shared_state_normalizer = extra[0] + self.state_normalizer = StaticNormalizer(self.task.state_dim) + # self.replay = config.replay_fn() + self.replay = extra[-1] + + def soft_update(self, target, src): + for target_param, param in zip(target.parameters(), src.parameters()): + target_param.data.copy_(target_param.data * (1.0 - self.config.target_network_mix) + + param.data * self.config.target_network_mix) + + def episode(self, deterministic=False): + self.random_process.reset_states() + state = self.task.reset() + state = self.state_normalizer(state) + + config = self.config + actor = self.worker_network.actor + critic = self.worker_network.critic + target_actor = self.target_network.actor + target_critic = self.target_network.critic + + steps = 0 + total_reward = 0.0 + while True: + actor.eval() + action = actor.predict(np.stack([state])).flatten() + if not deterministic: + action += self.random_process.sample() + next_state, reward, done, info = self.task.step(action) + done = (done or (config.max_episode_length and steps >= config.max_episode_length)) + next_state = self.state_normalizer(next_state) + total_reward += reward + + if not deterministic: + self.replay.feed([state, action, reward, next_state, int(done)]) + with config.steps_lock: + config.total_steps.value += 1 + + steps += 1 + state = next_state + + if done: + break + + if not deterministic and self.replay.size() >= config.min_memory_size: + self.worker_network.train() + experiences = self.replay.sample() + states, actions, rewards, next_states, terminals = experiences + q_next = target_critic.predict(next_states, target_actor.predict(next_states)) + terminals = critic.to_torch_variable(terminals).unsqueeze(1) + rewards = critic.to_torch_variable(rewards).unsqueeze(1) + q_next = config.discount * q_next * (1 - terminals) + q_next.add_(rewards) + q_next = q_next.detach() + q = critic.predict(states, actions) + critic_loss = self.criterion(q, q_next) + + critic.zero_grad() + self.critic_opt.zero_grad() + critic_loss.backward() + with config.network_lock: + for param, worker_param in zip(self.shared_network.critic.parameters(), critic.parameters()): + if param.grad is not None: + break + param._grad = worker_param.grad + self.critic_opt.step() + + actions = actor.predict(states, False) + var_actions = Variable(actions.data, requires_grad=True) + q = critic.predict(states, var_actions) + q.backward(torch.ones(q.size())) + + actor.zero_grad() + self.actor_opt.zero_grad() + actions.backward(-var_actions.grad.data) + with config.network_lock: + for param, worker_param in zip(self.shared_network.actor.parameters(), actor.parameters()): + if param.grad is not None: + break + param._grad = worker_param.grad + self.actor_opt.step() + + self.worker_network.load_state_dict(self.shared_network.state_dict()) + + self.soft_update(self.target_network, self.worker_network) + + return steps, total_reward diff --git a/async_worker/ppo.py b/async_worker/ppo.py index 9224c93..b0992ee 100644 --- a/async_worker/ppo.py +++ b/async_worker/ppo.py @@ -72,7 +72,6 @@ class ProximalPolicyOptimization: values.append(value) state, reward, done, _ = self.task.step(action) state = self.state_normalizer(state) - # print state done = (done or (config.max_episode_length and episode_length > config.max_episode_length)) batched_rewards += reward diff --git a/component/replay.py b/component/replay.py index ff07329..48fb2d0 100644 --- a/component/replay.py +++ b/component/replay.py @@ -7,6 +7,7 @@ import numpy as np import torch import random +import torch.multiprocessing as mp class Replay: def __init__(self, memory_size, batch_size, dtype=np.float32): @@ -95,6 +96,62 @@ class HybridRewardReplay: self.next_states[sampled_indices], self.terminals[sampled_indices]] +class SharedReplay: + def __init__(self, memory_size, batch_size, state_shape, action_shape): + self.memory_size = memory_size + self.batch_size = batch_size + + self.states = torch.zeros((self.memory_size, ) + state_shape) + self.actions = torch.zeros((self.memory_size, ) + action_shape) + self.rewards = torch.zeros(self.memory_size) + self.next_states = torch.zeros((self.memory_size, ) + state_shape) + self.terminals = torch.zeros(self.memory_size) + + self.states.share_memory_() + self.actions.share_memory_() + self.rewards.share_memory_() + self.next_states.share_memory_() + self.terminals.share_memory_() + + self.pos = 0 + self.full = False + self.buffer_lock = mp.Lock() + + def feed_(self, experience): + state, action, reward, next_state, done = experience + self.states[self.pos][:] = torch.FloatTensor(state) + self.actions[self.pos][:] = torch.FloatTensor(action) + self.rewards[self.pos] = reward + self.next_states[self.pos][:] = torch.FloatTensor(next_state) + self.terminals[self.pos] = done + + self.pos += 1 + if self.pos == self.memory_size: + self.full = True + self.pos = 0 + + def size(self): + if self.full: + return self.memory_size + return self.pos + + def sample_(self): + upper_bound = self.memory_size if self.full else self.pos + sampled_indices = torch.LongTensor(np.random.randint(0, upper_bound, size=self.batch_size)) + return [self.states[sampled_indices], + self.actions[sampled_indices], + self.rewards[sampled_indices], + self.next_states[sampled_indices], + self.terminals[sampled_indices]] + + def feed(self, experience): + with self.buffer_lock: + self.feed_(experience) + + def sample(self): + with self.buffer_lock: + return self.sample_() + class HighDimActionReplay: def __init__(self, memory_size, batch_size, dtype=np.float32): self.memory_size = memory_size @@ -130,6 +187,11 @@ class HighDimActionReplay: self.full = True self.pos = 0 + def size(self): + if self.full: + return self.memory_size + return self.pos + def sample(self): upper_bound = self.memory_size if self.full else self.pos sampled_indices = np.random.randint(0, upper_bound, size=self.batch_size) diff --git a/main.py b/main.py index 87c6ce6..06aa69a 100644 --- a/main.py +++ b/main.py @@ -201,9 +201,9 @@ def a3c_continuous(): def p3o_continuous(): config = Config() - config.task_fn = lambda: Pendulum() + # config.task_fn = lambda: Pendulum() # config.task_fn = lambda: BipedalWalkerHardcore() - # config.task_fn = lambda: Roboschool('RoboschoolInvertedPendulum-v1') + config.task_fn = lambda: Roboschool('RoboschoolInvertedPendulum-v1') # config.task_fn = lambda: Roboschool('RoboschoolAnt-v1') task = config.task_fn() config.actor_network_fn = lambda: GaussianActorNet(task.state_dim, task.action_dim, @@ -214,7 +214,8 @@ def p3o_continuous(): config.critic_optimizer_fn = lambda params: torch.optim.Adam(params, 0.001) config.policy_fn = lambda: GaussianPolicy() - config.replay_fn = lambda: GeneralReplay(memory_size=2048, batch_size=2048) + # config.replay_fn = lambda: GeneralReplay(memory_size=2048, batch_size=2048) + config.replay_fn = lambda: GeneralReplay(memory_size=2048, batch_size=64) config.worker = ProximalPolicyOptimization config.discount = 0.99 config.gae_tau = 0.97 @@ -225,7 +226,7 @@ def p3o_continuous(): config.entropy_weight = 0 config.gradient_clip = 20 config.rollout_length = 10000 - config.optimize_epochs = 1 + config.optimize_epochs = 10 config.ppo_ratio_clip = 0.2 config.logger = Logger('./log', gym.logger) agent = AsyncAgent(config) @@ -261,16 +262,56 @@ def ddpg_continuous(): config.logger = Logger('./log', gym.logger) run_episodes(DDPGAgent(config)) +def addpg_continuous(): + config = Config() + # config.task_fn = lambda: Pendulum() + # config.task_fn = lambda: ContinuousLunarLander() + config.task_fn = lambda: Roboschool('RoboschoolInvertedPendulum-v1') + # config.task_fn = lambda: Roboschool('RoboschoolReacher-v1') + # config.task_fn = lambda: BipedalWalker() + task = config.task_fn() + config.actor_network_fn = lambda: DeterministicActorNet( + task.state_dim, task.action_dim, F.tanh, 2, non_linear=F.relu, batch_norm=False) + config.critic_network_fn = lambda: DeterministicCriticNet( + task.state_dim, task.action_dim, non_linear=F.relu, batch_norm=False) + config.network_fn = lambda: DisjointActorCriticNet(config.actor_network_fn, config.critic_network_fn) + config.actor_optimizer_fn = lambda params: torch.optim.Adam(params, lr=1e-4) + config.critic_optimizer_fn =\ + lambda params: torch.optim.Adam(params, lr=1e-4) + # config.replay_fn = lambda: HighDimActionReplay(memory_size=1000000, batch_size=64) + config.replay_fn = lambda: SharedReplay(memory_size=1000000, batch_size=64, + state_shape=(task.state_dim, ), action_shape=(task.action_dim, )) + # config.replay_fn = lambda: GeneralReplay(memory_size=256, batch_size=64) + config.discount = 0.99 + config.max_episode_length = task.max_episode_steps + config.random_process_fn = \ + lambda: OrnsteinUhlenbeckProcess(size=task.action_dim, theta=0.15, sigma=0.2, + n_steps_annealing=100000) + config.worker = DeterministicPolicyGradient + config.num_workers = 6 + config.min_memory_size = 50 + config.target_network_mix = 0.001 + # config.update_interval = 10 + config.test_interval = 500 + config.test_repetitions = 1 + config.gradient_clip = 20 + config.rollout_length = 16 + config.optimize_epochs = 1 + config.logger = Logger('./log', gym.logger) + agent = AsyncAgent(config) + agent.run() + if __name__ == '__main__': - # gym.logger.setLevel(logging.DEBUG) - gym.logger.setLevel(logging.INFO) + gym.logger.setLevel(logging.DEBUG) + # gym.logger.setLevel(logging.INFO) # dqn_cart_pole() # async_cart_pole() # a3c_cart_pole() # a3c_continuous() # p3o_continuous() - ddpg_continuous() + # ddpg_continuous() + addpg_continuous() # dqn_fruit() # hrdqn_fruit()