1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

17

18

19

20

21

22

23

24

25

26

27

28

29

30

31

32

33

34

35

36

37

38

39

40

41

42

43

44

45

46

47

48

49

50

51

52

53

54

55

56

57

58

59

60

61

62

63

64

65

66

67

68

69

70

71

72

73

74

75

76

77

78

79

80

81

82

83

84

85

86

87

88

89

90

91

92

93

94

95

96

97

98

99

100

101

102

103

104

105

106

107

108

109

110

111

112

113

114

115

116

117

118

119

120

121

122

123

124

125

126

127

128

129

130

131

132

133

134

135

136

137

138

139

140

141

142

143

144

145

146

147

148

149

150

151

152

153

154

155

156

157

158

159

160

161

162

163

164

165

166

167

# http://pyrocko.org - GPLv3 

# 

# The Pyrocko Developers, 21st Century 

# ---|P------/S----------~Lg---------- 

 

try: 

import queue 

except ImportError: 

import Queue as queue 

 

 

import multiprocessing 

import traceback 

import errno 

 

 

def worker( 

q_in, q_out, function, eprintignore, pshared, 

startup, startup_args, cleanup): 

 

kwargs = {} 

if pshared is not None: 

kwargs['pshared'] = pshared 

 

if startup is not None: 

startup(*startup_args) 

 

while True: 

i, args = q_in.get() 

if i is None: 

if cleanup is not None: 

cleanup() 

 

break 

 

res, exception = None, None 

try: 

res = function(*args, **kwargs) 

except Exception as e: 

if eprintignore is None or not isinstance(e, eprintignore): 

traceback.print_exc() 

exception = e 

q_out.put((i, res, exception)) 

 

 

def parimap(function, *iterables, **kwargs): 

assert all( 

k in ( 

'nprocs', 'eprintignore', 'pshared', 'startup', 'startup_args', 

'cleanup') 

for k in kwargs.keys()) 

 

nprocs = kwargs.get('nprocs', None) 

eprintignore = kwargs.get('eprintignore', 'all') 

pshared = kwargs.get('pshared', None) 

startup = kwargs.get('startup', None) 

startup_args = kwargs.get('startup_args', ()) 

cleanup = kwargs.get('cleanup', None) 

 

if eprintignore == 'all': 

eprintignore = None 

 

if nprocs == 1: 

iterables = list(map(iter, iterables)) 

kwargs = {} 

if pshared is not None: 

kwargs['pshared'] = pshared 

 

while True: 

args = [] 

for it in iterables: 

try: 

args.append(next(it)) 

except StopIteration: 

return 

 

yield function(*args, **kwargs) 

 

return 

 

if nprocs is None: 

nprocs = multiprocessing.cpu_count() 

 

q_in = multiprocessing.Queue(1) 

q_out = multiprocessing.Queue() 

 

procs = [] 

 

results = [] 

nrun = 0 

nwritten = 0 

iout = 0 

all_written = False 

error_ahead = False 

iterables = list(map(iter, iterables)) 

while True: 

if nrun < nprocs and not all_written and not error_ahead: 

args = [] 

for it in iterables: 

try: 

args.append(next(it)) 

except StopIteration: 

pass 

 

if len(args) == len(iterables): 

if len(procs) < nrun + 1: 

p = multiprocessing.Process( 

target=worker, 

args=(q_in, q_out, function, eprintignore, pshared, 

startup, startup_args, cleanup)) 

p.daemon = True 

p.start() 

procs.append(p) 

 

q_in.put((nwritten, args)) 

nwritten += 1 

nrun += 1 

else: 

all_written = True 

[q_in.put((None, None)) for p in procs] 

q_in.close() 

 

try: 

while nrun > 0: 

if nrun < nprocs and not all_written and not error_ahead: 

results.append(q_out.get_nowait()) 

else: 

while True: 

try: 

results.append(q_out.get()) 

break 

except IOError as e: 

if e.errno != errno.EINTR: 

raise 

 

nrun -= 1 

 

except queue.Empty: 

pass 

 

if results: 

results.sort() 

# check for error ahead to prevent further enqueuing 

if any(exc for (_, _, exc) in results): 

error_ahead = True 

 

while results: 

(i, r, exc) = results[0] 

if i == iout: 

results.pop(0) 

if exc is not None: 

if not all_written: 

[q_in.put((None, None)) for p in procs] 

q_in.close() 

raise exc 

else: 

yield r 

 

iout += 1 

else: 

break 

 

if all_written and nrun == 0: 

break 

 

[p.join() for p in procs] 

return