(self)
| 124 | self._items_queue.put(item) |
| 125 | |
| 126 | def flush(self): |
| 127 | try: |
| 128 | items = [] |
| 129 | update_items = [] |
| 130 | requests = [] |
| 131 | callbacks = [] |
| 132 | items_fingerprints = [] |
| 133 | data_count = 0 |
| 134 | |
| 135 | while not self._items_queue.empty(): |
| 136 | data = self._items_queue.get_nowait() |
| 137 | data_count += 1 |
| 138 | |
| 139 | # data 分类 |
| 140 | if callable(data): |
| 141 | callbacks.append(data) |
| 142 | |
| 143 | elif isinstance(data, UpdateItem): |
| 144 | update_items.append(data) |
| 145 | |
| 146 | elif isinstance(data, Item): |
| 147 | items.append(data) |
| 148 | if setting.ITEM_FILTER_ENABLE: |
| 149 | items_fingerprints.append(data.fingerprint) |
| 150 | |
| 151 | else: # request-redis |
| 152 | requests.append(data) |
| 153 | |
| 154 | if data_count >= setting.ITEM_UPLOAD_BATCH_MAX_SIZE: |
| 155 | self.__add_item_to_db( |
| 156 | items, update_items, requests, callbacks, items_fingerprints |
| 157 | ) |
| 158 | |
| 159 | items = [] |
| 160 | update_items = [] |
| 161 | requests = [] |
| 162 | callbacks = [] |
| 163 | items_fingerprints = [] |
| 164 | data_count = 0 |
| 165 | |
| 166 | if data_count: |
| 167 | self.__add_item_to_db( |
| 168 | items, update_items, requests, callbacks, items_fingerprints |
| 169 | ) |
| 170 | |
| 171 | except Exception as e: |
| 172 | log.exception(e) |
| 173 | |
| 174 | def get_items_count(self): |
| 175 | return self._items_queue.qsize() |
no test coverage detected