MCPcopy Create free account
hub / github.com/Boris-code/feapder / flush

Method flush

feapder/buffer/item_buffer.py:126–172  ·  view source on GitHub ↗
(self)

Source from the content-addressed store, hash-verified

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

Callers 1

runMethod · 0.95

Calls 3

__add_item_to_dbMethod · 0.95
emptyMethod · 0.80
exceptionMethod · 0.80

Tested by

no test coverage detected