1 远程控制树莓派使用进程运行Python文件,如何实现?

1.1 subprocess

subprocess 模块主要用于创建子进程,并连接它们的输入、输出和错误管道,获取它们的返回状态。通俗地说就是通过这个模块,你可以在 Python 的代码里执行操作系统级别的命令,比如ipconfig、du -sh等。 它替代了一些老的模块和函数,比如:os.system、os.spawn*等。 subprocess 过去版本中的call(),check_call()和check_output()已经被3.5版本中新增的run()方法取代了。

大多数情况下,使用run()方法调用子进程,执行操作系统命令。在更高级的使用场景,你还可以使用 Popen 接口,其实run()方法在底层调用的就是 Popen 接口。

1.1 subprocess.run

要实现和 os.system()命令相同的方式,运行外部命令而不与之交互时候,我们可以使用 run()函数。

先看一下其语法结构。

subprocess.run(args, *, stdin=None, input=None, stdout=None, stderr=None, shell=False, timeout=None, check=False, encoding=None, errors=None)

功能:执行 args 参数所表示的命令,等待命令结束,并返回一个 CompletedProcess 类型对象,而不是我们想要的执行结果或相关信息

一般来说用到的常用参数:

args:表示要执行的命令,字符串或字符串参数列表。

  • stdin、stdout 和 stderr:子进程的标准输入、输出和错误。其值可以是subprocess.PIPE、subprocess.DEVNULL、一个已经存在的文件描述符、已经打开的文件对象或者 None。

  • subprocess.PIPE表示为子进程创建新的管道,subprocess.DEVNULL表示使用os.devnull,对于不应该显示或捕获输出的情况,使用 DEVNULL 来抑制输出流 。默认使用的是 None,表示什么都不做。另外,stderr=subprocess.STDOUT是特殊值,可传递给 stderr 参数,表示 stdout 和 stderr 合并输出。

结合常规输出和错误输出

要将进程的错误输出定向到其标准输出通道,可以使用 STDOUT 代替 stderr 而不是 PIPE。

    import subprocess
    
    print('popen4:')
    proc = subprocess.Popen(
        'ls -l; echo "to stderr" 1>&2',
        shell=True,
        stdin=subprocess.PIPE,
        stdout=subprocess.PIPE,
        stderr=subprocess.STDOUT,
    )
    msg = 'through stdin to stdout\n'.encode('utf-8')
    stdout_value, stderr_value = proc.communicate(msg)
    print('combined output:', repr(stdout_value.decode('utf-8')))
    print('stderr value   :', repr(stderr_value))

timeout:设置命令超时时间。如果命令执行时间超时,子进程将被杀死,并弹出TimeoutExpired异常。

    import subprocess
    
    try:
        subprocess.run(['false'], check=True)
    except subprocess.CalledProcessError as err:
        print('ERROR:', err)

check:如果该参数设置为 True,并且进程退出状态码不是 0,则弹出CalledProcessError异常。

encoding:如果指定了该参数,则 stdin、stdout 和 stderr 可以接收字符串数据,并以该编码方式编码。否则只接收 bytes 类型的数据,需要用decode对stdout进行解码。

shell:如果该参数为 True,将通过操作系统的 shell 执行指定的命令,默认为False。在 Linux 环境下,当 args 是个字符串时,必须指定 shell=True。当 args 是个列表的时候,shell 保持默认的 False。

返回

run方法会返回一个CompletedProcess对象,包括:

  • args 用户输入的启动进程的参数,字符串列表或字符串。
  • returncode 进程结束状态返回码。0表示成功状态。
  • stdout 获取子进程的标准输出。通常为 bytes,None 表示没有捕获值。如果你在调用 run() 方法时,设置了参数stderr=subprocess.STDOUT,则错误信息会和 stdout 一起输出。
  • stderr 获取子进程的错误信息。通常为 bytes 类型序列,None 表示没有捕获值。
  • check_returncode() 用于检查返回码。如果返回状态码不为零,弹出CalledProcessError异常。

