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

目录新增文件时无法向RabbitMQ发送消息排查求助

问题排查:SFTP挂载目录新增文件时Watchdog未触发RabbitMQ消息发送

我在Kubernetes Pod中运行以下Python代码,Pod可正常连接RabbitMQ,测试消息channel.basic_publish(exchange='', routing_key='myqueue', body='Test message')能成功发送。需求是当/mnt/目录(挂载了SFTP服务的Persistent Volume)新增文件时,向RabbitMQ队列发送消息,但实际新增文件时无消息发送。代码如下:

import pika
import time
import os
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler

class FileEventHandler(FileSystemEventHandler):
    def __init__(self, channel):
        self.channel = channel

    def on_created(self, event):
        if not event.is_directory:
            message = 'New file created: %s' % event.src_path
            print("Detected new file, sending message: ", message)
            is_delivered = self.channel.basic_publish(exchange='', routing_key='myqueue', body=message)
            print("Message delivery status: ", is_delivered)

def main():
    credentials = pika.PlainCredentials('LBO64L', 'test')
    print('Logged')

    parameters = pika.ConnectionParameters('rabbitmq.developpement-dev-01.svc.cluster.local',
                                   5672,
                                   '/',
                                   credentials,
                                   socket_timeout=2)
    print(parameters)

    connection = pika.BlockingConnection(parameters)
    print("Connection...")
   
    channel = connection.channel()

    channel.queue_declare(queue='myqueue')

    print('Queue created !')

    channel.basic_publish(exchange='', routing_key='myqueue', body='Test message')
    print("Test message sent")

    event_handler = FileEventHandler(channel)

    if os.path.isdir('/mnt/lb064l/data/'):
        print("Directory exists and is accessible")
    else:
        print("Directory does not exist or is not accessible")

    print('Starting observer...')
    observer = Observer()
    observer.schedule(event_handler, path='/mnt/lb064l/data/', recursive=True)
    observer.start()

    try:
        print('Running...')
        while True:
            time.sleep(1)
    except KeyboardInterrupt:
        observer.stop()
    observer.join()

    print('Message sent to queue !')

if __name__ == "__main__":
    main()

可能的问题原因

  • Watchdog不支持远程文件系统事件监听:Watchdog依赖操作系统原生的文件系统事件通知(比如Linux的inotify),但SFTP挂载的远程目录属于网络文件系统,服务器端的文件变化不会主动触发本地Pod的文件系统事件,自然无法触发on_created回调。
  • 文件上传流程触发的是移动事件而非创建事件:很多SFTP客户端上传文件时,会先传一个临时文件(比如带.tmp后缀),上传完成后再重命名为目标文件名。这种情况下,Watchdog会捕捉到临时文件的on_created和on_deleted,但最终的目标文件是通过on_moved事件产生的,你的代码只监听了创建事件,所以会漏掉。
  • 权限不足导致事件无法被感知:虽然代码检测到目录存在,但Pod的运行用户可能没有读取目录内文件变化的权限。检查Pod的SecurityContext配置,确保运行用户对/mnt/lb064l/data/目录有足够的权限(至少读权限)。
  • 递归监听失效:部分网络文件系统不支持递归的事件监听,尝试把recursive=True改成False,只监听一级目录,测试是否能触发事件。
  • Pika连接已断开:初始连接正常,但长时间运行后,Pika的BlockingConnection可能因为网络波动或心跳超时断开,导致后续basic_publish失败。可以在on_created方法里先检查连接状态,或者添加自动重连逻辑。
  • 事件未触发(无日志输出):先查看Pod的日志,看是否有"Detected new file"的打印。如果没有,说明Watchdog根本没捕捉到事件,大概率是文件系统不支持的问题;如果有打印但消息没发出去,再排查RabbitMQ连接或发布逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 11:27:39