perf: Optimize JobQueueService for performance and clarity improvements (#6690)

* ⚡️ Speed up method `JobQueueService.create_queue` by 630% in PR #5940 (`refactor-build-flow`)
Certainly! To optimize the provided `JobQueueService` class to run faster, we should ensure the code efficiency in creating event managers and handling queues. We can speed up the process by making the following changes.

1. Remove redundant logger calls if they don't provide essential debugging or monitoring information.
2. Inline `create_default_event_manager` function to avoid redundant function calls while ensuring that the `EventManager` setup remains efficient.
3. Ensure faster dictionary access and error message handling.

Here is the rewritten and optimized code.

In this optimization.
1. The `create_queue` method was streamlined to reduce redundancy checks and fewer logger messages.
2. The `create_default_event_manager` function is integrated directly as `_create_default_event_manager` private method within the class.
3. Utilized list iteration instead of multiple function calls to register events efficiently in the `EventManager`.

These changes aim to maintain functionality while improving performance where possible, particularly in the setup process.

* fix: improve error messages in JobQueueService for better clarity

* [autofix.ci] apply automated fixes

* refactor: Relax event type constraint in EventManager

Change event type from a strict literal to a more flexible string type, allowing for greater extensibility in event registration

* [autofix.ci] apply automated fixes

* [autofix.ci] apply automated fixes

* fix: update send_event method to use str for event_type

* fix: change event_type parameter to str and handle invalid types in create_event_by_type

* fix: Add noqa comment to suppress linting warning for access_host assignment

---------

Co-authored-by: codeflash-ai[bot] <148906541+codeflash-ai[bot]@users.noreply.github.com>
Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
This commit is contained in:
Gabriel Luiz Freitas Almeida 2025-06-03 12:57:25 -03:00 • committed by GitHub
commit f21d23fb58
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 375 additions and 19 deletions

View file

@ -415,7 +415,7 @@ def print_banner(host: str, port: int, protocol: str) -> None:
"To contribute, set: [bold]DO_NOT_TRACK=false[/bold] in your environment."
)
)
access_host = host if host != "0.0.0.0" else "localhost"
access_host = host if host != "0.0.0.0" else "localhost" # noqa: S104
access_link = f"[bold]🟢 Open Langflow →[/bold] [link={protocol}://{access_host}:{port}]{protocol}://{access_host}:{port}[/link]"
message = f"{title}\n{info_text}\n\n{telemetry_text}\n\n{access_link}"

View file

@ -5,7 +5,7 @@ import json
import time
import uuid
from functools import partial
from typing import TYPE_CHECKING, Literal
from typing import TYPE_CHECKING
from fastapi.encoders import jsonable_encoder
from loguru import logger
@ -50,7 +50,7 @@ class EventManager:
def register_event(
self,
name: str,
event_type: Literal["message", "error", "warning", "info", "token"],
event_type: str,
callback: EventCallback | None = None,
) -> None:
if not name:
@ -65,7 +65,7 @@ class EventManager:
callback_ = partial(callback, manager=self, event_type=event_type)
self.events[name] = callback_
def send_event(self, *, event_type: Literal["message", "error", "warning", "info", "token"], data: LoggableType):
def send_event(self, *, event_type: str, data: LoggableType):
try:
if isinstance(data, dict) and event_type in {"message", "error", "warning", "info", "token"}:
data = create_event_by_type(event_type, **data)

View file