1.3 交互式输入

并不是所有的操作系统命令都像dir或者ipconfig那样单纯地返回执行结果,还有很多像python这种交互式的命令,你要输入点什么,然后它返回执行的结果。使用run()方法怎么向stdin里输入?

run()方法的stdin参数可以接收一个文件句柄。比如在一个text.txt文件中写入print(‘hello Python’)。然后使用:

    import subprocess
    
    fd = open("d:\\1.txt")
    ret = subprocess.run("python", stdin=fd, stdout=subprocess.PIPE,shell=True)
    print(ret.stdout)
    fd.close()

这样做,虽然可以达到目的,但是很不方便,也不是以代码驱动的方式。这个时候,我们可以使用Popen类。

1.4 Popen

返回值是一个Popen对象,而不是CompletedProcess对象。

上述python的交互式命令功能:

    import subprocess
    
    s = subprocess.Popen("python", stdout=subprocess.PIPE, stdin=subprocess.PIPE, shell=True)
    s.stdin.write(b"import os\n")
    s.stdin.write(b"print(os.environ)")
    s.stdin.close()
    
    out = s.stdout.read().decode("GBK")
    s.stdout.close()
    print(out)
1.4.1 与进程的双向通信

要同时设置 Popen 实例进行读写。

    import subprocess
    
    print('popen2:')
    
    proc = subprocess.Popen(
        ['ls', '-l'],
        stdin=subprocess.PIPE,
        stdout=subprocess.PIPE,
    )
    msg = 'through stdin to stdout'.encode('utf-8')
    stdout_value = proc.communicate(msg)[0].decode('utf-8')
    print('pass through:', repr(stdout_value))

使用 communicate() 而非 .stdin.write, .stdout.read 或者 .stderr.read 来避免由于任意其他 OS 管道缓冲区被子进程填满阻塞而导致的死锁。

1.4.2 管道之间的连接

通过创建单独的 Popen 实例并将它们的输入和输出链接在一起,可以类似于 Unix shell 的工作方式将多个命令连接到管道中。

    import subprocess
    
    cat = subprocess.Popen(
        ['cat', 'subprocess_demo.py'],
        stdout=subprocess.PIPE,  # 提供输出的方式
    )
    
    grep = subprocess.Popen(
        ['grep', '公众号'],
        stdin=cat.stdout,  # cat 的输出最为输入
        stdout=subprocess.PIPE,
    )
    
    cut = subprocess.Popen(
        ['awk', '-F', ':', '{print $2}'],
        stdin=grep.stdout,
        stdout=subprocess.PIPE,
    )
    
    end_of_pipe = cut.stdout
    
    print(end_of_pipe.readline().decode('utf-8'))

输出

    python 学习开发

上面的内容就等价于下面的命令

    cat subprocess_demo.py |grep "公众号" |awk -F ':' '{print $2}'

2 使用subprocess的时候,怎么获取stdout和stderr

2.1 区分p.stdout下三个读函数

  • read(): 全部读,返回字符串,阻塞
  • readline():按行读,异步,会有多余空行
  • readlines():按行读完,返回列表,阻塞
    import subprocess
    p = subprocess.Popen(['tail','-10','/tmp/hosts.txt'],stdin=subprocess.PIPE,stdout=subprocess.PIPE,stderr=subprocess.PIPE,shell=False)
    
    stdout,stderr = p.communicate()
    print 'stdout : ',stdout
    print 'stderr : ',stder
  • popen调用的时候会在父进程和子进程建立管道,然后我们可以把子进程的标准输出和错误输出都重定向到管道,然后从父进程取出。上面的communicate会一直阻塞,直到子进程跑完。这种方式是不能及时获取子程序的stdout和stderr。
  • 只有当子进程结束才会打印结果,如果想要异步读,可以用p.stdout.readline()。
  • readline()每隔一行会多读出一个空行,可以用.splitlines()去掉空行,原理通split()函数。

