From f21d23fb58c26c4a94526f44dbd9a61a70c50975 Mon Sep 17 00:00:00 2001 From: Gabriel Luiz Freitas Almeida Date: Tue, 3 Jun 2025 12:57:25 -0300 Subject: [PATCH] perf: Optimize JobQueueService for performance and clarity improvements (#6690) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * ⚡️ 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> --- src/backend/base/langflow/__main__.py | 2 +- .../base/langflow/events/event_manager.py | 6 +- .../Financial Report Parser.json | 169 +++++++++++++++++- .../starter_projects/Meeting Summary.json | 165 ++++++++++++++++- .../base/langflow/schema/playground_events.py | 11 +- .../langflow/services/job_queue/service.py | 41 ++++- 6 files changed, 375 insertions(+), 19 deletions(-) diff --git a/src/backend/base/langflow/__main__.py b/src/backend/base/langflow/__main__.py index a1932baed..381628bc3 100644 --- a/src/backend/base/langflow/__main__.py +++ b/src/backend/base/langflow/__main__.py @@ -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}" diff --git a/src/backend/base/langflow/events/event_manager.py b/src/backend/base/langflow/events/event_manager.py index e499030d9..9d879809e 100644 --- a/src/backend/base/langflow/events/event_manager.py +++ b/src/backend/base/langflow/events/event_manager.py @@ -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) diff --git a/src/backend/base/langflow/initial_setup/starter_projects/Financial Report Parser.json b/src/backend/base/langflow/initial_setup/starter_projects/Financial Report Parser.json index 23d133b29..84cad8de9 100644 --- a/src/backend/base/langflow/initial_setup/starter_projects/Financial Report Parser.json +++ b/src/backend/base/langflow/initial_setup/starter_projects/Financial Report Parser.json @@ -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" diff --git a/src/backend/base/langflow/initial_setup/starter_projects/Meeting Summary.json b/src/backend/base/langflow/initial_setup/starter_projects/Meeting Summary.json index 676d739b5..3e5d29c40 100644 --- a/src/backend/base/langflow/initial_setup/starter_projects/Meeting Summary.json +++ b/src/backend/base/langflow/initial_setup/starter_projects/Meeting Summary.json @@ -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", diff --git a/src/backend/base/langflow/schema/playground_events.py b/src/backend/base/langflow/schema/playground_events.py index d7db91ddb..de58af959 100644 --- a/src/backend/base/langflow/schema/playground_events.py +++ b/src/backend/base/langflow/schema/playground_events.py @@ -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) diff --git a/src/backend/base/langflow/services/job_queue/service.py b/src/backend/base/langflow/services/job_queue/service.py index a05a94a51..cee9cbb0b 100644 --- a/src/backend/base/langflow/services/job_queue/service.py +++ b/src/backend/base/langflow/services/job_queue/service.py @@ -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