You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Celery任务终止时子进程未停止的问题排查请求

问题诊断与修复方案

核心问题分析

当前代码存在多个导致子进程无法被终止的问题:

  1. os.kill调用错误:os.getpid(subprocess_id)是完全错误的用法,os.getpid()无参数,返回当前进程PID,正确方式应直接传入子进程PID。且SIGUSR1是自定义信号,Scrapy默认不会响应此信号终止进程。
  2. 任务阻塞时无法响应终止:process.wait()会阻塞Celery任务进程,此时self.is_aborted()无法被检查,任务被标记终止后也无法执行子进程清理逻辑。
  3. Docker进程隔离限制:若stop接口与Celery Worker不在同一容器,直接使用容器内PID调用os.kill会失败——PID仅在各自容器的进程命名空间内有效。
  4. revoke信号过于强硬:使用SIGKILL直接杀死Celery任务进程,导致任务无机会执行子进程清理代码,子进程会沦为孤儿进程继续运行。

分步修复

1. 修正os.kill调用逻辑

修改stop接口中的进程终止代码:

if subprocess_id:
    try:
        # 先尝试优雅终止
        os.kill(subprocess_id, signal.SIGTERM)
        # 等待1秒确认进程退出,否则强制杀死
        time.sleep(1)
        if os.path.exists(f"/proc/{subprocess_id}"):
            os.kill(subprocess_id, signal.SIGKILL)
        logger.info(f"Subprocess with ID {subprocess_id} terminated.")
    except OSError as e:
        logger.error(f"Failed to terminate subprocess {subprocess_id}: {str(e)}")

2. 改造Celery任务,主动响应终止信号

在任务中定期检查终止状态,避免process.wait()阻塞导致无法处理终止请求:

@celery_app.task(bind=True, name='jina_fetch_urls_celery_task', base=AbortableTask)
def jina_fetch_urls_celery_task(self, run_id: PydanticObjectId, site: str):
    process = None
    try:
        for i in range(5):
            if self.is_aborted():
                return 'Task stopped!'
            logger.info(f"Progress: {i}")
            self.update_state(state='PROGRESS', meta={'progress': i})
            sleep(1)
        
        process = subprocess.Popen(
            ['scrapy', 'crawl', 'base_spider', '-a', f'start_url={site}', '-a', f'task_id={str(run_id)}'],
            cwd='/code/jina_scrapy'
        )
        
        subprocess_id = process.pid
        self.update_state(state='PROGRESS', meta={'subprocess_id': subprocess_id})
        logger.info(f"Subprocess ID set in Celery: {subprocess_id}")
        
        # 循环检查任务状态与子进程状态
        while process.poll() is None:
            if self.is_aborted():
                logger.info(f"Task aborted, terminating subprocess {subprocess_id}")
                process.terminate()
                # 超时强制杀死
                process.wait(timeout=3)
                return {'status': 'ABORTED', 'subprocess_id': subprocess_id}
            sleep(0.5)
        
        return {'status': 'DONE', 'subprocess_id': subprocess_id}
    except Exception as e:
        logger.error(f"Task failed: {str(e)}")
        # 异常时清理子进程
        if process and process.poll() is None:
            process.kill()
        return {'status': 'FAILURE', 'error': str(e)}

3. 调整revoke策略,优先优雅终止

在stop接口中,先使用SIGTERM让任务有机会清理子进程,再强制杀死:

# 先发送SIGTERM,允许任务自行清理
celery_task.revoke(terminate=True, signal='SIGTERM')
# 等待2秒让任务处理终止逻辑
time.sleep(2)
# 若任务仍未终止,再强制杀死
if celery_task.state not in ['REVOKED', 'SUCCESS', 'FAILURE']:
    celery_task.revoke(terminate=True, signal='SIGKILL')

4. 跨容器场景的额外处理(若适用)

若stop接口与Celery Worker不在同一容器,需通过以下方式解决进程隔离问题:

  • 使用Docker API调用docker exec进入Worker容器执行kill命令;
  • 将子进程信息存储到共享存储(如Redis),由Worker提供内部清理接口,stop接口调用该接口终止子进程;
  • 调整部署架构,将stop接口与Worker部署在同一容器(仅适用于单容器场景)。

内容的提问来源于stack exchange,提问作者DarkHorse_906

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 04:45:22