@ -512,7 +512,174 @@
},
{
"data": {
"id": "ChatOutput-meCWe",
"id": "ParseData-hLbU0",
"node": {
"base_classes": [
"Data",
"Message"
],
"beta": false,
"category": "processing",
"conditional_paths": [],
"custom_fields": {},
"description": "Convert Data objects into Messages using any {field_name} from input data.",
"display_name": "Data to Message",
"documentation": "",
"edited": false,
"field_order": [
"data",
"template",
"sep"
],
"frozen": false,
"icon": "message-square",
"key": "ParseData",
"legacy": true,
"lf_version": "1.1.5",
"metadata": {
"legacy_name": "Parse Data"
},
"minimized": false,
"output_types": [],
"outputs": [
{
"allows_loop": false,
"cache": true,
"display_name": "Message",
"group_outputs": false,
"method": "parse_data",
"name": "text",
"selected": "Message",
"tool_mode": true,
"types": [
"Message"
],
"value": "__UNDEFINED__"
},
{
"allows_loop": false,
"cache": true,
"display_name": "Data List",
"group_outputs": false,
"method": "parse_data_as_list",
"name": "data_list",
"selected": "Data",
"tool_mode": true,
"types": [
"Data"
],
"value": "__UNDEFINED__"
}
],
"pinned": false,
"score": 0.23285358167685585,
"template": {
"_type": "Component",
"code": {
"advanced": true,
"dynamic": true,
"fileTypes": [],
"file_path": "",
"info": "",
"list": false,
"load_from_db": false,
"multiline": true,
"name": "code",
"password": false,
"placeholder": "",
"required": true,
"show": true,
"title_case": false,
"type": "code",
"value": "from langflow.custom import Component\nfrom langflow.helpers.data import data_to_text, data_to_text_list\nfrom langflow.io import DataInput, MultilineInput, Output, StrInput\nfrom langflow.schema import Data\nfrom langflow.schema.message import Message\n\n\nclass ParseDataComponent(Component):\n display_name = \"Data to Message\"\n description = \"Convert Data objects into Messages using any {field_name} from input data.\"\n icon = \"message-square\"\n name = \"ParseData\"\n legacy = True\n metadata = {\n \"legacy_name\": \"Parse Data\",\n }\n\n inputs = [\n DataInput(\n name=\"data\",\n display_name=\"Data\",\n info=\"The data to convert to text.\",\n is_list=True,\n required=True,\n ),\n MultilineInput(\n name=\"template\",\n display_name=\"Template\",\n info=\"The template to use for formatting the data. \"\n \"It can contain the keys {text}, {data} or any other key in the Data.\",\n value=\"{text}\",\n required=True,\n ),\n StrInput(name=\"sep\", display_name=\"Separator\", advanced=True, value=\"\\n\"),\n ]\n\n outputs = [\n Output(\n display_name=\"Message\",\n name=\"text\",\n info=\"Data as a single Message, with each input Data separated by Separator\",\n method=\"parse_data\",\n ),\n Output(\n display_name=\"Data List\",\n name=\"data_list\",\n info=\"Data as a list of new Data, each having `text` formatted by Template\",\n method=\"parse_data_as_list\",\n ),\n ]\n\n def _clean_args(self) -> tuple[list[Data], str, str]:\n data = self.data if isinstance(self.data, list) else [self.data]\n template = self.template\n sep = self.sep\n return data, template, sep\n\n def parse_data(self) -> Message:\n data, template, sep = self._clean_args()\n result_string = data_to_text(template, data, sep)\n self.status = result_string\n return Message(text=result_string)\n\n def parse_data_as_list(self) -> list[Data]:\n data, template, _ = self._clean_args()\n text_list, data_list = data_to_text_list(template, data)\n for item, text in zip(data_list, text_list, strict=True):\n item.set_text(text)\n self.status = data_list\n return data_list\n"
},
"data": {
"_input_type": "DataInput",
"advanced": false,
"display_name": "Data",
"dynamic": false,
"info": "The data to convert to text.",
"input_types": [
"Data"
],
"list": true,
"list_add_label": "Add More",
"name": "data",
"placeholder": "",
"required": true,
"show": true,
"title_case": false,
"tool_mode": false,
"trace_as_input": true,
"trace_as_metadata": true,
"type": "other",
"value": ""
},
"sep": {
"_input_type": "StrInput",
"advanced": true,
"display_name": "Separator",
"dynamic": false,
"info": "",
"list": false,
"list_add_label": "Add More",
"load_from_db": false,
"name": "sep",
"placeholder": "",
"required": false,
"show": true,
"title_case": false,
"tool_mode": false,
"trace_as_metadata": true,
"type": "str",
"value": "\n"
},
"template": {
"_input_type": "MultilineInput",
"advanced": false,
"display_name": "Template",
"dynamic": false,
"info": "The template to use for formatting the data. It can contain the keys {text}, {data} or any other key in the Data.",
"input_types": [
"Message"
],
"list": false,
"list_add_label": "Add More",
"load_from_db": false,
"multiline": true,
"name": "template",
"placeholder": "",
"required": true,
"show": true,
"title_case": false,
"tool_mode": false,
"trace_as_input": true,
"trace_as_metadata": true,
"type": "str",
"value": "EBITIDA: {EBITIDA} , \nNet Income: {NET_INCOME} ,\nGROSS_PROFIT: {GROSS_PROFIT}"
}
},
"tool_mode": false
},
"showNode": true,
"type": "ParseData"
},
"dragging": false,
"id": "ParseData-hLbU0",
"measured": {
"height": 342,
"width": 320
},
"position": {
"x": 1788.8753184818502,
"y": 247.15776278878775
},
"selected": false,
"type": "genericNode"
},
{
"data": {
"id": "ChatOutput-ZtLkt",
"node": {
"base_classes": [
"Message"

View file

@ -432,7 +432,170 @@
},
{
"data": {
"id": "OpenAIModel-ac6TO",
"id": "ParseData-LUfjb",
"node": {
"base_classes": [
"Data",
"Message"
],
"beta": false,
"conditional_paths": [],
"custom_fields": {},
"description": "Convert Data objects into Messages using any {field_name} from input data.",
"display_name": "Data to Message",
"documentation": "",
"edited": false,
"field_order": [
"data",
"template",
"sep"
],
"frozen": false,
"icon": "message-square",
"legacy": true,
"lf_version": "1.1.5",
"metadata": {
"legacy_name": "Parse Data"
},
"minimized": false,
"output_types": [],
"outputs": [
{
"allows_loop": false,
"cache": true,
"display_name": "Message",
"group_outputs": false,
"method": "parse_data",
"name": "text",
"selected": "Message",
"tool_mode": true,
"types": [
"Message"
],
"value": "__UNDEFINED__"
},
{
"allows_loop": false,
"cache": true,
"display_name": "Data List",
"group_outputs": false,
"method": "parse_data_as_list",
"name": "data_list",
"selected": "Data",
"tool_mode": true,
"types": [
"Data"
],
"value": "__UNDEFINED__"
}
],
"pinned": false,
"template": {
"_type": "Component",
"code": {
"advanced": true,
"dynamic": true,
"fileTypes": [],
"file_path": "",
"info": "",
"list": false,
"load_from_db": false,
"multiline": true,
"name": "code",
"password": false,
"placeholder": "",
"required": true,
"show": true,
"title_case": false,
"type": "code",
"value": "from langflow.custom import Component\nfrom langflow.helpers.data import data_to_text, data_to_text_list\nfrom langflow.io import DataInput, MultilineInput, Output, StrInput\nfrom langflow.schema import Data\nfrom langflow.schema.message import Message\n\n\nclass ParseDataComponent(Component):\n display_name = \"Data to Message\"\n description = \"Convert Data objects into Messages using any {field_name} from input data.\"\n icon = \"message-square\"\n name = \"ParseData\"\n legacy = True\n metadata = {\n \"legacy_name\": \"Parse Data\",\n }\n\n inputs = [\n DataInput(\n name=\"data\",\n display_name=\"Data\",\n info=\"The data to convert to text.\",\n is_list=True,\n required=True,\n ),\n MultilineInput(\n name=\"template\",\n display_name=\"Template\",\n info=\"The template to use for formatting the data. \"\n \"It can contain the keys {text}, {data} or any other key in the Data.\",\n value=\"{text}\",\n required=True,\n ),\n StrInput(name=\"sep\", display_name=\"Separator\", advanced=True, value=\"\\n\"),\n ]\n\n outputs = [\n Output(\n display_name=\"Message\",\n name=\"text\",\n info=\"Data as a single Message, with each input Data separated by Separator\",\n method=\"parse_data\",\n ),\n Output(\n display_name=\"Data List\",\n name=\"data_list\",\n info=\"Data as a list of new Data, each having `text` formatted by Template\",\n method=\"parse_data_as_list\",\n ),\n ]\n\n def _clean_args(self) -> tuple[list[Data], str, str]:\n data = self.data if isinstance(self.data, list) else [self.data]\n template = self.template\n sep = self.sep\n return data, template, sep\n\n def parse_data(self) -> Message:\n data, template, sep = self._clean_args()\n result_string = data_to_text(template, data, sep)\n self.status = result_string\n return Message(text=result_string)\n\n def parse_data_as_list(self) -> list[Data]:\n data, template, _ = self._clean_args()\n text_list, data_list = data_to_text_list(template, data)\n for item, text in zip(data_list, text_list, strict=True):\n item.set_text(text)\n self.status = data_list\n return data_list\n"
},
"data": {
"_input_type": "DataInput",
"advanced": false,
"display_name": "Data",
"dynamic": false,
"info": "The data to convert to text.",
"input_types": [
"Data"
],
"list": true,
"list_add_label": "Add More",
"name": "data",
"placeholder": "",
"required": true,
"show": true,
"title_case": false,
"tool_mode": false,
"trace_as_input": true,
"trace_as_metadata": true,
"type": "other",
"value": ""
},
"sep": {
"_input_type": "StrInput",
"advanced": true,
"display_name": "Separator",
"dynamic": false,
"info": "",
"list": false,
"list_add_label": "Add More",
"load_from_db": false,
"name": "sep",
"placeholder": "",
"required": false,
"show": true,
"title_case": false,
"tool_mode": false,
"trace_as_metadata": true,
"type": "str",
"value": "\n"
},
"template": {
"_input_type": "MultilineInput",
"advanced": false,
"display_name": "Template",
"dynamic": false,
"info": "The template to use for formatting the data. It can contain the keys {text}, {data} or any other key in the Data.",
"input_types": [
"Message"
],
"list": false,
"list_add_label": "Add More",
"load_from_db": false,
"multiline": true,
"name": "template",
"placeholder": "",
"required": true,
"show": true,
"title_case": false,
"tool_mode": false,
"trace_as_input": true,
"trace_as_metadata": true,
"type": "str",
"value": "{text}"
}
},
"tool_mode": false
},
"showNode": true,
"type": "ParseData"
},
"id": "ParseData-LUfjb",
"measured": {
"height": 342,
"width": 320
},
"position": {
"x": 1330.927281184057,
"y": 382.3516758942169
},
"selected": false,
"type": "genericNode"
},
{
"data": {
"id": "OpenAIModel-iudDZ",
"node": {
"base_classes": [
"LanguageModel",

View file

@ -170,12 +170,13 @@ _EVENT_CREATORS: dict[str, tuple[Callable, inspect.Signature]] = {
}
def create_event_by_type(
event_type: Literal["message", "error", "warning", "info", "token"], **kwargs
) -> PlaygroundEvent | dict:
def create_event_by_type(event_type: str, **kwargs) -> PlaygroundEvent | dict:
if event_type not in _EVENT_CREATORS:
return kwargs
creator_func, signature = _EVENT_CREATORS[event_type]
try:
creator_func, signature = _EVENT_CREATORS[event_type]
except KeyError as e:
msg = f"Invalid event type: {event_type}"
raise ValueError(msg) from e
valid_params = {k: v for k, v in kwargs.items() if k in signature.parameters}
return creator_func(**valid_params)

View file

@ -4,7 +4,7 @@ import asyncio
from loguru import logger
from langflow.events.event_manager import EventManager, create_default_event_manager
from langflow.events.event_manager import EventManager
from langflow.services.base import Service
@ -133,18 +133,17 @@ class JobQueueService(Service):
- The asyncio.Queue instance for handling the job's tasks or messages.
- The EventManager instance for event handling tied to the queue.
"""
if job_id in self._queues:
msg = f"Queue for job_id {job_id} already exists"
logger.error(msg)
raise ValueError(msg)
if self._closed:
msg = "Queue service is closed"
logger.error(msg)
raise RuntimeError(msg)
existing_queue = self._queues.get(job_id)
if existing_queue:
msg = f"Queue for job_id {job_id} already exists"
raise ValueError(msg)
main_queue: asyncio.Queue = asyncio.Queue()
event_manager = create_default_event_manager(main_queue)
event_manager: EventManager = self._create_default_event_manager(main_queue)
# Register the queue without an active task.
self._queues[job_id] = (main_queue, event_manager, None, None)
@ -300,3 +299,29 @@ class JobQueueService(Service):
# Enough time has passed, perform the actual cleanup
logger.debug(f"Cleaning up job_id {job_id} after grace period")
await self.cleanup_job(job_id)
def _create_default_event_manager(self, queue: asyncio.Queue) -> EventManager:
"""Creates the default event manager with predefined events.
Args:
queue (asyncio.Queue): The queue to be associated with the event manager.
Returns:
EventManager: The configured EventManager instance.
"""
manager = EventManager(queue)
# Registering predefined events
event_names_types = [
("on_token", "token"),
("on_vertices_sorted", "vertices_sorted"),
("on_error", "error"),
("on_end", "end"),
("on_message", "add_message"),
("on_remove_message", "remove_message"),
("on_end_vertex", "end_vertex"),
("on_build_start", "build_start"),
("on_build_end", "build_end"),
]
for name, event_type in event_names_types:
manager.register_event(name, event_type)
return manager