From 2126468f7f97e33bac526204de57f6f58a692c0a Mon Sep 17 00:00:00 2001 From: guohelu <19503896967@163.com> Date: Mon, 15 Dec 2025 19:36:17 +0800 Subject: [PATCH 1/4] =?UTF-8?q?feat:=20=E6=96=B0=E5=A2=9E=E6=B5=81?= =?UTF-8?q?=E7=A8=8B=E5=88=9B=E5=BB=BA=E4=BB=BB=E5=8A=A1=E5=B9=B6=E5=8F=91?= =?UTF-8?q?=E6=8E=A7=E5=88=B6=20--story=3D128208486=20#=20Reviewed,=20tran?= =?UTF-8?q?saction=20id:=2068146?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../collections/subprocess_plugin/v1_0_0.py | 10 +- bkflow/space/configs.py | 18 +++ bkflow/task/celery/tasks.py | 12 +- bkflow/task/domains/callback.py | 14 ++- bkflow/task/models.py | 3 +- bkflow/task/serializers.py | 6 +- bkflow/task/signals/handlers.py | 27 ++++- bkflow/task/utils.py | 111 ++++++++++++++++++ bkflow/task/views.py | 41 +++++-- bkflow/utils/space.py | 75 ++++++++++++ config/default.py | 2 + env.py | 2 + 12 files changed, 297 insertions(+), 24 deletions(-) create mode 100644 bkflow/utils/space.py diff --git a/bkflow/pipeline_plugins/components/collections/subprocess_plugin/v1_0_0.py b/bkflow/pipeline_plugins/components/collections/subprocess_plugin/v1_0_0.py index 9d2ee0bd79..0e8cf94d86 100644 --- a/bkflow/pipeline_plugins/components/collections/subprocess_plugin/v1_0_0.py +++ b/bkflow/pipeline_plugins/components/collections/subprocess_plugin/v1_0_0.py @@ -32,6 +32,7 @@ from bkflow.contrib.api.collections.interface import InterfaceModuleClient from bkflow.exceptions import ValidationError from bkflow.pipeline_plugins.components.collections.base import BKFlowBaseService +from bkflow.task.utils import push_task_to_queue class Subprocess(BaseModel): @@ -143,7 +144,7 @@ def _create_subprocess_task_instance(self, subprocess, template, pipeline_tree, from bkflow.task.utils import extract_extra_info with transaction.atomic(): - time_zone = timezone.pytz.timezone(settings.TIME_ZONE) or "Asia/Shanghai" + time_zone = timezone.pytz.timezone(settings.TIME_ZONE) time_stamp = datetime.datetime.now(tz=time_zone).strftime("%Y%m%d%H%M%S") create_task_data = { "name": f"{subprocess.subprocess_name}_子流程_{time_stamp}", @@ -190,7 +191,7 @@ def _create_subprocess_task_instance(self, subprocess, template, pipeline_tree, except TaskFlowRelation.DoesNotExist: root_task_id = parent_task.id - relate_info = {"node_id": self.id, "node_version": self.version} + relate_info = {"node_id": self.id, "node_version": self.version, "parent_task_id": parent_task.id} TaskFlowRelation.objects.create( task_id=task_instance.id, parent_task_id=parent_task.id, @@ -214,6 +215,7 @@ def _create_subprocess_task_instance(self, subprocess, template, pipeline_tree, def plugin_execute(self, data, parent_data): from bkflow.task.models import TaskInstance from bkflow.task.operations import TaskOperation + from bkflow.task.utils import task_concurrency_limit_reached parent_task_id = parent_data.get_one_of_inputs("task_id") try: @@ -235,6 +237,10 @@ def plugin_execute(self, data, parent_data): # 设置输出并启动任务 data.set_outputs("task_id", task_instance.id) + + if task_concurrency_limit_reached(task_instance.space_id, task_instance.template_id): + push_task_to_queue(task_instance, "start") + return True task_operation = TaskOperation(task_instance=task_instance, queue=settings.BKFLOW_MODULE.code) operation_method = getattr(task_operation, "start", None) if operation_method is None: diff --git a/bkflow/space/configs.py b/bkflow/space/configs.py index 83edcaf046..eb3e9ea1e6 100644 --- a/bkflow/space/configs.py +++ b/bkflow/space/configs.py @@ -438,6 +438,24 @@ def validate(cls, value: str): return True +# 流程并发控制 +class ConcurrencyControlConfig(BaseSpaceConfig): + name = "concurrency_control" + desc = _("流程并发控制") + default_value = 0 + LEAST_NUMBER = 1 + control = True + + @classmethod + def validate(cls, value: str): + if int(value) < cls.LEAST_NUMBER: + raise ValidationError( + f"[validate concurrency control error]: concurrency control only support {cls.LEAST_NUMBER}" + ) + + return True + + # 定义 SCHEMA_V1 对应的模型 class SchemaV1Model(BaseModel): meta_apis: str diff --git a/bkflow/task/celery/tasks.py b/bkflow/task/celery/tasks.py index 2dc2bb16e9..66759a69e2 100644 --- a/bkflow/task/celery/tasks.py +++ b/bkflow/task/celery/tasks.py @@ -40,7 +40,13 @@ from bkflow.task.node_timeout import node_timeout_handler from bkflow.task.operations import TaskNodeOperation, TaskOperation from bkflow.task.serializers import CreateTaskInstanceSerializer -from bkflow.task.utils import ATOM_FAILED, redis_inst_check, send_task_instance_message +from bkflow.task.utils import ( + ATOM_FAILED, + push_task_to_queue, + redis_inst_check, + send_task_instance_message, + task_concurrency_limit_reached, +) logger = logging.getLogger("celery") @@ -199,6 +205,10 @@ def bkflow_periodic_task_start(*args, **kwargs): } ) + if task_concurrency_limit_reached(task_instance.space_id, task_instance.template_id): + push_task_to_queue(task_instance, "start") + return + task_operation = TaskOperation(task_instance=task_instance, queue=settings.BKFLOW_MODULE.code) operation_method = getattr(task_operation, "start") if operation_method is None: diff --git a/bkflow/task/domains/callback.py b/bkflow/task/domains/callback.py index d4504cee28..ab6f46661a 100644 --- a/bkflow/task/domains/callback.py +++ b/bkflow/task/domains/callback.py @@ -27,16 +27,17 @@ from pipeline.eri.models import Schedule as DBSchedule from pipeline.eri.runtime import BambooDjangoRuntime -from bkflow.task.models import TaskFlowRelation +from bkflow.task.models import TaskInstance +from bkflow.task.utils import push_task_to_queue, task_concurrency_limit_reached from bkflow.utils.redis_lock import redis_lock logger = logging.getLogger("root") class TaskCallBacker: - def __init__(self, task_id, *args, **kwargs): + def __init__(self, task_id, task_relate, *args, **kwargs): self.task_id = task_id - self.task_relate = TaskFlowRelation.objects.filter(task_id=self.task_id).first() + self.task_relate = task_relate self.extra_info = {"task_id": self.task_id, **self.task_relate.extra_info, **kwargs} def check_record_existence(self): @@ -68,6 +69,13 @@ def subprocess_callback(self): runtime.set_state(node_id=node_id, version=version, to_state=states.READY) runtime.set_state(node_id=node_id, version=version, to_state=states.RUNNING) + parent_task_id = self.extra_info["parent_task_id"] + parent_task = TaskInstance.objects.filter(id=parent_task_id).first() + + if task_concurrency_limit_reached(parent_task.space_id, parent_task.template_id, is_exemption=True): + push_task_to_queue(parent_task, "callback") + return True + bamboo_engine_api.callback(runtime=runtime, node_id=node_id, version=version, data=self.extra_info) except Exception as e: diff --git a/bkflow/task/models.py b/bkflow/task/models.py index 61f4c41d17..506565fc71 100644 --- a/bkflow/task/models.py +++ b/bkflow/task/models.py @@ -330,8 +330,7 @@ def change_parent_task_node_state_to_running(self): runtime.set_execution_data_outputs(parent_node_id, data_outputs) # 仅当父流程的节点状态为失败时,才需要唤醒父流程的节点 - parent_task_id = TaskFlowRelation.objects.filter(task_id=self.id).first().parent_task_id - parent_task = TaskInstance.objects.get(id=parent_task_id) + parent_task = TaskInstance.objects.get(id=record.parent_task_id) parent_task.change_parent_task_node_state_to_running() diff --git a/bkflow/task/serializers.py b/bkflow/task/serializers.py index 24fe6bb1d5..4e6121f93f 100644 --- a/bkflow/task/serializers.py +++ b/bkflow/task/serializers.py @@ -118,10 +118,11 @@ class TaskInstanceSerializer(serializers.ModelSerializer): create_time = serializers.DateTimeField(format="%Y-%m-%d %H:%M:%S%z") start_time = serializers.DateTimeField(format="%Y-%m-%d %H:%M:%S%z") finish_time = serializers.DateTimeField(format="%Y-%m-%d %H:%M:%S%z") + is_waiting = serializers.SerializerMethodField() class Meta: model = TaskInstance - fields = "__all__" + exclude = ["extra_info"] read_only_fields = ( "id", "instance_id", @@ -141,6 +142,9 @@ class Meta: "tree_info_id", ) + def get_is_waiting(self, instance): + return instance.extra_info.get("is_waiting", False) + class RetrieveTaskInstanceSerializer(TaskInstanceSerializer): pipeline_tree = serializers.SerializerMethodField() diff --git a/bkflow/task/signals/handlers.py b/bkflow/task/signals/handlers.py index 54b6ccdb24..988d7f8561 100644 --- a/bkflow/task/signals/handlers.py +++ b/bkflow/task/signals/handlers.py @@ -33,7 +33,12 @@ TaskInstance, TimeoutNodeConfig, ) -from bkflow.task.utils import ATOM_FAILED, TASK_FINISHED, redis_inst_check +from bkflow.task.utils import ( + ATOM_FAILED, + TASK_FINISHED, + process_task_from_queue, + redis_inst_check, +) logger = logging.getLogger("root") @@ -99,12 +104,14 @@ def bamboo_engine_eri_post_set_state_handler(sender, node_id, to_state, version, queue=f"task_common_{settings.BKFLOW_MODULE.code}", routing_key=f"task_common_{settings.BKFLOW_MODULE.code}", ) + _process_task_from_queue(root_id) elif to_state == bamboo_engine_states.REVOKED and node_id == root_id: try: TaskInstance.objects.set_revoked(root_id) except Exception as e: logger.exception(f"TaskInstance set revoked error: {e}") _check_and_callback(root_id, task_success=False) + _process_task_from_queue(root_id) elif to_state == bamboo_engine_states.FINISHED and node_id == root_id: try: TaskInstance.objects.set_finished(root_id) @@ -119,6 +126,7 @@ def bamboo_engine_eri_post_set_state_handler(sender, node_id, to_state, version, routing_key=f"task_common_{settings.BKFLOW_MODULE.code}", ) _check_and_callback(root_id, task_success=True) + _process_task_from_queue(root_id) try: _node_timeout_info_update(settings.redis_inst, to_state, node_id, version) @@ -126,6 +134,21 @@ def bamboo_engine_eri_post_set_state_handler(sender, node_id, to_state, version, logger.exception(f"node_timeout_info_update error: {e}") +def _process_task_from_queue(root_id): + try: + template_id = TaskInstance.objects.get(instance_id=root_id).template_id + process_task_from_queue.apply_async( + kwargs={ + "template_id": template_id, + }, + queue=f"task_common_{settings.BKFLOW_MODULE.code}", + routing_key=f"task_common_{settings.BKFLOW_MODULE.code}", + ) + except Exception as e: + logger.exception(f"TaskInstance get template_id error: {e}") + return + + def _check_and_callback(instance_id, *args, **kwargs): try: task_id = TaskInstance.objects.get(instance_id=instance_id).id @@ -143,7 +166,7 @@ def task_callback(task_id, retry_times=0, *args, **kwargs): task_relate = TaskFlowRelation.objects.filter(task_id=task_id).first() if not task_relate: return - tcb = TaskCallBacker(task_id, *args, **kwargs) + tcb = TaskCallBacker(task_id, task_relate, *args, **kwargs) if not tcb.check_record_existence(): message = f"[task_callback] task_id {task_id} does not in TaskCallBackRecord." logger.error(message) diff --git a/bkflow/task/utils.py b/bkflow/task/utils.py index 8606b08343..66fdd8fc59 100644 --- a/bkflow/task/utils.py +++ b/bkflow/task/utils.py @@ -22,6 +22,7 @@ from functools import wraps from bamboo_engine import states as bamboo_engine_states +from celery import current_app from django.conf import settings from django.utils.translation import ugettext_lazy as _ from pipeline.core import constants as pipeline_constants @@ -30,6 +31,8 @@ from bkflow.utils.dates import format_datetime from bkflow.utils.message import send_message +from bkflow.utils.space import space_config_manager +from bkflow.utils.trace import start_trace logger = logging.getLogger("root") @@ -166,3 +169,111 @@ def extract_extra_info(constants, keys=None): for key in list(constants.keys()) if not keys else keys: extra_info.update({key: {"name": constants[key]["name"], "value": constants[key]["value"]}}) return json.dumps(extra_info, ensure_ascii=False) + + +@redis_inst_check +def push_task_to_queue(task, operation, node_id=None, data=None): + template_id = task.template_id + redis_key = f"task_wait_{template_id}" + + # 检查队列大小限制(最大1000个任务) + queue_size = settings.redis_inst.llen(redis_key) + + if queue_size >= settings.TASK_QUEUE_MAX_SIZE: + logger.error(f"Task queue for template {template_id} is full (size: {queue_size}), cannot add more tasks") + return False + task_data = {"operation": operation, "task_id": task.id} + if node_id: + task_data.update({"node_id": node_id}) + if data: + task_data.update({"node_data": data}) + + with start_trace("push_task_to_queue", queue_len=queue_size, **task_data): + task_json = json.dumps(task_data) + settings.redis_inst.rpush(redis_key, task_json) + + task.extra_info.update({"is_waiting": True}) + task.save() + return True + + +@current_app.task(bind=True) +@redis_inst_check +def process_task_from_queue(template_id): + from bkflow.task.models import TaskInstance + from bkflow.task.operations import TaskNodeOperation, TaskOperation + + redis_key = f"task_wait_{template_id}" + task_json = settings.redis_cli.lpop(redis_key) + if not task_json: + return None + + task_data = json.loads(task_json) + operation = task_data.get("operation") + task_instance = TaskInstance.objects.get(id=task_data.get("task_id")) + + for invoke_num in range(1, settings.TASK_QUEUE_MAX_SIZE + 1): + task_instance.extra_info.update({"is_waiting": False}) + try: + if operation in ["start", "resume"]: + task_operation = TaskOperation(task_instance, settings.BKFLOW_MODULE.code) + operation_method = getattr(task_operation, operation, None) + else: + node_operation = TaskNodeOperation(task_instance, task_data.get("node_id")) + operation_method = getattr(node_operation, operation, None) + + operation_result = operation_method(operator=operation, **task_data.get("node_data", {})) + if operation_result.result: + break + opera_error = operation_result.message + except Exception as e: + logger.error(f"Failed to process task {task_instance.id} from queue (attempt {invoke_num}): {e}") + opera_error = e + + if invoke_num == settings.TASK_QUEUE_MAX_SIZE and opera_error: + logger.error(f"Failed to process task {task_instance.id} kwargs {task_data} from queue") + task_instance.extra_info.update({"operation_failed": opera_error}) + + task_instance.save() + return task_instance + + +def get_running_task_count(space_id, template_id): + """统计当前正在执行的任务数量""" + from bkflow.task.models import TaskInstance + from bkflow.task.operations import TaskOperation + + task_instances = TaskInstance.objects.filter( + space_id=space_id, template_id=template_id, is_deleted=False, is_started=True, is_finished=False + ) + + task_operations = [ + {"task_id": task_instance.id, "operation": TaskOperation(task_instance=task_instance).get_task_states()} + for task_instance in task_instances + ] + + task_count = 0 + for task_operation in task_operations: + if task_operation["operation"].result is False: + continue + if task_operation["operation"].data.get("state") == "RUNNING": + task_count += 1 + + return task_count + + +def task_concurrency_limit_reached(space_id, template_id, is_exemption=False): + """判断是否超出并发限制""" + concurrency_control = space_config_manager.get_concurrency_control(space_id) + if not concurrency_control: + return False + + redis_key = f"task_wait_{template_id}" + queue_size = settings.redis_inst.llen(redis_key) + if queue_size != 0: + return True + + running_count = get_running_task_count(space_id, template_id) + if is_exemption: + return running_count > concurrency_control + return running_count >= concurrency_control diff --git a/bkflow/task/views.py b/bkflow/task/views.py index 42232af140..5d7ca915c2 100644 --- a/bkflow/task/views.py +++ b/bkflow/task/views.py @@ -68,6 +68,7 @@ TaskOperationRecordSerializer, UpdatePeriodicTaskSerializer, ) +from bkflow.task.utils import push_task_to_queue, task_concurrency_limit_reached from bkflow.utils.handlers import handle_plain_log from bkflow.utils.mixins import BKFLOWCommonMixin from bkflow.utils.permissions import AdminPermission, AppInternalPermission @@ -194,6 +195,11 @@ def operate(self, request, operation, *args, **kwargs): template_id=task_instance.template_id, executor=task_instance.executor, ): + if operation in ["start", "resume"] and task_concurrency_limit_reached( + task_instance.space_id, task_instance.template_id + ): + push_task_to_queue(task_instance, operation) + return Response({"result": True, "data": None, "message": "success"}) task_operation = TaskOperation(task_instance=task_instance, queue=settings.BKFLOW_MODULE.code) operation_method = getattr(task_operation, operation, None) if operation_method is None: @@ -222,12 +228,19 @@ def node_operate(self, request, node_id, operation, *args, **kwargs): ): if task_instance.trigger_method == TaskTriggerMethod.subprocess.name and operation in ["skip", "retry"]: task_instance.change_parent_task_node_state_to_running() + data = request.data + operator = data.pop("operator", request.user.username) + if operation in ["skip", "retry"] and task_concurrency_limit_reached( + task_instance.space_id, task_instance.template_id + ): + push_task_to_queue(task_instance, operation, node_id, data) + return Response({"result": True, "data": None, "message": "success"}) + node_operation = TaskNodeOperation(task_instance=task_instance, node_id=node_id) operation_method = getattr(node_operation, operation, None) if operation_method is None: raise ValidationError("node operation not found") - data = request.data - operator = data.pop("operator", request.user.username) + operation_result = operation_method(operator=operator, **data) return Response(dict(operation_result)) @@ -238,6 +251,8 @@ def get_states(self, request, *args, **kwargs): task_instance = self.get_object() task_operation = TaskOperation(task_instance=task_instance, queue=settings.BKFLOW_MODULE.code) states = task_operation.get_task_states() + if task_instance.extra_info.get("is_waiting"): + states.data["state"] = "waiting" return Response(dict(states)) @swagger_auto_schema(methods=["post"], operation_description="任务状态查询", request_body=GetTasksStatesBodySerializer) @@ -250,19 +265,19 @@ def get_tasks_states(self, request, *args, **kwargs): space_id = ser.validated_data["space_id"] task_instances = TaskInstance.objects.filter(id__in=task_ids, space_id=space_id) task_operations = [ - {"task_id": task_instance.id, "operation": TaskOperation(task_instance=task_instance).get_task_states()} + {"task": task_instance, "operation": TaskOperation(task_instance=task_instance).get_task_states()} for task_instance in task_instances ] - task_states = { - task_operation["task_id"]: { - "state": ( - task_operation["operation"].data.get("state") - if task_operation["operation"].result is True - else None - ) - } - for task_operation in task_operations - } + task_states = {} + for task_operation in task_operations: + if task_operation["operation"].result is True: + state = task_operation["operation"].data.get("state") + elif task_operation["task"].extra_info.get("is_waiting"): + state = "waiting" + else: + state = None + task_states[task_operation["task"].id] = {"state": state} + return Response(task_states) @swagger_auto_schema(methods=["get"], operation_description="获取任务 mock 数据") diff --git a/bkflow/utils/space.py b/bkflow/utils/space.py new file mode 100644 index 0000000000..801b8909a2 --- /dev/null +++ b/bkflow/utils/space.py @@ -0,0 +1,75 @@ +""" +TencentBlueKing is pleased to support the open source community by making +蓝鲸流程引擎服务 (BlueKing Flow Engine Service) available. +Copyright (C) 2024 THL A29 Limited, +a Tencent company. All rights reserved. +Licensed under the MIT License (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at http://opensource.org/licenses/MIT +Unless required by applicable law or agreed to in writing, +software distributed under the License is distributed on +an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, +either express or implied. See the License for the +specific language governing permissions and limitations under the License. + +We undertake not to change the open source license (MIT license) applicable + +to the current version of the project delivered to anyone in the future. +""" + + +import logging +import time +from typing import Dict + +from bkflow.contrib.api.collections.interface import InterfaceModuleClient +from bkflow.utils.singleton import Singleton + +logger = logging.getLogger("root") + + +class SpaceConfigManager(metaclass=Singleton): + """空间配置管理器""" + + def __init__(self): + self._cache: Dict[str, Dict] = {} + self._cache_time: Dict[str, float] = {} + self._cache_duration = 60 + self._interface_client = InterfaceModuleClient() + + def get_space_config(self, space_id: str, config_names: str) -> Dict: + """获取空间配置,支持缓存""" + cache_key = f"{space_id}:{config_names}" + + current_time = time.time() + if cache_key in self._cache: + cache_time = self._cache_time.get(cache_key, 0) + if current_time - cache_time < self._cache_duration: + return self._cache[cache_key] + + try: + space_infos_result = self._interface_client.get_space_infos( + {"space_id": space_id, "config_names": config_names} + ) + + if space_infos_result.get("result"): + space_configs = space_infos_result.get("data", {}).get("configs", {}) + + self._cache[cache_key] = space_configs + self._cache_time[cache_key] = current_time + + return space_configs + else: + logger.error(f"获取空间配置失败: space_id={space_id}, error={space_infos_result.get('message')}") + return {} + + except Exception as e: + logger.error(f"获取空间配置异常: space_id={space_id}, error={e}") + return {} + + def get_concurrency_control(self, space_id: str) -> int: + space_configs = self.get_space_config(space_id, "concurrency_control") + return int(space_configs.get("concurrency_control", 0)) + + +space_config_manager = SpaceConfigManager() diff --git a/config/default.py b/config/default.py index 7997f63684..90afeaa9ef 100644 --- a/config/default.py +++ b/config/default.py @@ -448,3 +448,5 @@ def handler_filter_injection(filters: list): TEMPLATE_MAX_RECURSIVE_NUMBER = env.TEMPLATE_MAX_RECURSIVE_NUMBER REQUEST_RETRY_NUMBER = env.REQUEST_RETRY_NUMBER + +TASK_QUEUE_MAX_SIZE = env.TASK_QUEUE_MAX_SIZE diff --git a/env.py b/env.py index fde9f5dea1..fabeec9e5d 100644 --- a/env.py +++ b/env.py @@ -181,3 +181,5 @@ TEMPLATE_MAX_RECURSIVE_NUMBER = int(os.getenv("TEMPLATE_MAX_RECURSIVE_NUMBER", 10)) REQUEST_RETRY_NUMBER = int(os.getenv("REQUEST_RETRY_NUMBER", 3)) + +TASK_QUEUE_MAX_SIZE = int(os.getenv("TASK_QUEUE_MAX_SIZE", 3)) From afeae91919596f45e96b0fee3ffc201906b97311 Mon Sep 17 00:00:00 2001 From: guohelu <19503896967@163.com> Date: Mon, 15 Dec 2025 20:14:21 +0800 Subject: [PATCH 2/4] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E9=80=BB=E8=BE=91?= =?UTF-8?q?=E5=A4=84=E7=90=86=E9=97=AE=E9=A2=98=20--story=3D128208486=20#?= =?UTF-8?q?=20Reviewed,=20transaction=20id:=2068150?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app_desc.yaml | 2 +- bkflow/task/domains/callback.py | 5 ++++- bkflow/task/signals/handlers.py | 4 ++-- bkflow/task/utils.py | 9 ++++----- config/default.py | 1 + env.py | 3 ++- 6 files changed, 14 insertions(+), 10 deletions(-) diff --git a/app_desc.yaml b/app_desc.yaml index 1f1ee0c02e..e7749efc24 100644 --- a/app_desc.yaml +++ b/app_desc.yaml @@ -127,7 +127,7 @@ modules: plan: 4C4G5R replicas: 5 cworker: - command: celery -A blueapps.core.celery worker -P threads -Q celery,pipeline_additional_task,pipeline_additional_task_priority,task_common_${BKFLOW_MODULE_CODE},task_callback_${BKFLOW_MODULE_CODE},node_auto_retry_${BKFLOW_MODULE_CODE},timeout_node_execute_${BKFLOW_MODULE_CODE},timeout_node_record_${BKFLOW_MODULE_CODE} -n common_worker@%h -c 10 -l info + command: celery -A blueapps.core.celery worker -P threads -Q celery,pipeline_additional_task,pipeline_additional_task_priority,task_common_${BKFLOW_MODULE_CODE},task_callback_${BKFLOW_MODULE_CODE},task_process_queue_${BKFLOW_MODULE_CODE},node_auto_retry_${BKFLOW_MODULE_CODE},timeout_node_execute_${BKFLOW_MODULE_CODE},timeout_node_record_${BKFLOW_MODULE_CODE} -n common_worker@%h -c 10 -l info plan: 4C1G5R replicas: 5 timeout: diff --git a/bkflow/task/domains/callback.py b/bkflow/task/domains/callback.py index ab6f46661a..9fcc6a6fb1 100644 --- a/bkflow/task/domains/callback.py +++ b/bkflow/task/domains/callback.py @@ -72,7 +72,10 @@ def subprocess_callback(self): parent_task_id = self.extra_info["parent_task_id"] parent_task = TaskInstance.objects.filter(id=parent_task_id).first() - if task_concurrency_limit_reached(parent_task.space_id, parent_task.template_id, is_exemption=True): + if ( + task_concurrency_limit_reached(parent_task.space_id, parent_task.template_id, is_exemption=True) + and self.extra_info["task_success"] is True + ): push_task_to_queue(parent_task, "callback") return True diff --git a/bkflow/task/signals/handlers.py b/bkflow/task/signals/handlers.py index 988d7f8561..1b22339072 100644 --- a/bkflow/task/signals/handlers.py +++ b/bkflow/task/signals/handlers.py @@ -141,8 +141,8 @@ def _process_task_from_queue(root_id): kwargs={ "template_id": template_id, }, - queue=f"task_common_{settings.BKFLOW_MODULE.code}", - routing_key=f"task_common_{settings.BKFLOW_MODULE.code}", + queue=f"task_process_queue_{settings.BKFLOW_MODULE.code}", + routing_key=f"task_process_queue_{settings.BKFLOW_MODULE.code}", ) except Exception as e: logger.exception(f"TaskInstance get template_id error: {e}") diff --git a/bkflow/task/utils.py b/bkflow/task/utils.py index 66fdd8fc59..cfe995fd74 100644 --- a/bkflow/task/utils.py +++ b/bkflow/task/utils.py @@ -176,7 +176,6 @@ def push_task_to_queue(task, operation, node_id=None, data=None): template_id = task.template_id redis_key = f"task_wait_{template_id}" - # 检查队列大小限制(最大1000个任务) queue_size = settings.redis_inst.llen(redis_key) if queue_size >= settings.TASK_QUEUE_MAX_SIZE: @@ -197,14 +196,14 @@ def push_task_to_queue(task, operation, node_id=None, data=None): return True -@current_app.task(bind=True) +@current_app.task() @redis_inst_check def process_task_from_queue(template_id): from bkflow.task.models import TaskInstance from bkflow.task.operations import TaskNodeOperation, TaskOperation redis_key = f"task_wait_{template_id}" - task_json = settings.redis_cli.lpop(redis_key) + task_json = settings.redis_inst.lpop(redis_key) if not task_json: return None @@ -212,7 +211,7 @@ def process_task_from_queue(template_id): operation = task_data.get("operation") task_instance = TaskInstance.objects.get(id=task_data.get("task_id")) - for invoke_num in range(1, settings.TASK_QUEUE_MAX_SIZE + 1): + for invoke_num in range(1, settings.TASK_MAX_RETRY_FREQUENCY + 1): task_instance.extra_info.update({"is_waiting": False}) try: if operation in ["start", "resume"]: @@ -230,7 +229,7 @@ def process_task_from_queue(template_id): logger.error(f"Failed to process task {task_instance.id} from queue (attempt {invoke_num}): {e}") opera_error = e - if invoke_num == settings.TASK_QUEUE_MAX_SIZE and opera_error: + if invoke_num == settings.TASK_MAX_RETRY_FREQUENCY and opera_error: logger.error(f"Failed to process task {task_instance.id} kwargs {task_data} from queue") task_instance.extra_info.update({"operation_failed": opera_error}) diff --git a/config/default.py b/config/default.py index 90afeaa9ef..4c409a6af5 100644 --- a/config/default.py +++ b/config/default.py @@ -450,3 +450,4 @@ def handler_filter_injection(filters: list): REQUEST_RETRY_NUMBER = env.REQUEST_RETRY_NUMBER TASK_QUEUE_MAX_SIZE = env.TASK_QUEUE_MAX_SIZE +TASK_MAX_RETRY_FREQUENCY = env.TASK_MAX_RETRY_FREQUENCY diff --git a/env.py b/env.py index fabeec9e5d..ddea4baeaa 100644 --- a/env.py +++ b/env.py @@ -182,4 +182,5 @@ TEMPLATE_MAX_RECURSIVE_NUMBER = int(os.getenv("TEMPLATE_MAX_RECURSIVE_NUMBER", 10)) REQUEST_RETRY_NUMBER = int(os.getenv("REQUEST_RETRY_NUMBER", 3)) -TASK_QUEUE_MAX_SIZE = int(os.getenv("TASK_QUEUE_MAX_SIZE", 3)) +TASK_QUEUE_MAX_SIZE = int(os.getenv("TASK_QUEUE_MAX_SIZE", 100)) +TASK_MAX_RETRY_FREQUENCY = int(os.getenv("TASK_MAX_RETRY_FREQUENCY", 3)) From 8f2374f2272f5dc94f3e999b74a58ba946092b3a Mon Sep 17 00:00:00 2001 From: guohelu <19503896967@163.com> Date: Mon, 15 Dec 2025 20:33:33 +0800 Subject: [PATCH 3/4] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8Dreview=E9=97=AE?= =?UTF-8?q?=E9=A2=98=20--story=3D128208486=20#=20Reviewed,=20transaction?= =?UTF-8?q?=20id:=2068153?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../collections/subprocess_plugin/v1_0_0.py | 6 +++- bkflow/task/celery/tasks.py | 5 ++- bkflow/task/domains/callback.py | 5 ++- bkflow/task/signals/handlers.py | 1 - bkflow/task/utils.py | 32 ++++++++++++++----- bkflow/task/views.py | 10 ++++-- bkflow/utils/space.py | 23 ++++++------- 7 files changed, 54 insertions(+), 28 deletions(-) diff --git a/bkflow/pipeline_plugins/components/collections/subprocess_plugin/v1_0_0.py b/bkflow/pipeline_plugins/components/collections/subprocess_plugin/v1_0_0.py index 0e8cf94d86..7d92bef98c 100644 --- a/bkflow/pipeline_plugins/components/collections/subprocess_plugin/v1_0_0.py +++ b/bkflow/pipeline_plugins/components/collections/subprocess_plugin/v1_0_0.py @@ -239,7 +239,11 @@ def plugin_execute(self, data, parent_data): data.set_outputs("task_id", task_instance.id) if task_concurrency_limit_reached(task_instance.space_id, task_instance.template_id): - push_task_to_queue(task_instance, "start") + try: + push_task_to_queue(task_instance, "start") + except Exception as e: + data.set_outputs("ex_data", str(e)) + return False return True task_operation = TaskOperation(task_instance=task_instance, queue=settings.BKFLOW_MODULE.code) operation_method = getattr(task_operation, "start", None) diff --git a/bkflow/task/celery/tasks.py b/bkflow/task/celery/tasks.py index 66759a69e2..d1bd7a0805 100644 --- a/bkflow/task/celery/tasks.py +++ b/bkflow/task/celery/tasks.py @@ -206,7 +206,10 @@ def bkflow_periodic_task_start(*args, **kwargs): ) if task_concurrency_limit_reached(task_instance.space_id, task_instance.template_id): - push_task_to_queue(task_instance, "start") + try: + push_task_to_queue(task_instance, "start") + except Exception as e: + logger.exception(f"[bkflow_periodic_task_start] push task to queue failed: {e}") return task_operation = TaskOperation(task_instance=task_instance, queue=settings.BKFLOW_MODULE.code) diff --git a/bkflow/task/domains/callback.py b/bkflow/task/domains/callback.py index 9fcc6a6fb1..fb72a75e49 100644 --- a/bkflow/task/domains/callback.py +++ b/bkflow/task/domains/callback.py @@ -76,7 +76,10 @@ def subprocess_callback(self): task_concurrency_limit_reached(parent_task.space_id, parent_task.template_id, is_exemption=True) and self.extra_info["task_success"] is True ): - push_task_to_queue(parent_task, "callback") + try: + push_task_to_queue(parent_task, "callback") + except Exception as e: + logger.exception(f"[TaskCallBacker _subprocess_callback] push task to queue error: {e}") return True bamboo_engine_api.callback(runtime=runtime, node_id=node_id, version=version, data=self.extra_info) diff --git a/bkflow/task/signals/handlers.py b/bkflow/task/signals/handlers.py index 1b22339072..e402efbec0 100644 --- a/bkflow/task/signals/handlers.py +++ b/bkflow/task/signals/handlers.py @@ -146,7 +146,6 @@ def _process_task_from_queue(root_id): ) except Exception as e: logger.exception(f"TaskInstance get template_id error: {e}") - return def _check_and_callback(instance_id, *args, **kwargs): diff --git a/bkflow/task/utils.py b/bkflow/task/utils.py index cfe995fd74..17abdbf461 100644 --- a/bkflow/task/utils.py +++ b/bkflow/task/utils.py @@ -176,20 +176,36 @@ def push_task_to_queue(task, operation, node_id=None, data=None): template_id = task.template_id redis_key = f"task_wait_{template_id}" - queue_size = settings.redis_inst.llen(redis_key) - - if queue_size >= settings.TASK_QUEUE_MAX_SIZE: - logger.error(f"Task queue for template {template_id} is full (size: {queue_size}), cannot add more tasks") - return False + # 准备任务数据 task_data = {"operation": operation, "task_id": task.id} if node_id: task_data.update({"node_id": node_id}) if data: task_data.update({"node_data": data}) + task_json = json.dumps(task_data) + + lua_script = """ + local queue_key = KEYS[1] + local max_size = tonumber(ARGV[1]) + local task_data = ARGV[2] + + local current_size = redis.call('llen', queue_key) + if current_size >= max_size then + return -1 -- 队列已满 + end + + redis.call('rpush', queue_key, task_data) + return current_size + 1 -- 返回新队列大小 + """ + + with start_trace("push_task_to_queue", operation=operation, task_id=task.id): + result = settings.redis_inst.eval(lua_script, 1, redis_key, settings.TASK_QUEUE_MAX_SIZE, task_json) + + if result == -1: + logger.error(f"Task queue for template {template_id} is full, cannot add more tasks") + raise Exception(f"Task queue for template {template_id} is full, cannot add more tasks") - with start_trace("push_task_to_queue", queue_len=queue_size, **task_data): - task_json = json.dumps(task_data) - settings.redis_inst.rpush(redis_key, task_json) + logger.info(f"Task {task.id} added to queue for template {template_id}, new queue size: {result}") task.extra_info.update({"is_waiting": True}) task.save() diff --git a/bkflow/task/views.py b/bkflow/task/views.py index 5d7ca915c2..10d18b48e7 100644 --- a/bkflow/task/views.py +++ b/bkflow/task/views.py @@ -198,7 +198,10 @@ def operate(self, request, operation, *args, **kwargs): if operation in ["start", "resume"] and task_concurrency_limit_reached( task_instance.space_id, task_instance.template_id ): - push_task_to_queue(task_instance, operation) + try: + push_task_to_queue(task_instance, operation) + except Exception as e: + return Response({"result": False, "data": None, "message": str(e)}) return Response({"result": True, "data": None, "message": "success"}) task_operation = TaskOperation(task_instance=task_instance, queue=settings.BKFLOW_MODULE.code) operation_method = getattr(task_operation, operation, None) @@ -233,7 +236,10 @@ def node_operate(self, request, node_id, operation, *args, **kwargs): if operation in ["skip", "retry"] and task_concurrency_limit_reached( task_instance.space_id, task_instance.template_id ): - push_task_to_queue(task_instance, operation, node_id, data) + try: + push_task_to_queue(task_instance, operation, node_id, data) + except Exception as e: + return Response({"result": False, "data": None, "message": str(e)}) return Response({"result": True, "data": None, "message": "success"}) node_operation = TaskNodeOperation(task_instance=task_instance, node_id=node_id) diff --git a/bkflow/utils/space.py b/bkflow/utils/space.py index 801b8909a2..1f8c6556eb 100644 --- a/bkflow/utils/space.py +++ b/bkflow/utils/space.py @@ -16,12 +16,12 @@ to the current version of the project delivered to anyone in the future. """ - - +import json import logging -import time from typing import Dict +from django.conf import settings + from bkflow.contrib.api.collections.interface import InterfaceModuleClient from bkflow.utils.singleton import Singleton @@ -32,20 +32,17 @@ class SpaceConfigManager(metaclass=Singleton): """空间配置管理器""" def __init__(self): - self._cache: Dict[str, Dict] = {} - self._cache_time: Dict[str, float] = {} self._cache_duration = 60 self._interface_client = InterfaceModuleClient() def get_space_config(self, space_id: str, config_names: str) -> Dict: """获取空间配置,支持缓存""" - cache_key = f"{space_id}:{config_names}" + # cache_key = f"{space_id}:{config_names}" + cache_key = f"space_config:{space_id}:{config_names}" - current_time = time.time() - if cache_key in self._cache: - cache_time = self._cache_time.get(cache_key, 0) - if current_time - cache_time < self._cache_duration: - return self._cache[cache_key] + cached_data = settings.redis_inst.get(cache_key) + if cached_data: + return json.loads(cached_data) try: space_infos_result = self._interface_client.get_space_infos( @@ -55,9 +52,7 @@ def get_space_config(self, space_id: str, config_names: str) -> Dict: if space_infos_result.get("result"): space_configs = space_infos_result.get("data", {}).get("configs", {}) - self._cache[cache_key] = space_configs - self._cache_time[cache_key] = current_time - + settings.redis_inst.setex(cache_key, self._cache_duration, json.dumps(space_configs)) return space_configs else: logger.error(f"获取空间配置失败: space_id={space_id}, error={space_infos_result.get('message')}") From d52c2f96cddf200f6d82940d82e110d067c8108b Mon Sep 17 00:00:00 2001 From: guohelu <19503896967@163.com> Date: Mon, 15 Dec 2025 20:38:04 +0800 Subject: [PATCH 4/4] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E5=8D=95=E4=BE=A7?= =?UTF-8?q?=E5=A4=B1=E8=B4=A5=E9=97=AE=E9=A2=98=20--story=3D128208486=20#?= =?UTF-8?q?=20Reviewed,=20transaction=20id:=2068155?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tests/interface/space/test_space_config.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/interface/space/test_space_config.py b/tests/interface/space/test_space_config.py index 70ba55c93b..91959388f1 100644 --- a/tests/interface/space/test_space_config.py +++ b/tests/interface/space/test_space_config.py @@ -39,9 +39,9 @@ class TestSpaceConfigHandler: def test_get_all_configs(self): configs = SpaceConfigHandler.get_all_configs() - assert len(configs) == 12 + assert len(configs) == 13 configs = SpaceConfigHandler.get_all_configs(only_public=True) - assert len(configs) == 11 + assert len(configs) == 12 def test_get_config(self): # valid cases