1. 项目概述:Controller/Worker模式在任务分解与调度中的应用
在分布式系统设计中,Controller/Worker架构是一种经典的任务调度模式,特别适合需要将复杂任务拆解为多个子任务并行执行的场景。这个模式的核心思想是将控制逻辑(Controller)与执行逻辑(Worker)分离,Controller负责任务分解、调度和结果汇总,Worker则专注于具体任务的执行。
最近我在一个结合SQL查询和文档RAG(检索增强生成)的项目中实践了这种模式。项目需求是要实现一个智能问答系统:用户输入自然语言问题后,系统需要先查询结构化数据库获取基础数据,同时检索相关文档片段作为上下文,最后综合这两类信息生成回答。这个过程中涉及多种不同类型的任务——SQL解析与执行、文档检索、文本生成等,非常适合采用Controller/Worker架构来实现。
2. 核心架构设计
2.1 Controller组件设计
Controller是整个系统的大脑,主要承担以下职责:
任务接收与解析:接收用户请求,解析出需要执行的任务类型和参数。在我们的案例中,用户问题会被分析是否需要SQL查询、文档检索或两者都需要。
任务分解:将复合任务拆解为原子性子任务。例如,一个包含数据查询和文档检索的请求会被分解为:
- SQL子任务:提取问题中的查询条件,生成SQL语句
- RAG子任务:提取问题中的关键词,准备文档检索
任务调度:根据子任务类型和当前系统负载,将任务分配给合适的Worker。我们维护了一个Worker能力注册表,记录每个Worker支持的任务类型和当前负载。
结果聚合:收集各Worker返回的结果,进行必要的后处理。比如将SQL查询结果和文档片段合并,作为大模型生成回答的上下文。
class Controller: def __init__(self): self.worker_pool = WorkerPool() self.task_queue = PriorityQueue() def handle_request(self, user_query): # 任务分解 subtasks = self.analyze_query(user_query) # 任务分发 futures = [] for task in subtasks: worker = self.worker_pool.get_available_worker(task.type) future = worker.execute(task) futures.append(future) # 等待结果 results = [f.result() for f in futures] # 结果聚合 final_result = self.aggregate_results(results) return final_result2.2 Worker组件设计
Worker是任务的具体执行者,我们设计了两种专用Worker:
SQL Worker:
- 接收包含SQL查询参数的任务
- 连接数据库执行查询
- 对结果进行初步处理(如格式化、过滤)
- 返回结构化数据
RAG Worker:
- 接收包含检索关键词的任务
- 从文档库中检索相关片段
- 对文档进行相关性排序
- 返回最相关的几个文档片段
每种Worker都实现了统一的接口,便于Controller调度:
class SQLWorker: def __init__(self, db_connection): self.db = db_connection def execute(self, task): try: # 执行SQL查询 cursor = self.db.cursor() cursor.execute(task.sql) results = cursor.fetchall() # 简单的结果处理 processed = self.process_results(results) return processed except Exception as e: raise WorkerException(f"SQL执行失败: {str(e)}") class RAGWorker: def __init__(self, vector_db): self.vector_db = vector_db def execute(self, task): # 向量相似度检索 docs = self.vector_db.search(task.keywords, top_k=3) # 相关性过滤 filtered = [doc for doc in docs if doc.score > 0.7] return filtered3. 任务调度实现细节
3.1 任务队列管理
我们采用优先级队列来管理待处理任务,考虑以下因素确定优先级:
- 任务类型(SQL查询通常比文档检索更紧急)
- 任务预估耗时(短任务优先)
- 用户等级(VIP用户任务优先级更高)
class Task: def __init__(self, task_type, params, priority=0): self.type = task_type # 'sql' or 'rag' self.params = params self.priority = priority self.create_time = time.time() def __lt__(self, other): # 优先级比较逻辑 if self.priority != other.priority: return self.priority > other.priority return self.create_time < other.create_time3.2 Worker负载均衡
为了避免某些Worker过载而其他Worker闲置,我们实现了动态负载均衡:
- 每个Worker定期向Controller报告其负载状态(CPU、内存使用率,待处理任务数)
- Controller根据负载情况分配新任务
- 如果所有Worker都处于高负载状态,Controller可以:
- 拒绝新请求(返回"系统繁忙")
- 将任务排队等待
- 动态创建新的Worker实例(如果支持弹性伸缩)
class WorkerPool: def __init__(self): self.workers = { 'sql': [], 'rag': [] } self.load_metrics = {} # worker_id -> {cpu, memory, queue_size} def get_available_worker(self, task_type): # 找出同类型Worker中负载最低的一个 candidates = self.workers[task_type] if not candidates: raise NoAvailableWorkerError() # 选择负载最低的Worker return min(candidates, key=lambda w: self.calculate_load(w.id)) def calculate_load(self, worker_id): metrics = self.load_metrics.get(worker_id, {}) cpu = metrics.get('cpu', 0) memory = metrics.get('memory', 0) queue = metrics.get('queue_size', 0) # 综合负载计算公式 return 0.6*cpu + 0.3*memory + 0.1*queue4. SQL与RAG的协同处理
4.1 任务依赖处理
有些任务之间存在依赖关系,比如需要先获取SQL查询结果,然后基于这些结果进行文档检索。我们通过任务图(DAG)来表示这种依赖:
- Controller解析出任务间的依赖关系
- 无依赖的任务可以并行执行
- 有依赖的任务按顺序执行
def build_task_graph(user_query): # 分析查询,构建任务依赖图 tasks = [] dependencies = {} # 示例:先执行SQL查询,然后用查询结果作为文档检索条件 sql_task = Task('sql', extract_sql_params(user_query)) rag_task = Task('rag', extract_rag_params(user_query)) tasks.append(sql_task) tasks.append(rag_task) dependencies[rag_task] = [sql_task] # RAG任务依赖SQL任务 return tasks, dependencies4.2 结果合并策略
当SQL结果和文档片段都获取后,需要将它们合并作为大模型生成回答的上下文。我们设计了多种合并策略:
- 简单拼接:将SQL结果转为文本,与文档片段直接拼接
- 结构化合并:保持SQL结果的结构化特性,文档作为补充说明
- 智能过滤:使用小型模型评估各信息的相关性,过滤掉不相关部分
def aggregate_results(results): sql_result = next(r for r in results if r.type == 'sql') rag_result = next(r for r in results if r.type == 'rag') # 使用结构化合并策略 context = { "data": sql_result.data, "documents": rag_result.documents, "metadata": { "sql_query": sql_result.metadata, "rag_keywords": rag_result.metadata } } return context5. 性能优化实践
5.1 Worker预热
为了避免冷启动问题,我们实现了Worker预热机制:
- 系统启动时预创建一定数量的Worker
- 定期保持最低数量的空闲Worker
- 对数据库连接、模型加载等耗时操作提前完成
def pre_start_workers(): # 预启动SQL Worker for _ in range(MIN_SQL_WORKERS): worker = SQLWorker(create_db_connection()) worker_pool.register(worker) # 预启动RAG Worker for _ in range(MIN_RAG_WORKERS): worker = RAGWorker(load_vector_db()) worker_pool.register(worker)5.2 查询缓存
对于频繁出现的相似查询,我们实现了两级缓存:
- SQL结果缓存:对参数化查询的结果缓存5分钟
- 文档检索缓存:对相同关键词的检索结果缓存10分钟
缓存键考虑了查询参数和用户上下文,确保不同用户看到的缓存结果可能不同。
class ResultCache: def __init__(self): self.sql_cache = LRUCache(maxsize=1000, ttl=300) self.rag_cache = LRUCache(maxsize=5000, ttl=600) def get_sql_result(self, query, params): cache_key = self.make_sql_key(query, params) return self.sql_cache.get(cache_key) def get_rag_result(self, keywords, user_context): cache_key = self.make_rag_key(keywords, user_context) return self.rag_cache.get(cache_key)6. 错误处理与容灾
6.1 Worker故障处理
Worker可能因各种原因失败,我们实现了以下容错机制:
- 心跳检测:Worker定期发送心跳,超时视为故障
- 任务重试:失败的任务会重试最多2次
- Worker替换:故障Worker会被新Worker替代
def handle_worker_failure(worker_id): # 从池中移除故障Worker worker_pool.remove(worker_id) # 重新调度该Worker未完成的任务 for task in task_queue.get_pending_tasks(worker_id): task_queue.put(task) # 重新入队 # 启动新Worker替代 if worker_pool.worker_type(worker_id) == 'sql': new_worker = SQLWorker(create_db_connection()) else: new_worker = RAGWorker(load_vector_db()) worker_pool.register(new_worker)6.2 降级策略
当部分功能不可用时,系统可以降级运行:
- SQL不可用:仅使用文档检索结果
- RAG不可用:仅使用SQL查询结果
- 两者都不可用:返回缓存的通用回答或错误信息
def fallback_strategy(available_components): if 'sql' not in available_components and 'rag' not in available_components: return {"error": "系统维护中,请稍后再试"} if 'sql' not in available_components: return {"warning": "数据查询不可用,仅提供文档信息"} if 'rag' not in available_components: return {"warning": "文档检索不可用,仅提供数据结果"}7. 监控与日志
完善的监控对系统运维至关重要,我们实现了:
性能指标收集:
- 任务处理耗时
- Worker负载
- 队列长度
业务指标收集:
- SQL查询成功率
- 文档检索召回率
- 结果生成质量
分布式追踪:
- 每个请求的唯一ID
- 跨组件的调用链
- 各阶段耗时分析
class Monitor: def record_metrics(self, metric_type, values): # 记录到时间序列数据库 tsdb.write(metric_type, values, timestamp=time.time()) def log_trace(self, trace_id, component, event, duration=None): # 记录分布式追踪信息 logger.info(f"{trace_id} | {component} | {event} | {duration}") def alert(self, condition): # 触发告警 if condition: alert_system.notify()8. 实际部署考量
8.1 资源分配
根据我们的经验,不同Worker类型的资源需求差异很大:
SQL Worker:
- 需要较多数据库连接
- CPU密集型
- 内存需求中等
RAG Worker:
- 需要GPU加速(如果使用神经网络模型)
- 内存需求高(加载向量索引)
- 磁盘I/O密集(文档检索)
建议部署时:
- SQL Worker和RAG Worker分开部署
- 根据监控数据动态调整Worker比例
- 为RAG Worker配置更强大的硬件
8.2 安全考虑
SQL注入防护:
- 使用参数化查询
- 限制查询复杂度
- 设置查询超时
文档访问控制:
- 实现基于角色的文档访问
- 检索结果根据用户权限过滤
- 敏感文档特殊处理
def safe_execute_sql(worker, query, params): # 检查查询复杂度 if estimate_query_complexity(query) > MAX_COMPLEXITY: raise QueryTooComplexError() # 设置执行超时 try: with timeout(QUERY_TIMEOUT): return worker.execute(query, params) except TimeoutError: worker.cancel_query() raise QueryTimeoutError()9. 扩展性与演进
9.1 支持新Worker类型
系统设计支持轻松添加新Worker类型:
- 定义新Worker类,实现统一接口
- 在Controller中注册新类型
- 更新任务分发逻辑
例如,要添加图像处理Worker:
class ImageWorker: def execute(self, task): # 实现图像处理逻辑 pass # 注册新类型 worker_pool.register_worker_type('image', ImageWorker)9.2 横向扩展
当负载增加时,可以通过以下方式扩展:
- 增加Worker数量:启动更多Worker实例
- 分区处理:根据任务特征分区,不同分区由不同Worker组处理
- 层级化调度:引入多个Controller,形成树状调度结构
10. 经验总结与避坑指南
在实际实现过程中,我们积累了一些宝贵经验:
任务粒度要适中:
- 太细会导致调度开销过大
- 太粗会降低并行度和资源利用率
- 建议:每个任务执行时间在100ms-5s之间
避免Worker状态共享:
- Worker之间尽量不要共享状态
- 必须共享的状态通过中央存储管理
- 这样可以简化容错和扩展
合理设置超时:
- 任务执行超时
- Worker心跳超时
- 数据库查询超时
- 建议:根据实际场景测试确定最佳值
完善的监控必不可少:
- 系统级监控(CPU、内存等)
- 业务级监控(成功率、延迟等)
- 及时告警能避免小问题演变成故障
文档和示例很重要:
- 为每个Worker类型提供详细文档
- 包含常见使用示例
- 记录已知问题和解决方案
这个Controller/Worker模式的实现已经在我们生产环境稳定运行了6个月,日均处理超过50万次任务,平均延迟控制在200ms以内。特别是在处理混合了SQL查询和文档检索的复杂请求时,相比传统串行处理方式,性能提升了3-5倍。