1. 工作流技能的定义与核心价值
工作流中的Skill(技能)是指在工作流程中执行特定任务或功能的最小可复用单元。它相当于一个封装好的功能模块,可以在不同工作流中被重复调用。举个例子,就像乐高积木中的基础模块,单独使用时能完成特定功能,组合起来又能构建复杂系统。
在实际应用中,工作流Skill通常表现为以下几种形式:
- 自动化脚本(如Python函数、Shell命令)
- API接口调用封装
- 可视化工具中的预制节点
- 低代码平台中的功能组件
这类技能的核心价值在于:
- 复用性:一次开发多次使用,避免重复造轮子
- 解耦性:各技能模块相互独立,修改不影响整体流程
- 组合性:通过技能的自由组合实现复杂业务逻辑
- 可维护性:问题定位和修复可以精确到具体技能单元
提示:好的工作流Skill应该像Unix哲学倡导的那样——"只做一件事,并做到极致"
2. 工作流Skill的设计原则
2.1 单一职责原则
每个Skill应该只解决一个具体问题。比如:
- 一个专门处理日期格式转换的Skill
- 一个专门调用某API获取数据的Skill
- 一个专门进行数据校验的Skill
反例是将数据获取、转换、校验全部写在一个Skill里。这样会导致:
- 复用性降低(其他流程可能只需要其中部分功能)
- 问题排查困难(出错时难以定位具体环节)
- 版本管理复杂(修改一个功能需要整体发布)
2.2 明确输入输出
定义清晰的接口规范:
python复制# 好的示例
def calculate_tax(amount: float, tax_rate: float) -> float:
"""计算税费"""
return amount * tax_rate
# 差的示例
def process_data(data):
"""模糊的输入输出"""
# 各种混合处理逻辑
return result
建议使用类型注解(如Python的typing)或接口描述文档明确:
- 输入参数的名称、类型、取值范围
- 返回值的结构、可能的状态码
- 可能抛出的异常类型
2.3 无状态设计
Skill应该尽量避免依赖或修改外部状态。理想情况下:
- 输出只由输入决定
- 不依赖全局变量
- 不修改外部数据存储
这样能确保:
- 相同输入永远得到相同输出(幂等性)
- 可以安全地并行调用
- 更容易进行单元测试
3. 工作流Skill的实现方式
3.1 脚本类Skill实现
以Python为例,一个规范的邮件发送Skill:
python复制import smtplib
from email.mime.text import MIMEText
from typing import Tuple
def send_email(
sender: str,
receivers: list[str],
subject: str,
content: str,
smtp_config: dict
) -> Tuple[bool, str]:
"""
发送邮件Skill
参数:
sender: 发件人邮箱
receivers: 收件人邮箱列表
subject: 邮件主题
content: 邮件正文
smtp_config: SMTP配置字典,包含:
- host: SMTP服务器地址
- port: 端口号
- username: 用户名
- password: 密码
- use_tls: 是否启用TLS
返回:
(成功状态, 错误信息)
"""
try:
msg = MIMEText(content)
msg['Subject'] = subject
msg['From'] = sender
msg['To'] = ', '.join(receivers)
with smtplib.SMTP(smtp_config['host'], smtp_config['port']) as server:
if smtp_config['use_tls']:
server.starttls()
server.login(smtp_config['username'], smtp_config['password'])
server.sendmail(sender, receivers, msg.as_string())
return True, ""
except Exception as e:
return False, str(e)
关键实现要点:
- 使用类型注解明确接口
- 详细的文档字符串说明
- 使用上下文管理器确保资源释放
- 统一的错误处理返回格式
3.2 API类Skill封装
对于调用第三方API的场景,建议封装为:
python复制import requests
from datetime import datetime
from typing import Optional
class WeatherAPI:
def __init__(self, api_key: str):
self.base_url = "https://api.weather.com/v3"
self.api_key = api_key
def get_current_weather(
self,
location: str,
units: str = "metric"
) -> Optional[dict]:
"""
获取当前天气情况
参数:
location: 城市名称或经纬度(如"39.9042,116.4074")
units: 单位制(metric/imperial)
返回:
天气数据字典或None(失败时)
"""
try:
params = {
"location": location,
"units": units,
"apiKey": self.api_key
}
response = requests.get(
f"{self.base_url}/weather/current",
params=params,
timeout=5
)
response.raise_for_status()
return response.json()
except requests.exceptions.RequestException:
return None
封装时的注意事项:
- 将API密钥等敏感信息通过构造方法注入
- 设置合理的请求超时时间
- 处理各种网络异常情况
- 对返回数据不做过多处理,保持原始结构
3.3 可视化工具中的Skill配置
以常见的低代码平台为例,配置一个文件处理Skill的典型参数:
| 参数项 | 类型 | 必填 | 说明 | 示例值 |
|---|---|---|---|---|
| 操作类型 | 下拉选择 | 是 | 文件操作类型 | "移动文件" |
| 源文件路径 | 字符串 | 是 | 支持通配符 | "/data/input/*.csv" |
| 目标路径 | 字符串 | 是 | 目标目录 | "/data/archive/" |
| 冲突处理 | 单选按钮 | 否 | 文件存在时的处理方式 | "跳过" |
| 超时时间 | 数字 | 否 | 操作超时(秒) | 30 |
可视化Skill的设计要点:
- 每个配置项都要有明确的说明和示例
- 提供合理的默认值
- 对关键参数进行校验(如路径是否存在)
- 在界面上展示Skill的执行日志
4. 工作流Skill的测试与调试
4.1 单元测试实现
为上述邮件发送Skill编写测试用例:
python复制import pytest
from unittest.mock import patch
from email_skill import send_email
@pytest.fixture
def mock_smtp():
with patch('smtplib.SMTP') as mock:
yield mock
def test_send_email_success(mock_smtp):
# 配置mock
mock_instance = mock_smtp.return_value.__enter__.return_value
mock_instance.login.return_value = None
mock_instance.sendmail.return_value = None
# 测试输入
smtp_config = {
"host": "smtp.example.com",
"port": 587,
"username": "user",
"password": "pass",
"use_tls": True
}
# 调用被测Skill
success, error = send_email(
sender="from@example.com",
receivers=["to@example.com"],
subject="Test",
content="Hello",
smtp_config=smtp_config
)
# 验证结果
assert success is True
assert error == ""
mock_instance.starttls.assert_called_once()
mock_instance.login.assert_called_once_with("user", "pass")
def test_send_email_failure(mock_smtp):
mock_instance = mock_smtp.return_value.__enter__.return_value
mock_instance.login.side_effect = Exception("Auth failed")
success, error = send_email(
sender="from@example.com",
receivers=["to@example.com"],
subject="Test",
content="Hello",
smtp_config={}
)
assert success is False
assert "Auth failed" in error
测试要点:
- 使用mock替代真实SMTP服务
- 测试正常流程和异常流程
- 验证关键方法是否被正确调用
- 检查错误信息的准确性
4.2 集成测试方案
在本地搭建测试工作流:
bash复制# 使用Docker模拟测试环境
docker run -d --name workflow-test \
-v $(pwd)/skills:/skills \
-v $(pwd)/test-data:/data \
python:3.9-slim
# 安装依赖
docker exec workflow-test pip install pytest requests
# 执行测试
docker exec -w /skills workflow-test pytest -v
集成测试的关键检查项:
- Skill在工作流引擎中的加载情况
- 参数传递是否正确
- 与其他Skill的协同工作
- 错误处理和工作流中断情况
- 性能基准测试(执行时间、资源占用)
4.3 调试技巧与工具
常用调试方法:
- 日志记录:
python复制import logging
logger = logging.getLogger(__name__)
def example_skill(input):
logger.debug("收到输入: %s", input)
try:
result = process(input)
logger.info("处理成功,结果: %s", result)
return result
except Exception as e:
logger.error("处理失败: %s", str(e), exc_info=True)
raise
- 交互式调试:
- 使用Python的pdb或ipdb设置断点
- 在Jupyter Notebook中逐步测试Skill
- 使用Postman测试API类Skill
- 可视化追踪:
mermaid复制graph TD
A[开始] --> B[Skill1]
B --> C{条件判断}
C -->|是| D[Skill2]
C -->|否| E[Skill3]
D --> F[结束]
E --> F
注意:实际开发中应该避免将敏感信息(如密码、API密钥)记录到日志中
5. 工作流Skill的版本管理与发布
5.1 版本控制策略
推荐使用语义化版本控制(SemVer):
code复制MAJOR.MINOR.PATCH
- MAJOR:不兼容的接口修改
- MINOR:向下兼容的功能新增
- PATCH:向下兼容的问题修正
示例版本迭代:
code复制1.0.0 - 初始版本
1.0.1 - 修复时区处理bug
1.1.0 - 新增多语言支持
2.0.0 - 重构接口,移除旧参数
5.2 发布流程
标准化的发布检查清单:
- [ ] 更新版本号(version)
- [ ] 更新CHANGELOG.md
- [ ] 通过所有单元测试
- [ ] 通过集成测试
- [ ] 更新文档(参数说明、示例)
- [ ] 打Git标签(git tag v1.0.0)
- [ ] 推送到仓库(git push origin v1.0.0)
- [ ] 发布到内部Skill仓库
5.3 依赖管理
使用requirements.txt或pyproject.toml明确定义依赖:
toml复制# pyproject.toml示例
[project]
dependencies = [
"requests>=2.28.0",
"pydantic>=1.10.0",
]
[project.optional-dependencies]
test = [
"pytest>=7.0.0",
"pytest-mock>=3.0.0"
]
关键原则:
- 指定主要依赖的最低兼容版本
- 区分核心依赖和测试依赖
- 避免使用过于宽泛的版本范围
- 定期更新依赖版本(安全补丁)
6. 高级技巧与最佳实践
6.1 性能优化策略
- 缓存常用结果:
python复制from functools import lru_cache
@lru_cache(maxsize=128)
def get_config(key: str) -> dict:
"""带缓存的配置读取"""
return query_database(key) # 耗时的数据库查询
- 批量处理替代循环调用:
python复制# 差的做法
for user_id in user_ids:
send_notification(user_id)
# 好的做法
def batch_send_notifications(user_ids: list):
"""批量发送通知"""
# 使用更高效的批量API
- 异步处理:
python复制import asyncio
async def async_fetch_data(urls: list[str]) -> list:
"""异步获取多个URL数据"""
async with aiohttp.ClientSession() as session:
tasks = [fetch_url(session, url) for url in urls]
return await asyncio.gather(*tasks)
6.2 安全注意事项
- 参数校验:
python复制from pydantic import BaseModel, HttpUrl
class APIParams(BaseModel):
url: HttpUrl
retries: int = Field(ge=0, le=5)
timeout: float = Field(gt=0.0)
- 敏感信息处理:
python复制import os
from dotenv import load_dotenv
load_dotenv()
# 从环境变量读取配置
api_key = os.getenv("API_KEY")
db_url = os.getenv("DB_URL")
- 权限控制:
python复制def require_role(role: str):
"""装饰器检查用户角色"""
def decorator(func):
@functools.wraps(func)
def wrapper(*args, **kwargs):
if current_user.role != role:
raise PermissionError("无权访问")
return func(*args, **kwargs)
return wrapper
return decorator
6.3 监控与指标
使用Prometheus客户端收集指标:
python复制from prometheus_client import Counter, Histogram
REQUEST_COUNT = Counter(
'skill_requests_total',
'Total request count',
['skill_name', 'status']
)
REQUEST_TIME = Histogram(
'skill_request_duration_seconds',
'Request processing time',
['skill_name']
)
def track_metrics(func):
"""装饰器记录执行指标"""
@functools.wraps(func)
def wrapper(*args, **kwargs):
start_time = time.time()
try:
result = func(*args, **kwargs)
REQUEST_COUNT.labels(
skill_name=func.__name__,
status='success'
).inc()
return result
except Exception:
REQUEST_COUNT.labels(
skill_name=func.__name__,
status='failure'
).inc()
raise
finally:
REQUEST_TIME.labels(
skill_name=func.__name__
).observe(time.time() - start_time)
return wrapper
关键监控指标应该包括:
- 执行次数(按成功/失败分类)
- 执行耗时分布
- 资源使用情况(CPU、内存)
- 队列等待时间(如有)
7. 实际案例解析
7.1 电商订单处理Skill
一个完整的订单处理Skill实现:
python复制from typing import Literal
from datetime import datetime
from pydantic import BaseModel, validator
class Order(BaseModel):
order_id: str
user_id: str
items: list[dict]
status: Literal["pending", "paid", "shipped", "completed"]
created_at: datetime
@validator('order_id')
def validate_order_id(cls, v):
if not v.startswith('ORD-'):
raise ValueError("订单ID格式不正确")
return v
def process_order(order_data: dict) -> dict:
"""
处理电商订单
参数:
order_data: 原始订单数据字典
返回:
处理后的订单信息,包含:
- order_id: 订单编号
- status: 处理状态
- shipping_info: 物流信息(如已发货)
- payment_info: 支付信息
"""
# 1. 数据校验
try:
order = Order(**order_data)
except ValueError as e:
return {
"status": "invalid",
"error": str(e)
}
# 2. 支付处理
payment_result = process_payment(order)
if not payment_result["success"]:
return {
"order_id": order.order_id,
"status": "payment_failed",
"error": payment_result["message"]
}
# 3. 库存检查
inventory_check = check_inventory(order.items)
if not inventory_check["all_available"]:
return {
"order_id": order.order_id,
"status": "inventory_shortage",
"unavailable_items": inventory_check["unavailable"]
}
# 4. 物流处理
shipping_info = create_shipping(order)
return {
"order_id": order.order_id,
"status": "completed",
"shipping_info": shipping_info,
"payment_info": payment_result
}
这个案例展示了:
- 使用Pydantic进行严格输入校验
- 清晰的业务流程分步处理
- 全面的错误处理
- 结构化的返回结果
7.2 数据分析流水线Skill
一个典型的数据处理Skill链:
python复制def data_processing_pipeline(source_path: str) -> dict:
"""
完整的数据处理流水线
参数:
source_path: 原始数据文件路径
返回:
处理结果统计信息
"""
# 1. 数据抽取
raw_data = extract_data(source_path)
# 2. 数据清洗
cleaned_data = clean_data(raw_data)
# 3. 数据转换
transformed_data = transform_data(cleaned_data)
# 4. 数据分析
analysis_result = analyze_data(transformed_data)
# 5. 结果存储
save_result(analysis_result)
return {
"input_records": len(raw_data),
"output_records": len(analysis_result),
"processing_time": time.time() - start_time
}
每个子Skill都可以独立开发和测试:
python复制def extract_data(path: str) -> list:
"""数据抽取Skill"""
# 根据文件类型调用相应解析器
if path.endswith('.csv'):
return parse_csv(path)
elif path.endswith('.json'):
return parse_json(path)
else:
raise ValueError("不支持的格式")
def clean_data(data: list) -> list:
"""数据清洗Skill"""
# 处理缺失值
# 去除重复项
# 标准化格式
return processed_data
7.3 跨平台Skill适配
设计可跨工作流引擎使用的Skill:
python复制class BaseSkill:
"""Skill抽象基类"""
def __init__(self, config: dict):
self.config = config
self.validate_config()
def validate_config(self):
"""验证配置"""
required_keys = self.get_required_config()
for key in required_keys:
if key not in self.config:
raise ValueError(f"缺少必要配置: {key}")
@classmethod
def get_required_config(cls) -> list:
"""返回必要配置项"""
return []
def execute(self, input_data: dict) -> dict:
"""执行Skill"""
raise NotImplementedError
class FileMoveSkill(BaseSkill):
"""文件移动Skill"""
@classmethod
def get_required_config(cls) -> list:
return ["source_path", "target_path"]
def execute(self, input_data: dict) -> dict:
import shutil
source = self.config["source_path"]
target = self.config["target_path"]
try:
shutil.move(source, target)
return {
"status": "success",
"moved_files": [source]
}
except Exception as e:
return {
"status": "failed",
"error": str(e)
}
这种设计允许:
- 统一接口跨平台使用
- 集中的配置验证
- 标准化的执行结果格式
- 易于扩展新的Skill类型
8. 常见问题解决方案
8.1 技能执行超时问题
问题现象:
- Skill执行时间超过工作流引擎设置的最大超时
- 导致整个工作流中断
解决方案:
- 优化Skill性能:
python复制# 使用更高效的算法
def process_large_data(data):
# 差的做法:逐行处理
# result = [transform(line) for line in data]
# 好的做法:批量处理
batch_size = 1000
batches = [data[i:i+batch_size] for i in range(0, len(data), batch_size)]
result = []
for batch in batches:
result.extend(batch_transform(batch))
return result
- 设置合理的超时时间:
python复制import signal
def run_with_timeout(func, args=(), kwargs={}, timeout=30):
"""带超时控制的执行"""
def handler(signum, frame):
raise TimeoutError("执行超时")
signal.signal(signal.SIGALRM, handler)
signal.alarm(timeout)
try:
result = func(*args, **kwargs)
finally:
signal.alarm(0)
return result
- 拆分大任务:
python复制def process_in_chunks(data, chunk_size=1000):
"""分块处理大数据集"""
for i in range(0, len(data), chunk_size):
chunk = data[i:i+chunk_size]
yield process_chunk(chunk)
# 在工作流中逐个处理分块结果
8.2 技能版本兼容性问题
问题场景:
- 工作流A使用Skill v1.0
- 工作流B使用Skill v2.0
- 两个版本接口不兼容
解决方案:
- 语义化版本控制:
code复制技能仓库结构:
/skills
/email_sender
/v1
__init__.py
/v2
__init__.py
- 适配器模式:
python复制class SkillAdapter:
def __init__(self, skill_version):
if skill_version == "v1":
from .v1 import Skill
elif skill_version == "v2":
from .v2 import Skill
else:
raise ValueError("不支持的版本")
self.skill = Skill()
def execute(self, input):
# 统一不同版本的接口差异
return self.skill.run(input)
- 自动化迁移工具:
python复制def migrate_workflow(workflow_def, from_version, to_version):
"""自动迁移工作流定义"""
# 1. 解析工作流中的技能引用
# 2. 查找对应的版本映射
# 3. 更新参数格式
# 4. 返回迁移后的定义
return updated_workflow
8.3 技能调试困难问题
常见痛点:
- 复杂工作流中难以定位问题Skill
- 缺少执行上下文信息
- 测试环境与生产环境差异
调试工具箱:
- 上下文日志:
python复制import logging
from contextlib import contextmanager
@contextmanager
def log_context(context_info):
"""添加上下文信息的日志"""
logger = logging.getLogger(__name__)
logger.info("开始执行: %s", context_info)
try:
yield
except Exception as e:
logger.error("执行失败: %s - %s", context_info, str(e))
raise
finally:
logger.info("执行完成: %s", context_info)
# 使用示例
with log_context({"skill": "email_sender", "flow": "order_confirm"}):
send_email(...)
- 本地模拟器:
python复制def local_simulator(skill_class, test_cases):
"""本地测试技能"""
for case in test_cases:
print(f"测试用例: {case['name']}")
skill = skill_class(case["config"])
try:
result = skill.execute(case["input"])
assert case["expected"] == result
print("✓ 测试通过")
except Exception as e:
print(f"✗ 测试失败: {str(e)}")
- 差异对比工具:
python复制from deepdiff import DeepDiff
def compare_results(actual, expected):
"""详细比较结果差异"""
diff = DeepDiff(actual, expected, ignore_order=True)
if not diff:
print("结果完全匹配")
else:
print("发现差异:")
for key, value in diff.items():
print(f"- {key}: {value}")
9. 技能开发工作台搭建
9.1 本地开发环境配置
推荐使用VSCode开发环境配置:
json复制// .vscode/settings.json
{
"python.pythonPath": ".venv/bin/python",
"python.linting.enabled": true,
"python.linting.pylintEnabled": true,
"python.formatting.provider": "black",
"python.testing.pytestEnabled": true,
"files.exclude": {
"**/.git": true,
"**/.DS_Store": true,
"**/__pycache__": true
}
}
配套的Docker开发环境:
dockerfile复制# Dockerfile.dev
FROM python:3.9-slim
WORKDIR /app
# 安装系统依赖
RUN apt-get update && apt-get install -y \
git \
&& rm -rf /var/lib/apt/lists/*
# 安装Python依赖
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# 复制代码
COPY . .
# 开发模式启动
CMD ["bash", "-c", "while true; do sleep 1; done"]
9.2 自动化测试流水线
GitHub Actions示例:
yaml复制# .github/workflows/test.yml
name: Skill Tests
on: [push, pull_request]
jobs:
test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v2
- name: Set up Python
uses: actions/setup-python@v2
with:
python-version: '3.9'
- name: Install dependencies
run: |
python -m pip install --upgrade pip
pip install -r requirements.txt
pip install pytest pytest-cov
- name: Run tests
run: |
pytest --cov=./ --cov-report=xml
- name: Upload coverage
uses: codecov/codecov-action@v1
9.3 技能打包与分发
使用setuptools打包:
python复制# setup.py
from setuptools import setup, find_packages
setup(
name="workflow-skills",
version="0.1.0",
packages=find_packages(),
install_requires=[
"requests>=2.25.0",
"pydantic>=1.8.0"
],
extras_require={
"test": ["pytest", "pytest-mock"],
"dev": ["ipython", "black"]
},
entry_points={
"workflow.skills": [
"email_sender = skills.email:EmailSender",
"data_processor = skills.data:DataProcessor"
]
}
)
发布到私有仓库:
bash复制# 构建包
python setup.py sdist bdist_wheel
# 上传到私有仓库
twine upload --repository-url https://your-pypi.org/simple/ dist/*
10. 技能仓库管理实践
10.1 技能分类体系
建议的目录结构:
code复制skills/
├── communication/
│ ├── email_sender/
│ ├── sms_gateway/
│ └── notification/
├── data/
│ ├── extraction/
│ ├── transformation/
│ └── analysis/
├── integration/
│ ├── api_client/
│ └── db_connector/
└── utils/
├── datetime/
└── file_ops/
每个技能的标准结构:
code复制skill_name/
├── __init__.py # 主实现
├── tests/ # 单元测试
│ └── test_skill.py
├── docs/ # 文档
│ ├── usage.md
│ └── examples/
├── CHANGELOG.md # 变更日志
└── README.md # 使用说明
10.2 技能发现机制
实现一个简单的技能注册表:
python复制# registry.py
from typing import Dict, Type
from importlib import import_module
class SkillRegistry:
_skills: Dict[str, Type] = {}
@classmethod
def register(cls, name: str):
"""装饰器注册技能"""
def decorator(skill_class):
cls._skills[name] = skill_class
return skill_class
return decorator
@classmethod
def get_skill(cls, name: str):
"""获取技能类"""
if name not in cls._skills:
raise ValueError(f"未注册的技能: {name}")
return cls._skills[name]
@classmethod
def discover_skills(cls, package: str):
"""自动发现并注册技能"""
pkg = import_module(package)
for attr in dir(pkg):
if attr.endswith("Skill"):
skill_class = getattr(pkg, attr)
if isinstance(skill_class, type):
cls.register(attr[:-5].lower(), skill_class)
# 使用示例
@SkillRegistry.register("email")
class EmailSender:
pass
# 或者自动发现
SkillRegistry.discover_skills("skills.communication")
10.3 技能文档规范
每个技能应包含标准文档:
markdown复制# 邮件发送技能
## 功能描述
提供发送电子邮件的功能,支持HTML内容和附件
## 配置参数
| 参数名 | 类型 | 必填 | 说明 |
|--------|------|------|------|
| smtp_host | string | 是 | SMTP服务器地址 |
| smtp_port | number | 是 | 端口号 |
| username | string | 是 | 认证用户名 |
| password | string | 是 | 认证密码 |
## 输入参数
```json
{
"to": ["recipient@example.com"],
"subject": "邮件主题",
"content": "<p>邮件内容</p>",
"attachments": ["/path/to/file"]
}
```
## 输出结果
```json
{
"success": true,
"message_id": "<20220301123456.12345@example.com>"
}
```
## 使用示例
```python
from skills.email import EmailSender
config = {
"smtp_host": "smtp.example.com",
"smtp_port": 587,
"username": "user",
"password": "pass"
}
sender = EmailSender(config)
result = sender.execute({
"to": ["user@example.com"],
"subject": "测试邮件",
"content": "这是一封测试邮件"
})
```
## 变更历史
- v1.0.0 (2022-01-01): 初始版本
- v1.1.0 (2022-02-01): 增加附件支持
11. 技能组合与编排模式
11.1 顺序执行模式
基础链式调用:
python复制def sequential_flow(input_data):
"""顺序执行多个技能"""
# 步骤1: 数据准备
prepared = prepare_data(input_data)
# 步骤2: 数据处理
processed = process_data(prepared)
# 步骤3: 结果存储
stored = store_result(processed)
return stored
带错误处理的改进版:
python复制def safe_sequential(steps, input_data):
"""带错误处理的顺序执行"""
result = input_data
for step in steps:
try:
result = step.execute(result)
except Exception as e:
logger.error("步骤执行失败: %s", step.name)
raise WorkflowError(f"{step.name}执行失败") from e
return result
11.2 并行执行模式
使用线程池并行执行:
python复制from concurrent.futures import ThreadPoolExecutor
def parallel_execute(tasks: list, input_data: dict) -> dict:
"""并行执行多个技能"""
results = {}
with ThreadPoolExecutor(max_workers=5) as executor:
future_to_key = {
executor.submit(task.execute, input_data): task.name
for task in tasks
}
for future in concurrent.futures.as_completed(future_to_key):
key = future_to_key[future]
try:
results[key] = future.result()
except Exception as e:
results[key] = {"error": str(e)}
return results
11.3 条件分支模式
基于条件的动态路由:
python复制def conditional_flow(input_data):
"""条件分支工作流"""
# 初始处理
processed = preprocess(input_data)
# 条件判断
if processed["type"] == "A":
result = handle_type_a(processed)
elif processed["type"] == "B":
result = handle_type_b(processed)
else:
result = handle_default(processed)
# 后续处理
return postprocess(result)
使用策略模式改进:
python复制class FlowStrategy:
def should_handle(self, data) -> bool:
raise NotImplementedError
def handle(self, data):
raise NotImplementedError
class TypeAStrategy(FlowStrategy):
def should_handle(self, data):
return data["type"] == "A"
def handle(self, data):
return handle_type_a(data)
def strategic_flow(strategies: list, input_data):
"""策略模式工作流"""
processed = preprocess(input_data)
for strategy in strategies:
if strategy.should_handle(processed):
return strategy.handle(processed)
return handle_default(processed)
12. 性能优化进阶技巧
12.1 内存优化策略
处理大数据集时的技巧:
python复制def process_large_file(path):
"""流式处理大文件"""
with open(path, 'r') as f:
for line in f: # 逐行读取,不加载整个文件
yield process_line(line)
def batch_process(iterable, batch_size=1000):
"""分批处理可迭代对象"""
batch = []
for item in iterable:
batch.append(item)
if len(batch) >= batch_size:
yield process_batch(batch)
batch = []
if batch: # 处理剩余部分
yield process_batch(batch)
使用内存映射文件:
python复制import mmap
def search_large_file(path, pattern):
"""在超大文件中搜索"""
with open(path, 'r+') as f:
# 内存映射文件
mm = mmap.mmap(f.fileno(), 0)
# 直接在内存中搜索
index = mm.find(pattern.encode())
mm.close()
return index
12.2 CPU密集型任务优化
使用多进程池:
python复制from multiprocessing import Pool
def cpu_intensive_task(data):
"""CPU密集型任务"""
return heavy_computation(data)
def parallel
