运行引擎并返回流式响应的迭代器。 Args: req_ids (list[str]): 请求ID列表 prompts: 原始提示词列表,用于设置到输出中 use_tqdm (bool, optional): 是否使用tqdm进度条 topk_logprobs (Optional[int]): 返回的top-k logprobs数量 Yields: list[RequestOutput]: 包含增量更新的部分响应列表
(
self,
req_ids: list[str],
prompts,
use_tqdm: bool,
topk_logprobs: Optional[int] = None,
chat_template_kwargs: Optional[dict[str, Any]] = None,
)
| 596 | return output |
| 597 | |
| 598 | def _run_engine_stream( |
| 599 | self, |
| 600 | req_ids: list[str], |
| 601 | prompts, |
| 602 | use_tqdm: bool, |
| 603 | topk_logprobs: Optional[int] = None, |
| 604 | chat_template_kwargs: Optional[dict[str, Any]] = None, |
| 605 | ): |
| 606 | """ |
| 607 | 运行引擎并返回流式响应的迭代器。 |
| 608 | |
| 609 | Args: |
| 610 | req_ids (list[str]): 请求ID列表 |
| 611 | prompts: 原始提示词列表,用于设置到输出中 |
| 612 | use_tqdm (bool, optional): 是否使用tqdm进度条 |
| 613 | topk_logprobs (Optional[int]): 返回的top-k logprobs数量 |
| 614 | |
| 615 | Yields: |
| 616 | list[RequestOutput]: 包含增量更新的部分响应列表 |
| 617 | """ |
| 618 | # Initialize tqdm |
| 619 | if use_tqdm: |
| 620 | num_requests = len(req_ids) |
| 621 | pbar = tqdm( |
| 622 | total=num_requests, |
| 623 | desc="Processed prompts", |
| 624 | dynamic_ncols=True, |
| 625 | postfix=(f"est. speed input: {0:.2f} toks/s, " f"output: {0:.2f} toks/s"), |
| 626 | ) |
| 627 | |
| 628 | num_requests = len(req_ids) |
| 629 | original_num_requests = len(req_ids) # Keep track of original count |
| 630 | output = [None] * original_num_requests |
| 631 | req_ids_with_pos = [(pos, req_id) for pos, req_id in enumerate(req_ids)] |
| 632 | |
| 633 | # Track previous token counts for each request to identify new tokens |
| 634 | previous_token_counts = {req_id: 0 for req_id in req_ids} |
| 635 | |
| 636 | while num_requests > 0: |
| 637 | has_new_tokens = False |
| 638 | finished = [] |
| 639 | |
| 640 | for i, (pos, req_id) in enumerate(req_ids_with_pos): |
| 641 | with self.mutex: |
| 642 | if req_id not in self.req_output: |
| 643 | continue |
| 644 | |
| 645 | current_result = self.req_output[req_id] |
| 646 | current_token_count = ( |
| 647 | len(current_result.outputs.token_ids) if current_result.outputs.token_ids else 0 |
| 648 | ) |
| 649 | previous_count = previous_token_counts[req_id] |
| 650 | |
| 651 | # Check if there are new tokens since last yield |
| 652 | if current_token_count > previous_count: |
| 653 | has_new_tokens = True |
| 654 | # Create incremental output with only new tokens |
| 655 | incremental_result = self._create_incremental_result( |