Kaggle / Kaggle/kaggle-environments

Threading not working

Open
#104 1 comment 1 reaction 0 assignees View on GitHub
Dominant language
JavaScript
Stars
452
Forks
190
Avg merge
31m
Merged PRs (30d)
2

Description

Hi,

When using multiple threads to collect experience each thread always get lock on call of the step function.
I have been facing this problem since I updated the API. I tried both halite and football environment.

`new_obs, reward, done, info = self.env.step(actions)`

Full code:

`class EpisodeCollector(threading.Thread):
n_episode = 0
reward_sum = 0
max_episode = 0

def __init__(self, env: FootEnv, policy: Policy, result_queue=None, replays_dir=None):
super().__init__()
self.result_queue = result_queue
self.env = env
self.policy = policy
self.replays_dir = replays_dir
self.n_episode = -1

def clone(self):
obj = EpisodeCollector(self.env, self.policy)
obj.result_queue = self.result_queue
obj.replays_dir = self.replays_dir
obj.n_episode = self.n_episode
return obj

def run(self):
self.result_queue.put(self.collect(1))

def collect(self, n=1):
n = max(n, self.n_episode)
return [self.collect_() for _ in range(n)]

def collect_(self):
memory = Memory()
done = False
EpisodeCollector.n_episode += 1
obs = self.env.reset()
i = 0
total_reward = 0
state = None
while not done:
actions, state = self.policy.get_action(obs, state=state)
new_obs, reward, done, info = self.env.step(actions[0])
total_reward = reward
# store data
memory.store(obs, actions, reward, done)

if done or i % 100 == 0:
with lock:
print(
f"Episode: {EpisodeCollector.n_episode}/{EpisodeCollector.max_episode} | "
f"Step: {i} | "
f"Env ID: {self.env.env_id} | "
f"Reward: {total_reward} | "
f"Done: {done} | "
f"Total Rewards: {EpisodeCollector.reward_sum} | "
)
print(info)

obs = new_obs
i += 1
EpisodeCollector.reward_sum += total_reward
if self.replays_dir:
with open(os.path.join(self.replays_dir, f'replay-{uuid.uuid4().hex}.dill'), 'wb') as f:
dill.dump(memory, f)
return memory

class ParallelEpisodeCollector:

def __init__(self, env_fn, n_jobs, policy: Policy, replays_dir=None, ):
self.n_jobs = n_jobs
self.policy: Policy
self.envs = []
self.result_queue = Queue()
self.replays_dir = replays_dir
for i in range(n_jobs):
self.envs.append(env_fn(env_id=i))
self.collectors = [EpisodeCollector(env,
policy=policy,
result_queue=self.result_queue,
replays_dir=replays_dir) for env in self.envs]

def collect(self, n_steps=1):
if not n_steps: n_steps = 1
result_queue = self.result_queue
for i, collector in enumerate(self.collectors):
collector = collector.clone()
self.collectors[i] = collector
collector.n_episode = max(1, int(n_steps / len(self.collectors)))
print("Starting collector {}".format(i))
collector.start()
tmp = []
for _ in self.collectors:
res = result_queue.get()
tmp.extend(res)
[collector.join() for collector in self.collectors]
return tmp`

Contributor guide

Open the contributing guide

Research direction

Start by reproducing the reported concurrent EpisodeCollector behavior with the provided code in both the halite and football environments, focusing on the env.step(actions) call. Done means multiple collectors can progress without each thread blocking at that call; the issue mentions no repository file or test to run.

Written by the indexing model from the issue text.

Assessment

Tech stack
python
Domain
backend
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
28/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.