2.2 实时获取输出/异常结果

使用poll轮询p的状态,如果没结束就读取一行。这里把错误输出重定向到PIPE对应的标准输出,也就是说现在stderr都stdout是一起的了,下面是一个while,poll回去不断查看子进程是否已经被终止,如果程序没有终止,就一直返回None,但是子进程终止了就返回状态码,甚至于调用多次poll都会返回状态码。上面的demo就是可以获取子进程的标准输出和标准错误输出。

    p = subprocess.Popen("/etc/service/tops-cmos/module/hadoop/test.sh", shell=True, stdout=subprocess.PIPE, stderr=subprocess.STDOUT)
    returncode = p.poll()
    while returncode is None:
            line = p.stdout.readline()
            returncode = p.poll()
            line = line.strip()
            print line
    print returncode

也可以这样,但是不够上面那种严谨

    msg = p.stdout.readline()
    while msg:
    	msg=p.stdout.readline()
        print(msg)

2.3 子线程运行完毕后怎么通知主线程呢

我们可以使用全局变量实现,在全局设置一个flag=False,子线程中运行完毕后将其设置为True,与此同时主线程中一致轮询flag,当为TRUE的时候执行下一步操作。

这里有使用队列的例子作为参考:

    import threading
    import queue
    import time
    import random
    
    '''
    需求:主线程开启了多个线程去干活,每个线程需要完成的时间
    不同,但是在干完活以后都要通知给主线程
    多线程和queue配合使用,实现子线程和主线程相互通信的例子
    '''
    q = queue.Queue()
    threads=[]
    class MyThread(threading.Thread):
        def __init__(self,q,t,j):
            super(MyThread,self).__init__()
            self.q=q
            self.t=t
            self.j=j
    
        def run(self):
            time.sleep(self.j)
            # 通过q.put()方法,将每个子线程要返回给主线程的消息,存到队列中
            self.q.put("我是第%d个线程,我睡眠了%d秒,当前时间是%s" % (self.t, self.j,time.ctime()))
    
    '''
    # 生成15个子线程,加入到线程组里,
        # 每个线程随机睡眠1-8秒(模拟每个线程干活时间的长短不同)
    '''
    
    for i in range(15):
       j=random.randint(1,8)
       threads.append(MyThread(q,i,j))
    
    #    循环开启所有子线程
    for mt in threads:
        mt.start()
    print('进程开启时间:%s'%(time.ctime()))
    '''
    通过一个while循环,当q队列中不为空时,通过q.get()方法,
    循环读取队列q中的消息,每次计数器加一,当计数器到15时,
    证明所有子线程的消息都已经拿到了,此时循环停止
    '''
    count = 0
    while True:
        if not q.empty():
            print(q.get())
            count+=1
        if count==15:
            break

2.4 如何在主线程杀死子线程呢?

  • 方法1 使用退出标记,这里使用threading.Event()创建一个事件管理标记flag,这种方法是比较安全的。
    # encoding:utf-8
    import time
    import threading
    
    
    class StoppableThread(threading.Thread):
        """Thread class with a stop() method. The thread itself has to check
        regularly for the stopped() condition."""
    
        def __init__(self,  *args, **kwargs):
            super(StoppableThread, self).__init__(*args, **kwargs)
            self._stop_event = threading.Event()
    
        def stop(self):
            self._stop_event.set()
    
        def stopped(self):
            return self._stop_event.is_set()
    
        def run(self):
            print("begin run the child thread")
            while True:
                print("sleep 1s")
                time.sleep(1)
                if self.stopped():
                    # 做一些必要的收尾工作
                    break
    
    
    if __name__ == "__main__":
        print("begin run main thread")
        t = StoppableThread()
        t.start()
        time.sleep(3)
        # t.join()
        t.stop()
        print("main thread end")

