forked from binux/pyspider
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathtest_message_queue.py
232 lines (196 loc) · 8.1 KB
/
test_message_queue.py
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
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
#!/usr/bin/env python
# -*- encoding: utf-8 -*-
# vim: set et sw=4 ts=4 sts=4 ff=unix fenc=utf8:
# Author: Binux<[email protected]>
# http://binux.me
# Created on 2014-10-07 10:33:38
import os
import six
import time
import unittest
from pyspider.libs import utils
from six.moves import queue as Queue
class TestMessageQueue(object):
@classmethod
def setUpClass(self):
raise NotImplementedError
def test_10_put(self):
self.assertEqual(self.q1.qsize(), 0)
self.assertEqual(self.q2.qsize(), 0)
self.q1.put('TEST_DATA1', timeout=3)
self.q1.put('TEST_DATA2_中文', timeout=3)
time.sleep(0.01)
self.assertEqual(self.q1.qsize(), 2)
self.assertEqual(self.q2.qsize(), 2)
def test_20_get(self):
self.assertEqual(self.q1.get(timeout=0.01), 'TEST_DATA1')
self.assertEqual(self.q2.get_nowait(), 'TEST_DATA2_中文')
with self.assertRaises(Queue.Empty):
self.q2.get(timeout=0.01)
with self.assertRaises(Queue.Empty):
self.q2.get_nowait()
def test_30_full(self):
self.assertEqual(self.q1.qsize(), 0)
self.assertEqual(self.q2.qsize(), 0)
for i in range(2):
self.q1.put_nowait('TEST_DATA%d' % i)
for i in range(3):
self.q2.put('TEST_DATA%d' % i)
with self.assertRaises(Queue.Full):
self.q1.put('TEST_DATA6', timeout=0.01)
with self.assertRaises(Queue.Full):
self.q1.put_nowait('TEST_DATA6')
def test_40_multiple_threading_error(self):
def put(q):
for i in range(100):
q.put("DATA_%d" % i)
def get(q):
for i in range(100):
q.get()
t = utils.run_in_thread(put, self.q3)
get(self.q3)
t.join()
class BuiltinQueue(TestMessageQueue, unittest.TestCase):
@classmethod
def setUpClass(self):
from pyspider.message_queue import connect_message_queue
with utils.timeout(3):
self.q1 = self.q2 = connect_message_queue('test_queue', maxsize=5)
self.q3 = connect_message_queue('test_queue_for_threading_test')
#@unittest.skipIf(six.PY3, 'pika not suport python 3')
@unittest.skipIf(os.environ.get('IGNORE_RABBITMQ') or os.environ.get('IGNORE_ALL'), 'no rabbitmq server for test.')
class TestPikaRabbitMQ(TestMessageQueue, unittest.TestCase):
@classmethod
def setUpClass(self):
from pyspider.message_queue import rabbitmq
with utils.timeout(3):
self.q1 = rabbitmq.PikaQueue('test_queue', maxsize=5, lazy_limit=False)
self.q2 = rabbitmq.PikaQueue('test_queue', amqp_url='amqp://localhost:5672/%2F', maxsize=5, lazy_limit=False)
self.q3 = rabbitmq.PikaQueue('test_queue_for_threading_test', amqp_url='amqp://guest:guest@localhost:5672/', lazy_limit=False)
self.q2.delete()
self.q2.reconnect()
self.q3.delete()
self.q3.reconnect()
@classmethod
def tearDownClass(self):
self.q2.delete()
self.q3.delete()
del self.q1
del self.q2
del self.q3
def test_30_full(self):
self.assertEqual(self.q1.qsize(), 0)
self.assertEqual(self.q2.qsize(), 0)
for i in range(2):
self.q1.put_nowait('TEST_DATA%d' % i)
for i in range(3):
self.q2.put('TEST_DATA%d' % i)
print(self.q1.__dict__)
print(self.q1.qsize())
with self.assertRaises(Queue.Full):
self.q1.put_nowait('TEST_DATA6')
print(self.q1.__dict__)
print(self.q1.qsize())
with self.assertRaises(Queue.Full):
self.q1.put('TEST_DATA6', timeout=0.01)
@unittest.skipIf(six.PY3, 'Python 3 now using Pika')
@unittest.skipIf(os.environ.get('IGNORE_RABBITMQ') or os.environ.get('IGNORE_ALL'), 'no rabbitmq server for test.')
class TestAmqpRabbitMQ(TestMessageQueue, unittest.TestCase):
@classmethod
def setUpClass(self):
from pyspider.message_queue import connect_message_queue
with utils.timeout(3):
self.q1 = connect_message_queue('test_queue', 'amqp://localhost:5672/',
maxsize=5, lazy_limit=False)
self.q2 = connect_message_queue('test_queue', 'amqp://localhost:5672/%2F',
maxsize=5, lazy_limit=False)
self.q3 = connect_message_queue('test_queue_for_threading_test',
'amqp://guest:guest@localhost:5672/', lazy_limit=False)
self.q2.delete()
self.q2.reconnect()
self.q3.delete()
self.q3.reconnect()
@classmethod
def tearDownClass(self):
self.q2.delete()
self.q3.delete()
del self.q1
del self.q2
del self.q3
def test_30_full(self):
self.assertEqual(self.q1.qsize(), 0)
self.assertEqual(self.q2.qsize(), 0)
for i in range(2):
self.q1.put_nowait('TEST_DATA%d' % i)
for i in range(3):
self.q2.put('TEST_DATA%d' % i)
print(self.q1.__dict__)
print(self.q1.qsize())
with self.assertRaises(Queue.Full):
self.q1.put('TEST_DATA6', timeout=0.01)
print(self.q1.__dict__)
print(self.q1.qsize())
with self.assertRaises(Queue.Full):
self.q1.put_nowait('TEST_DATA6')
@unittest.skipIf(os.environ.get('IGNORE_REDIS') or os.environ.get('IGNORE_ALL'), 'no redis server for test.')
class TestRedisQueue(TestMessageQueue, unittest.TestCase):
@classmethod
def setUpClass(self):
from pyspider.message_queue import connect_message_queue
from pyspider.message_queue import redis_queue
with utils.timeout(3):
self.q1 = redis_queue.RedisQueue('test_queue', maxsize=5, lazy_limit=False)
self.q2 = redis_queue.RedisQueue('test_queue', maxsize=5, lazy_limit=False)
self.q3 = connect_message_queue('test_queue_for_threading_test',
'redis://localhost:6379/')
while not self.q1.empty():
self.q1.get()
while not self.q2.empty():
self.q2.get()
while not self.q3.empty():
self.q3.get()
@classmethod
def tearDownClass(self):
while not self.q1.empty():
self.q1.get()
while not self.q2.empty():
self.q2.get()
while not self.q3.empty():
self.q3.get()
class TestKombuQueue(TestMessageQueue, unittest.TestCase):
kombu_url = 'kombu+memory://'
@classmethod
def setUpClass(self):
from pyspider.message_queue import connect_message_queue
with utils.timeout(3):
self.q1 = connect_message_queue('test_queue', self.kombu_url, maxsize=5, lazy_limit=False)
self.q2 = connect_message_queue('test_queue', self.kombu_url, maxsize=5, lazy_limit=False)
self.q3 = connect_message_queue('test_queue_for_threading_test', self.kombu_url, lazy_limit=False)
while not self.q1.empty():
self.q1.get()
while not self.q2.empty():
self.q2.get()
while not self.q3.empty():
self.q3.get()
@classmethod
def tearDownClass(self):
while not self.q1.empty():
self.q1.get()
self.q1.delete()
while not self.q2.empty():
self.q2.get()
self.q2.delete()
while not self.q3.empty():
self.q3.get()
self.q3.delete()
@unittest.skip('test cannot pass, get is buffered')
@unittest.skipIf(os.environ.get('IGNORE_RABBITMQ') or os.environ.get('IGNORE_ALL'), 'no rabbitmq server for test.')
class TestKombuAmpqQueue(TestKombuQueue):
kombu_url = 'kombu+amqp://'
@unittest.skip('test cannot pass, put is buffered')
@unittest.skipIf(os.environ.get('IGNORE_REDIS') or os.environ.get('IGNORE_ALL'), 'no redis server for test.')
class TestKombuRedisQueue(TestKombuQueue):
kombu_url = 'kombu+redis://'
@unittest.skipIf(os.environ.get('IGNORE_MONGODB') or os.environ.get('IGNORE_ALL'), 'no mongodb server for test.')
class TestKombuMongoDBQueue(TestKombuQueue):
kombu_url = 'kombu+mongodb://'