| 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 |