主进程在3s后关闭子线程。

  • 方法2 使用抛出异常,导致线程退出
    import threading
    import time
    import inspect
    import ctypes
    
    
    def _async_raise(tid, exctype):
        """Raises an exception in the threads with id tid"""
        if not inspect.isclass(exctype):
            raise TypeError("Only types can be raised (not instances)")
        res = ctypes.pythonapi.PyThreadState_SetAsyncExc(ctypes.c_long(tid), ctypes.py_object(exctype))
        if res == 0:
            raise ValueError("invalid thread id")
        elif res != 1:
            # """if it returns a number greater than one, you're in trouble,
            # and you should call it again with exc=NULL to revert the effect"""
            ctypes.pythonapi.PyThreadState_SetAsyncExc(tid, None)
            raise SystemError("PyThreadState_SetAsyncExc failed")
    
    
    def stop_thread(thread):
        _async_raise(thread.ident, SystemExit)
    
    
    class TestThread(threading.Thread):
        def run(self):
            print("begin run the child thread")
            while True:
                print("sleep 1s")
                time.sleep(1)
    
    
    if __name__ == "__main__":
        print("begin run main thread")
        t = TestThread()
        t.start()
        time.sleep(3)
        stop_thread(t)
        print("main thread end")

运行结果都是

    begin run main thread
    begin run the child thread
    sleep 1s
    sleep 1s
    sleep 1s
    main thread end

这种方法是强制杀死线程,但是如果线程中涉及获取释放锁,可能会导致死锁。

结合我的上一篇博客,使用Tornado和subprocess实现和外部程序通信,可以实现在浏览器客户端远程控制树莓派,上传并执行python代码,并且可以通过websocket来建立连接,控制进程的启动和中断,并及时反馈输出信息。

部分代码:

    # 存放线程对象
    p = None
    t = None
    class StoppableThread(threading.Thread):
        """Thread class with a stop() method. The thread itself has to check
        regularly for the stopped() condition."""
    
        def __init__(self,  *args, **kwargs):
            super(StoppableThread, self).__init__(*args, **kwargs)
            self._stop_event = threading.Event()
    
        def stop(self):
            print("stop!!!")
            self._stop_event.set()
    
        def stopped(self):
            return self._stop_event.is_set()
    
        def run(self):
            print("begin run the child thread")
            global p
            head = "运行结果:<br>"
            msg = p.stdout.readline()
            while msg:
                if msg:
                    msg=p.stdout.readline()
                    print(msg)
                if self.stopped():
                    print("此线程结束",self)
                    break
    
    ##
    # 执行代码文件的类
    class SsHandler(BaseHandler):
        @asynchronous
        def get(self, cmd):
            '''
            cmd: 接收前端的正则表达式字符串
            '''
            print("Command: %s" % cmd)
            global p
            if cmd == 'start':
                # 将文件名修改
                p = subprocess.Popen(["python","./tmp/new_code.py"], shell=True, stdout=subprocess.PIPE, stderr = subprocess.STDOUT, bufsize=1)
                global t
                t = StoppableThread()
                t.start()
                # p.stdout.close()
                # p.wait()
            elif cmd == 'stop':
                t.stop()
                if p is not None:
                    poll = p.poll() #获取子进程的状态
    
                    if poll==0:
                        # print(p.stdout.read())
                        # print(p.stderr.read())
                        print("程序执行完毕!不需要手动结束!")
                    elif poll is None:
                        print("程序正在执行!马上退出...")
                        p.terminate()
                        # 获取子进程的pid
                        pid = p.pid
                        print("pid: {}".format(pid))
                        # 杀死进程
                        p.kill()
                        print("killed.")
                    else:
                        print("状态码:{}".format(poll))
                        print("程序异常退出!")
                else:
                    print("p is nothing")
            else:
                print("Command is not right!")
            
            print("Done!")
        
        # 没用
        def post(self):
            self.write("post-StartHandler")
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