Total Complexity | 35 |
Total Lines | 135 |
Duplicated Lines | 21.48 % |
Duplicate code is one of the most pungent code smells. A rule that is often used is to re-structure code once it is duplicated in three or more places.
Common duplication problems, and corresponding solutions are:
1 | #!/usr/bin/env python |
||
52 | class PikaQueue(object): |
||
53 | """ |
||
54 | A Queue like rabbitmq connector |
||
55 | """ |
||
56 | |||
57 | Empty = BaseQueue.Empty |
||
58 | Full = BaseQueue.Full |
||
59 | max_timeout = 0.3 |
||
60 | |||
61 | View Code Duplication | def __init__(self, name, amqp_url='amqp://guest:guest@localhost:5672/%2F', |
|
|
|||
62 | maxsize=0, lazy_limit=True): |
||
63 | """ |
||
64 | Constructor for a PikaQueue. |
||
65 | |||
66 | Not works with python 3. Default for python 2. |
||
67 | |||
68 | amqp_url: https://www.rabbitmq.com/uri-spec.html |
||
69 | maxsize: an integer that sets the upperbound limit on the number of |
||
70 | items that can be placed in the queue. |
||
71 | lazy_limit: as rabbitmq is shared between multipul instance, for a strict |
||
72 | limit on the number of items in the queue. PikaQueue have to |
||
73 | update current queue size before every put operation. When |
||
74 | `lazy_limit` is enabled, PikaQueue will check queue size every |
||
75 | max_size / 10 put operation for better performace. |
||
76 | """ |
||
77 | self.name = name |
||
78 | self.amqp_url = amqp_url |
||
79 | self.maxsize = maxsize |
||
80 | self.lock = threading.RLock() |
||
81 | |||
82 | self.lazy_limit = lazy_limit |
||
83 | if self.lazy_limit and self.maxsize: |
||
84 | self.qsize_diff_limit = int(self.maxsize * 0.1) |
||
85 | else: |
||
86 | self.qsize_diff_limit = 0 |
||
87 | self.qsize_diff = 0 |
||
88 | |||
89 | self.reconnect() |
||
90 | |||
91 | def reconnect(self): |
||
92 | """Reconnect to rabbitmq server""" |
||
93 | import pika |
||
94 | import pika.exceptions |
||
95 | |||
96 | self.connection = pika.BlockingConnection(pika.URLParameters(self.amqp_url)) |
||
97 | self.channel = self.connection.channel() |
||
98 | try: |
||
99 | self.channel.queue_declare(self.name) |
||
100 | except pika.exceptions.ChannelClosed: |
||
101 | self.connection = pika.BlockingConnection(pika.URLParameters(self.amqp_url)) |
||
102 | self.channel = self.connection.channel() |
||
103 | #self.channel.queue_purge(self.name) |
||
104 | |||
105 | @catch_error |
||
106 | def qsize(self): |
||
107 | with self.lock: |
||
108 | ret = self.channel.queue_declare(self.name, passive=True) |
||
109 | return ret.method.message_count |
||
110 | |||
111 | def empty(self): |
||
112 | if self.qsize() == 0: |
||
113 | return True |
||
114 | else: |
||
115 | return False |
||
116 | |||
117 | def full(self): |
||
118 | if self.maxsize and self.qsize() >= self.maxsize: |
||
119 | return True |
||
120 | else: |
||
121 | return False |
||
122 | |||
123 | @catch_error |
||
124 | def put(self, obj, block=True, timeout=None): |
||
125 | if not block: |
||
126 | return self.put_nowait() |
||
127 | |||
128 | start_time = time.time() |
||
129 | while True: |
||
130 | try: |
||
131 | return self.put_nowait(obj) |
||
132 | except BaseQueue.Full: |
||
133 | if timeout: |
||
134 | lasted = time.time() - start_time |
||
135 | if timeout > lasted: |
||
136 | time.sleep(min(self.max_timeout, timeout - lasted)) |
||
137 | else: |
||
138 | raise |
||
139 | else: |
||
140 | time.sleep(self.max_timeout) |
||
141 | |||
142 | @catch_error |
||
143 | def put_nowait(self, obj): |
||
144 | if self.lazy_limit and self.qsize_diff < self.qsize_diff_limit: |
||
145 | pass |
||
146 | elif self.full(): |
||
147 | raise BaseQueue.Full |
||
148 | else: |
||
149 | self.qsize_diff = 0 |
||
150 | with self.lock: |
||
151 | self.qsize_diff += 1 |
||
152 | return self.channel.basic_publish("", self.name, umsgpack.packb(obj)) |
||
153 | |||
154 | @catch_error |
||
155 | def get(self, block=True, timeout=None, ack=False): |
||
156 | if not block: |
||
157 | return self.get_nowait() |
||
158 | |||
159 | start_time = time.time() |
||
160 | while True: |
||
161 | try: |
||
162 | return self.get_nowait(ack) |
||
163 | except BaseQueue.Empty: |
||
164 | if timeout: |
||
165 | lasted = time.time() - start_time |
||
166 | if timeout > lasted: |
||
167 | time.sleep(min(self.max_timeout, timeout - lasted)) |
||
168 | else: |
||
169 | raise |
||
170 | else: |
||
171 | time.sleep(self.max_timeout) |
||
172 | |||
173 | @catch_error |
||
174 | def get_nowait(self, ack=False): |
||
175 | with self.lock: |
||
176 | method_frame, header_frame, body = self.channel.basic_get(self.name, not ack) |
||
177 | if method_frame is None: |
||
178 | raise BaseQueue.Empty |
||
179 | if ack: |
||
180 | self.channel.basic_ack(method_frame.delivery_tag) |
||
181 | return umsgpack.unpackb(body) |
||
182 | |||
183 | @catch_error |
||
184 | def delete(self): |
||
185 | with self.lock: |
||
186 | return self.channel.queue_delete(queue=self.name) |
||
187 | |||
271 |