【发布时间】:2016-11-17 23:05:55
【问题描述】:
我是进程间通信的新手,我正在尝试了解os.pipe 和os.fork 在Python 中的用法。
在下面的代码中,如果我取消注释“Broken Pipe”错误会出现,否则它工作正常。
想法是在子进程退出时拥有一个 SIGCHLD 处理程序,并在仅子函数 (run_child) 和仅父函数 (sigchld_handler) 执行时增加各自的计数器。由于分叉的进程会有自己的内存版本,并且更改不会反映在父进程中,因此尝试是让子进程通过管道向父进程发送消息,并让父进程更新计数器。
import os
import signal
import time
class A(object):
def __init__(self):
self.parent = 0
self.child = 0
self._child_pid = None
self.rd , self.wr = os.pipe()
print self.rd , self.wr
signal.signal(signal.SIGCHLD, self.sigchld_handler)
def sigchld_handler(self, a, b):
self.parent += 1
print "Main run count : (parent) ", self.parent
#rf = os.fdopen(self.rd, 'r')
#self.child = int(rf.read())
#rf.close()
self._child_pid = None
def run_child(self):
self.child += 1
print "Main run count (child) : ", self.child
print "Running in child : " , os.getpid()
wr = os.fdopen(self.wr,'w')
text = "%s" % (self.child)
print "C==>", text
wr.write(text)
wr.close()
os._exit(os.EX_OK)
def run(self):
if self._child_pid:
print "Child Process", self._child_pid, " already running."
else:
self._child_pid = os.fork()
if not self._child_pid:
self.run_child()
a = A()
i = 0
while True:
a.run()
time.sleep(4)
i += 1
if i > 5:
break
有趣的是,错误出现在最初的几次迭代之后。有人可以解释一下为什么会出现错误,我应该怎么做才能解决这个问题。
编辑 1: 有几个类似的例子:ex1,ex2,ex3。实际上,我只是将它们用于学习,但就我而言,我将示例扩展为循环运行,以更像生产者/消费者队列。我知道这可能不是一个好方法,因为 Python 中提供了多进程/队列模块,但我想了解我在这里犯的错误。
编辑 2(解决方案):
基于@S.kozlov's answer,修改代码为每次通信创建一个新管道。这是修改后的代码。
import os
import pdb
import signal
import time
class A(object):
def __init__(self):
self.parent = 0
self.child = 0
self._child_pid = None
signal.signal(signal.SIGCHLD, self.sigchld_handler)
def sigchld_handler(self, a, b):
self.parent += 1
os.close(self.wr)
print "Main run count : (parent) ", self.parent
rd = os.fdopen(self.rd, 'r')
self.child = int(rd.read())
self._child_pid = None
def run_child(self):
self.child += 1
print "Main run count (child) : ", self.child
print "Running in child : " , os.getpid()
os.close(self.rd)
wr = os.fdopen(self.wr, 'w')
text = "%s" % (self.child)
print "C==>", text
wr.write(text)
wr.close()
os._exit(os.EX_OK)
def run(self):
if self._child_pid:
print "Child Process", self._child_pid, " already running."
else:
self.rd , self.wr = os.pipe()
self._child_pid = os.fork()
if not self._child_pid:
self.run_child()
a = A()
i = 0
while True:
a.run()
time.sleep(4)
i += 1
if i > 5:
break
有了这个,输出应该是这样的。
Main run count (child) : 1
Running in child : 15752
C==> 1
Main run count : (parent) 1
Main run count (child) : 2
Running in child : 15753
C==> 2
Main run count : (parent) 2
Main run count (child) : 3
Running in child : 15754
C==> 3
Main run count : (parent) 3
Main run count (child) : 4
Running in child : 15755
C==> 4
Main run count : (parent) 4
Main run count (child) : 5
Running in child : 15756
C==> 5
Main run count : (parent) 5
Main run count (child) : 6
Running in child : 15757
C==> 6
Main run count : (parent) 6
【问题讨论】:
-
它对我有用,没有错误。我已经运行了很多次了。
-
@quantummind :您的意思是在 sigchld_handler 中“取消注释”行之后说。
-
在关闭读取端后尝试写入管道时会发生管道损坏。我认为尝试阅读时不会发生这种情况。
-
谢谢@Barmar。这是一个有用的位。这是否意味着,我应该从 sigchld_handler 函数中删除 rf.close() 行。因为整个过程是循环的一部分。
-
在这两种情况下对我来说都没有错误。然而结果不同。使用 cmets 它会运行 while 循环,直到 i>5。取消注释时,最后输出行是“Main run count : (parent) 1”并继续运行。
标签: python pipe multiprocessing fork ipc