| 117 | return self._is_adding_to_db |
| 118 | |
| 119 | def __add_request_to_db(self): |
| 120 | request_list = [] |
| 121 | prioritys = [] |
| 122 | callbacks = [] |
| 123 | |
| 124 | while self._requests_deque: |
| 125 | request = self._requests_deque.popleft() |
| 126 | self._is_adding_to_db = True |
| 127 | |
| 128 | if callable(request): |
| 129 | # 函数 |
| 130 | # 注意:应该考虑闭包情况。闭包情况可写成 |
| 131 | # def test(xxx = xxx): |
| 132 | # # TODO 业务逻辑 使用 xxx |
| 133 | # 这么写不会导致xxx为循环结束后的最后一个值 |
| 134 | callbacks.append(request) |
| 135 | continue |
| 136 | |
| 137 | priority = request.priority |
| 138 | |
| 139 | # 如果需要去重并且库中已重复 则continue |
| 140 | if self.is_exist_request(request): |
| 141 | continue |
| 142 | else: |
| 143 | request_list.append(str(request.to_dict)) |
| 144 | prioritys.append(priority) |
| 145 | |
| 146 | if len(request_list) > MAX_URL_COUNT: |
| 147 | self._db.zadd(self._table_request, request_list, prioritys) |
| 148 | request_list = [] |
| 149 | prioritys = [] |
| 150 | |
| 151 | # 入库 |
| 152 | if request_list: |
| 153 | self._db.zadd(self._table_request, request_list, prioritys) |
| 154 | |
| 155 | # 执行回调 |
| 156 | for callback in callbacks: |
| 157 | try: |
| 158 | callback() |
| 159 | except Exception as e: |
| 160 | log.exception(e) |
| 161 | |
| 162 | # 删除已做任务 |
| 163 | if self._del_requests_deque: |
| 164 | request_done_list = [] |
| 165 | while self._del_requests_deque: |
| 166 | request_done_list.append(self._del_requests_deque.popleft()) |
| 167 | |
| 168 | # 去掉request_list中的requests, 否则可能会将刚添加的request删除 |
| 169 | request_done_list = list(set(request_done_list) - set(request_list)) |
| 170 | |
| 171 | if request_done_list: |
| 172 | self._db.zrem(self._table_request, request_done_list) |
| 173 | |
| 174 | self._is_adding_to_db = False |