Skip to content

slurm worker --without-heartbeat causing issues? #519

Description

@pearcetm

python -m girder_worker -l info -Ofair --prefetch-multiplier=1 --without-heartbeat --concurrency=2

@willdunklin When creating slurm jobs on our HPC, I'm seeing job logs that contain broken pipe errors. The worker is started without a --pool argument so it defaults to prefork; we are using concurrency=8. It looks like what is happening is that once per process, the first job that gets picked up contains an error like the following:

BrokenPipeError: [Errno 32] Broken pipe
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/celery/app/trace.py", line 585, in trace_task
    R = retval = fun(*args, **kwargs)
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/slicer_cli_web/slicer_cli_web/singularity/slicer_cli_web_singularity/girder_worker_plugin/direct_singularity_run.py", line 27, in __call__
    super().__call__(*args, **kwargs)
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/girder_worker/girder_worker/docker/tasks/__init__.py", line 383, in __call__
    super().__call__(*args, **kwargs)
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/girder_worker/girder_worker/task.py", line 170, in __call__
    results = super().__call__(*_t_args, **_t_kwargs)
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/celery/app/trace.py", line 858, in __protected_call__
    return self.run(*args, **kwargs)
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/slicer_cli_web/slicer_cli_web/singularity/slicer_cli_web_singularity/girder_worker_plugin/direct_singularity_run.py", line 59, in run
    return singularity_run(task, **kwargs)
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/girder_worker/girder_worker/singularity/girder_worker_singularity/tasks/__init__.py", line 59, in singularity_run
    slurm_dispatch(task, container_args, run_kwargs, read_streams, write_streams, log_file_name)
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/girder_worker/girder_worker/slurm/girder_worker_slurm/__init__.py", line 36, in slurm_dispatch
    utils.select_loop(exit_condition=check_job_cancellation,
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/girder_worker/girder_worker/docker/utils.py", line 32, in select_loop
    exit = exit_condition()
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/girder_worker/girder_worker/slurm/girder_worker_slurm/__init__.py", line 27, in check_job_cancellation
    if task.canceled and monitor_thread.job_id:
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/girder_worker/girder_worker/task.py", line 121, in canceled
    return is_revoked(self)
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/girder_worker/girder_worker/utils.py", line 114, in is_revoked
    return task.request.id in _revoked_tasks(task)
  File "/home/wsigenie/milk/dsa_custom_config/worker/lib/girder_worker/girder_worker/utils.py", line 57, in _revoked_tasks
    _revoked = _worker_inspector(task).revoked()
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/celery/app/control.py", line 254, in revoked
    return self._request('revoked')
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/celery/app/control.py", line 106, in _request
    return self._prepare(self.app.control.broadcast(
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/celery/app/control.py", line 785, in broadcast
    return self.mailbox(conn)._broadcast(
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/kombu/pidbox.py", line 347, in _broadcast
    self._publish(command, arguments, destination=destination,
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/kombu/pidbox.py", line 309, in _publish
    maybe_declare(self.reply_queue(chan))
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/kombu/common.py", line 113, in maybe_declare
    return _maybe_declare(entity, channel)
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/kombu/common.py", line 155, in _maybe_declare
    entity.declare(channel=channel)
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/kombu/entity.py", line 616, in declare
    self._create_exchange(nowait=nowait, channel=channel)
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/kombu/entity.py", line 623, in _create_exchange
    self.exchange.declare(nowait=nowait, channel=channel)
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/kombu/entity.py", line 184, in declare
    return (channel or self.channel).exchange_declare(
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/amqp/channel.py", line 624, in exchange_declare
    self.send_method(
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/amqp/abstract_channel.py", line 70, in send_method
    conn.frame_writer(1, self.channel_id, sig, args, content)
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/amqp/method_framing.py", line 186, in write_frame
    write(buffer_store.view[:offset])
  File "/home/wsigenie/milk/dsa_custom_config/worker/venv/lib64/python3.9/site-packages/amqp/transport.py", line 350, in write
    self._write(s)

Then, the process will pick up a second job, while the first one is still running via slurm. The second job's log will pick up log messages from the first job.

It looks like stale connections to rabbitmq are a likely culprit; would --without-heartbeat impact that at all? Can we just add the heartbeat to help avoid this? Also, if I'm not mistaken, the downstream effect of the error is because monitor threads stay alive and continue to write messages even after the main thread exits (freeing the process to pick up another job). Any suggestions for catching such errors more gracefully, and/or killing the ENTIRE job if the main thread goes down?

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions