三、组件封装的实战技巧
3.1 组件定义与设计原则
什么是组件?
组件是可独立开发、测试、部署、复用的软件单元,具有清晰的接口和明确的职责。
组件的五大特征:
可重用性:可以在不同场景复用
可组合性:可以组合成更大的系统
可替换性:可以用不同实现替换
独立性:可以独立开发和测试
封装性:隐藏内部实现细节
组件设计原则 - REUSE/RELEASE等价原则
组件的粒度应该足够小以便于复用,同时足够大以便于独立发布。
3.2 通用组件的封装模式
3.2.1 配置管理组件
# config_manager.py - 统一配置管理组件
import os
import json
import yaml
from typing import Any, Optional, Dict
from functools import lru_cache
from dataclasses import dataclass, field
from pathlib import Path
class ConfigSource(ABC):
"""配置源抽象"""
@abstractmethod
def get(self, key: str, default: Any = None) -> Any:
pass
@abstractmethod
def get_all(self) -> Dict[str, Any]:
pass
class EnvironmentConfigSource(ConfigSource):
"""环境变量配置源"""
def get(self, key: str, default: Any = None) -> Any:
return os.getenv(key, default)
def get_all(self) -> Dict[str, Any]:
return dict(os.environ)
class FileConfigSource(ConfigSource):
"""文件配置源 - 支持JSON/YAML"""
def __init__(self, file_path: str):
self.file_path = Path(file_path)
self._data = self._load_file()
def _load_file(self) -> Dict[str, Any]:
if not self.file_path.exists():
return {}
with open(self.file_path) as f:
if self.file_path.suffix == '.json':
return json.load(f)
elif self.file_path.suffix in ['.yml', '.yaml']:
return yaml.safe_load(f)
else:
raise ValueError(f"Unsupported config file format: {self.file_path.suffix}")
def get(self, key: str, default: Any = None) -> Any:
keys = key.split('.')
value = self._data
for k in keys:
if isinstance(value, dict):
value = value.get(k)
if value is None:
return default
else:
return default
return value if value is not None else default
def get_all(self) -> Dict[str, Any]:
return self._data
class ConfigManager:
"""配置管理器 - 支持多源配置合并"""
def __init__(self):
self.sources: List[ConfigSource] = []
self.cache: Dict[str, Any] = {}
self._cache_enabled = True
def add_source(self, source: ConfigSource, priority: int = 0):
"""添加配置源,优先级高的覆盖优先级低的"""
self.sources.append((priority, source))
self.sources.sort(key=lambda x: x[0], reverse=True)
self._clear_cache()
def add_environment_source(self, prefix: str = None):
"""添加环境变量源"""
class PrefixedEnvSource(EnvironmentConfigSource):
def __init__(self, pfx):
self.prefix = pfx
def get(self, key: str, default: Any = None) -> Any:
env_key = f"{self.prefix}_{key}".upper() if self.prefix else key.upper()
return super().get(env_key, default)
self.add_source(PrefixedEnvSource(prefix) if prefix else EnvironmentConfigSource())
def add_file_source(self, file_path: str, required: bool = False):
"""添加文件配置源"""
source = FileConfigSource(file_path)
if required and not source.get_all():
raise FileNotFoundError(f"Required config file not found: {file_path}")
self.add_source(source)
def get(self, key: str, default: Any = None, type_hint: type = None) -> Any:
"""获取配置值,支持类型转换"""
if self._cache_enabled and key in self.cache:
value = self.cache[key]
else:
value = self._get_from_sources(key, default)
if self._cache_enabled:
self.cache[key] = value
if type_hint and value is not None:
value = self._convert_type(value, type_hint)
return value
def _get_from_sources(self, key: str, default: Any = None) -> Any:
for _, source in self.sources:
value = source.get(key)
if value is not None:
return value
return default
def get_int(self, key: str, default: int = 0) -> int:
return int(self.get(key, default))
def get_bool(self, key: str, default: bool = False) -> bool:
value = self.get(key, default)
if isinstance(value, str):
return value.lower() in ['true', '1', 'yes', 'on']
return bool(value)
def get_list(self, key: str, default: List = None) -> List:
value = self.get(key, default)
if isinstance(value, str):
# 支持逗号分隔的字符串
return [v.strip() for v in value.split(',') if v.strip()]
return value if isinstance(value, list) else (default or [])
def _convert_type(self, value: Any, target_type: type) -> Any:
"""类型转换"""
if target_type == bool:
return self.get_bool(key="", default=value) # Hack for bool conversion
try:
return target_type(value)
except (ValueError, TypeError):
return value
def reload(self):
"""重新加载所有配置源"""
for _, source in self.sources:
if hasattr(source, 'reload'):
source.reload()
self._clear_cache()
def _clear_cache(self):
self.cache.clear()
def enable_cache(self, enabled: bool):
self._cache_enabled = enabled
# 使用示例 - 支持配置热更新
config = ConfigManager()
# 添加多个配置源,优先级:环境变量 > 本地文件 > 默认配置
config.add_environment_source(prefix="APP")
config.add_file_source("/etc/app/config.yaml")
config.add_file_source("./config.json") # 本地配置优先级更低
# 获取配置
db_host = config.get("database.host", default="localhost")
db_port = config.get_int("database.port", default=3306)
debug_mode = config.get_bool("debug", default=False)
3.2.2 日志组件封装
# logger.py - 统一日志组件
import logging
import sys
import json
from datetime import datetime
from typing import Any, Dict, Optional
from contextvars import ContextVar
import traceback
# 请求追踪上下文
request_id_var: ContextVar[str] = ContextVar('request_id', default='')
user_id_var: ContextVar[str] = ContextVar('user_id', default='')
class JSONFormatter(logging.Formatter):
"""JSON格式化器 - 结构化日志"""
def format(self, record: logging.LogRecord) -> str:
log_entry = {
"timestamp": datetime.utcnow().isoformat(),
"level": record.levelname,
"logger": record.name,
"message": record.getMessage(),
"module": record.module,
"function": record.funcName,
"line": record.lineno,
"request_id": request_id_var.get(),
"user_id": user_id_var.get()
}
# 添加异常信息
if record.exc_info:
log_entry["exception"] = {
"type": record.exc_info[0].__name__,
"message": str(record.exc_info[1]),
"traceback": traceback.format_exception(*record.exc_info)
}
# 添加额外字段
if hasattr(record, 'extra_fields'):
log_entry.update(record.extra_fields)
return json.dumps(log_entry, ensure_ascii=False)
class ColorfulConsoleFormatter(logging.Formatter):
"""彩色控制台格式化器 - 开发环境使用"""
COLORS = {
'DEBUG': '\033[36m', # 青色
'INFO': '\033[32m', # 绿色
'WARNING': '\033[33m', # 黄色
'ERROR': '\033[31m', # 红色
'CRITICAL': '\033[35m' # 紫色
}
RESET = '\033[0m'
def format(self, record: logging.LogRecord) -> str:
color = self.COLORS.get(record.levelname, self.RESET)
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f")[:-3]
# 添加请求追踪信息
request_id = request_id_var.get()
req_info = f"[{request_id}]" if request_id else ""
log_line = f"{color}{timestamp} [{record.levelname}] {req_info} {record.name}:{record.lineno} - {record.getMessage()}{self.RESET}"
if record.exc_info:
log_line += "\n" + traceback.format_exc()
return log_line
class LoggerFactory:
"""日志工厂 - 创建和管理日志实例"""
_instance = None
_initialized = False
def __new__(cls):
if cls._instance is None:
cls._instance = super().__new__(cls)
return cls._instance
def __init__(self):
if not self._initialized:
self._loggers: Dict[str, logging.Logger] = {}
self._configure_root_logger()
LoggerFactory._initialized = True
def _configure_root_logger(self, level: str = "INFO", json_format: bool = False):
"""配置根日志器"""
root_logger = logging.getLogger()
root_logger.setLevel(getattr(logging, level.upper()))
# 清除现有处理器
root_logger.handlers.clear()
# 控制台处理器
console_handler = logging.StreamHandler(sys.stdout)
if json_format:
console_handler.setFormatter(JSONFormatter())
else:
console_handler.setFormatter(ColorfulConsoleFormatter())
root_logger.addHandler(console_handler)
# 文件处理器(可选)
if os.getenv('LOG_FILE'):
file_handler = logging.FileHandler(os.getenv('LOG_FILE'))
file_handler.setFormatter(JSONFormatter())
root_logger.addHandler(file_handler)
def get_logger(self, name: str, **kwargs) -> 'AppLogger':
"""获取日志实例"""
if name not in self._loggers:
logger = logging.getLogger(name)
self._loggers[name] = AppLogger(logger)
return self._loggers[name]
class AppLogger:
"""应用日志包装器 - 提供更友好的接口"""
def __init__(self, logger: logging.Logger):
self._logger = logger
def _add_context(self, extra: Dict[str, Any] = None) -> Dict[str, Any]:
"""添加上下文信息"""
if extra is None:
extra = {}
extra['request_id'] = request_id_var.get()
extra['user_id'] = user_id_var.get()
return extra
def debug(self, message: str, *args, **kwargs):
extra = self._add_context(kwargs.pop('extra', None))
self._logger.debug(message, *args, extra=extra, **kwargs)
def info(self, message: str, *args, **kwargs):
extra = self._add_context(kwargs.pop('extra', None))
self._logger.info(message, *args, extra=extra, **kwargs)
def warning(self, message: str, *args, **kwargs):
extra = self._add_context(kwargs.pop('extra', None))
self._logger.warning(message, *args, extra=extra, **kwargs)
def error(self, message: str, *args, **kwargs):
extra = self._add_context(kwargs.pop('extra', None))
self._logger.error(message, *args, extra=extra, **kwargs)
def exception(self, message: str, *args, **kwargs):
"""记录异常,自动捕获异常堆栈"""
extra = self._add_context(kwargs.pop('extra', None))
self._logger.exception(message, *args, extra=extra, **kwargs)
def critical(self, message: str, *args, **kwargs):
extra = self._add_context(kwargs.pop('extra', None))
self._logger.critical(message, *args, extra=extra, **kwargs)
def bind_context(self, **kwargs):
"""绑定上下文(装饰器用法)"""
def decorator(func):
from functools import wraps
@wraps(func)
def wrapper(*args, **fkwargs):
# 保存当前上下文
old_request_id = request_id_var.get()
old_user_id = user_id_var.get()
# 设置新上下文
if 'request_id' in kwargs:
request_id_var.set(kwargs['request_id'])
if 'user_id' in kwargs:
user_id_var.set(kwargs['user_id'])
try:
return func(*args, **fkwargs)
finally:
# 恢复上下文
request_id_var.set(old_request_id)
user_id_var.set(old_user_id)
return wrapper
return decorator
def with_request_id(self, request_id: str):
"""设置请求ID上下文管理器"""
class RequestIdContext:
def __enter__(self_ctx):
self.old_request_id = request_id_var.get()
request_id_var.set(request_id)
def __exit__(self_ctx, *args):
request_id_var.set(self.old_request_id)
return RequestIdContext()
# 使用示例
logger_factory = LoggerFactory()
# 在应用启动时配置
logger_factory._configure_root_logger(level="DEBUG", json_format=False)
# 获取日志实例
logger = logger_factory.get_logger("order_service")
# 使用方式1:直接记录
logger.info("用户登录成功", extra={"user_id": "123", "ip": "192.168.1.1"})
# 使用方式2:使用上下文绑定
@logger.bind_context(request_id="req-001", user_id="user-456")
def process_order():
logger.info("开始处理订单") # 自动包含request_id和user_id
# 业务逻辑...
# 使用方式3:手动上下文管理
def handle_request(request_id: str):
with logger.with_request_id(request_id):
logger.info("收到请求")
# 处理请求...
3.2.3 数据库连接池组件
# db_pool.py - 数据库连接池组件
import asyncio
import aiomysql
from contextlib import asynccontextmanager
from typing import Optional, Dict, Any, List
from datetime import datetime
import time
from dataclasses import dataclass, field
@dataclass
class ConnectionStats:
"""连接统计"""
created: int = 0
acquired: int = 0
released: int = 0
closed: int = 0
errors: int = 0
slow_queries: int = 0
total_query_time: float = 0.0
class ConnectionPool:
"""数据库连接池 - 支持连接复用、自动重连、健康检查"""
def __init__(
self,
host: str,
port: int,
user: str,
password: str,
database: str,
min_size: int = 5,
max_size: int = 20,
max_idle_time: int = 300, # 5分钟空闲关闭
connection_timeout: int = 10,
query_timeout: int = 30,
enable_stats: bool = True
):
self.config = {
'host': host,
'port': port,
'user': user,
'password': password,
'db': database,
'autocommit': True,
'charset': 'utf8mb4'
}
self.min_size = min_size
self.max_size = max_size
self.max_idle_time = max_idle_time
self.connection_timeout = connection_timeout
self.query_timeout = query_timeout
self._pool: Optional[aiomysql.Pool] = None
self._stats = ConnectionStats() if enable_stats else None
self._slow_query_threshold = 1.0 # 1秒以上为慢查询
self._health_check_interval = 60
self._last_health_check = 0
async def initialize(self):
"""初始化连接池"""
self._pool = await aiomysql.create_pool(
**self.config,
minsize=self.min_size,
maxsize=self.max_size,
pool_recycle=self.max_idle_time,
connect_timeout=self.connection_timeout,
echo=False
)
if self._stats:
self._stats.created = self.min_size
# 启动健康检查
asyncio.create_task(self._health_check_loop())
async def close(self):
"""关闭连接池"""
if self._pool:
self._pool.close()
await self._pool.wait_closed()
@asynccontextmanager
async def acquire(self):
"""获取连接(上下文管理器)"""
conn = None
start_time = time.time()
try:
# 获取连接
conn = await self._pool.acquire()
if self._stats:
self._stats.acquired += 1
yield conn
except Exception as e:
if self._stats:
self._stats.errors += 1
raise
finally:
if conn:
await self._pool.release(conn)
if self._stats:
self._stats.released += 1
# 记录获取连接的耗时
elapsed = time.time() - start_time
if elapsed > 1.0: # 超过1秒的获取连接操作
logger.warning(f"Slow connection acquire: {elapsed:.3f}s")
async def execute(
self,
query: str,
params: tuple = None,
retry: int = 3
) -> int:
"""执行写操作(INSERT/UPDATE/DELETE)"""
start_time = time.time()
last_error = None
for attempt in range(retry):
try:
async with self.acquire() as conn:
async with conn.cursor() as cursor:
# 设置查询超时
if self.query_timeout:
await cursor.execute(f"SET SESSION max_execution_time={self.query_timeout * 1000}")
await cursor.execute(query, params or ())
await conn.commit()
rows = cursor.rowcount
# 记录统计
self._record_query_stats(query, start_time)
return rows
except aiomysql.Error as e:
last_error = e
if attempt < retry - 1:
await asyncio.sleep(2 ** attempt) # 指数退避
continue
raise
raise last_error
async def query_one(
self,
query: str,
params: tuple = None
) -> Optional[Dict[str, Any]]:
"""查询单条记录"""
start_time = time.time()
async with self.acquire() as conn:
async with conn.cursor(aiomysql.DictCursor) as cursor:
await cursor.execute(query, params or ())
result = await cursor.fetchone()
self._record_query_stats(query, start_time)
return result
async def query_all(
self,
query: str,
params: tuple = None
) -> List[Dict[str, Any]]:
"""查询多条记录"""
start_time = time.time()
async with self.acquire() as conn:
async with conn.cursor(aiomysql.DictCursor) as cursor:
await cursor.execute(query, params or ())
results = await cursor.fetchall()
self._record_query_stats(query, start_time)
return results
async def query_paginated(
self,
query: str,
params: tuple,
page: int = 1,
page_size: int = 20
) -> Dict[str, Any]:
"""分页查询"""
offset = (page - 1) * page_size
# 执行原查询获取总数
count_query = f"SELECT COUNT(*) as total FROM ({query}) as t"
total_result = await self.query_one(count_query, params)
total = total_result['total'] if total_result else 0
# 添加分页
paginated_query = f"{query} LIMIT %s OFFSET %s"
paginated_params = list(params) + [page_size, offset]
items = await self.query_all(paginated_query, tuple(paginated_params))
return {
'items': items,
'total': total,
'page': page,
'page_size': page_size,
'total_pages': (total + page_size - 1) // page_size
}
async def transaction(self):
"""事务上下文管理器"""
@asynccontextmanager
async def _transaction():
async with self.acquire() as conn:
await conn.begin()
try:
yield conn
await conn.commit()
except Exception:
await conn.rollback()
raise
return _transaction()
def _record_query_stats(self, query: str, start_time: float):
"""记录查询统计"""
if not self._stats:
return
elapsed = time.time() - start_time
self._stats.total_query_time += elapsed
if elapsed > self._slow_query_threshold:
self._stats.slow_queries += 1
# 记录慢查询日志
logger.warning(
f"Slow query detected: {elapsed:.3f}s",
extra={"query": query[:200], "duration": elapsed}
)
async def _health_check_loop(self):
"""健康检查循环"""
while True:
await asyncio.sleep(self._health_check_interval)
now = time.time()
if now - self._last_health_check < self._health_check_interval:
continue
await self._health_check()
self._last_health_check = now
async def _health_check(self):
"""检查连接池健康状态"""
try:
# 执行简单查询测试连接
result = await self.query_one("SELECT 1 as health_check")
if not result or result.get('health_check') != 1:
logger.error("Database health check failed")
# 触发重连
await self._reconnect()
except Exception as e:
logger.error(f"Database health check error: {e}")
await self._reconnect()
async def _reconnect(self):
"""重连数据库"""
logger.warning("Reconnecting to database...")
await self.close()
await self.initialize()
logger.info("Database reconnected successfully")
def get_stats(self) -> Dict[str, Any]:
"""获取连接池统计信息"""
if not self._stats:
return {}
return {
'pool_size': self._pool.size if self._pool else 0,
'pool_freesize': self._pool.freesize if self._pool else 0,
'stats': {
'created': self._stats.created,
'acquired': self._stats.acquired,
'released': self._stats.released,
'closed': self._stats.closed,
'errors': self._stats.errors,
'slow_queries': self._stats.slow_queries,
'avg_query_time': self._stats.total_query_time / max(1, self._stats.acquired)
}
}
# 使用示例
async def main():
# 创建连接池
pool = ConnectionPool(
host='localhost',
port=3306,
user='app_user',
password='password',
database='app_db',
min_size=10,
max_size=50
)
await pool.initialize()
try:
# 使用事务
async with pool.transaction() as conn:
async with conn.cursor() as cursor:
await cursor.execute(
"INSERT INTO orders (user_id, amount) VALUES (%s, %s)",
("user_123", 99.99)
)
order_id = cursor.lastrowid
await cursor.execute(
"INSERT INTO order_items (order_id, product_id, quantity) VALUES (%s, %s, %s)",
(order_id, "prod_001", 2)
)
# 查询
orders = await pool.query_all(
"SELECT * FROM orders WHERE user_id = %s",
("user_123",)
)
# 分页查询
result = await pool.query_paginated(
"SELECT * FROM orders WHERE user_id = %s ORDER BY created_at DESC",
("user_123",),
page=1,
page_size=10
)
# 获取统计
stats = pool.get_stats()
logger.info(f"Connection pool stats: {stats}")
finally:
await pool.close()
3.3 组件发布与版本管理
# setup.py - 组件打包配置
from setuptools import setup, find_packages
setup(
name="myapp-common-components",
version="1.0.0",
packages=find_packages(),
install_requires=[
"aiomysql>=0.1.0",
"redis>=4.0.0",
"pydantic>=2.0.0"
],
extras_require={
"dev": ["pytest", "black", "mypy"],
"aws": ["boto3"],
"gcp": ["google-cloud-storage"]
},
entry_points={
"console_scripts": [
"myapp-cli=myapp.cli:main",
],
},
classifiers=[
"Development Status :: 4 - Beta",
"Intended Audience :: Developers",
"License :: OSI Approved :: MIT License",
"Programming Language :: Python :: 3.8",
"Programming Language :: Python :: 3.9",
"Programming Language :: Python :: 3.10",
],
python_requires=">=3.8",
)
# 语义化版本管理
"""
版本格式:主版本号.次版本号.修订号
- 主版本号:不兼容的API修改
- 次版本号:向下兼容的功能性新增
- 修订号:向下兼容的问题修正
"""
# 组件版本演进示例
VERSIONS = {
"1.0.0": "首次发布",
"1.1.0": "新增连接池健康检查功能",
"1.2.0": "支持异步操作",
"2.0.0": "重构API,移除废弃的同步接口",
"2.0.1": "修复连接泄漏bug"
}
四、架构落地的实践框架
4.1 ADR(架构决策记录)
每个重要的架构决策都应该被记录下来:
# ADR-001: 选择异步架构
## 状态
已采纳
## 背景
系统需要支持高并发场景(预估峰值5000 QPS),I/O密集型操作较多(数据库查询、外部API调用)。
## 决策
采用全异步架构,使用asyncio + async/await模式。
## 理由
1. 相比于多线程,协程开销更小,可以支持更高并发
2. 避免线程安全问题
3. 与Python生态中的异步库(aiomysql、aiohttp、asyncpg)集成良好
## 后果
### 正面
- 单机支持5000+并发连接
- 代码更清晰,没有回调地狱
### 负面
- 需要全链路异步化,同步库需要使用线程池执行
- 调试相对复杂(async stack traces)
- 团队成员需要学习异步编程
## 备选方案
1. 多线程架构
- 优点:简单,与现有同步库兼容
- 缺点:线程开销大,GIL限制
2. 进程池架构
- 优点:绕过GIL
- 缺点:IPC开销大,状态同步复杂
## 相关决策
- ADR-002: 使用asyncio + uvloop作为事件循环
- ADR-003: 数据库连接池选用aiomysql
4.2 架构度量指标
# architecture_metrics.py - 架构健康度度量
import ast
import networkx as nx
from pathlib import Path
from typing import Dict, Set, Tuple
class ArchitectureMetrics:
"""架构度量工具"""
def __init__(self, project_root: Path):
self.project_root = project_root
self.dependency_graph = nx.DiGraph()
def analyze_cyclomatic_complexity(self, file_path: Path) -> Dict:
"""计算圈复杂度"""
with open(file_path) as f:
tree = ast.parse(f.read())
complexity = 0
for node in ast.walk(tree):
if isinstance(node, (ast.If, ast.While, ast.For, ast.ExceptHandler)):
complexity += 1
elif isinstance(node, ast.BoolOp):
complexity += len(node.values) - 1
return {
"file": str(file_path),
"cyclomatic_complexity": complexity,
"threshold_breached": complexity > 10
}
def analyze_package_dependencies(self) -> Dict:
"""分析包依赖关系,检测循环依赖"""
# 构建依赖图
for py_file in self.project_root.rglob("*.py"):
if "__pycache__" in str(py_file):
continue
module = self._get_module_name(py_file)
imports = self._extract_imports(py_file)
for imp in imports:
self.dependency_graph.add_edge(module, imp)
# 检测循环依赖
cycles = list(nx.simple_cycles(self.dependency_graph))
# 计算依赖深度
longest_path = nx.dag_longest_path(self.dependency_graph) if not cycles else []
return {
"total_modules": len(self.dependency_graph.nodes),
"total_dependencies": len(self.dependency_graph.edges),
"circular_dependencies": len(cycles),
"circular_examples": cycles[:5],
"max_dependency_depth": len(longest_path),
"deprecated_dependencies": self._check_deprecated_patterns()
}
def analyze_layer_compliance(self) -> Dict:
"""检查分层架构合规性"""
# 定义允许的依赖方向
allowed_deps = {
"presentation": ["application"],
"application": ["domain", "infrastructure"],
"domain": [], # 领域层不应该依赖其他层
"infrastructure": ["domain"]
}
violations = []
# 实际分析代码中的依赖...
return {
"compliant": len(violations) == 0,
"violations": violations,
"compliance_rate": 1.0 - (len(violations) / max(1, self.dependency_graph.number_of_edges()))
}
def generate_report(self) -> str:
"""生成架构度量报告"""
complexity = self._average_complexity()
deps = self.analyze_package_dependencies()
layers = self.analyze_layer_compliance()
return f"""
========== 架构健康度报告 ==========
代码复杂度:
- 平均圈复杂度:{complexity['avg']:.2f}
- 高复杂度文件数:{complexity['high_complexity_count']}
- 建议:复杂度超过10的函数需要重构
依赖分析:
- 模块总数:{deps['total_modules']}
- 循环依赖数:{deps['circular_dependencies']}
- 最大依赖深度:{deps['max_dependency_depth']}
分层合规性:
- 合规率:{layers['compliance_rate']:.1%}
- 违规数:{len(layers['violations'])}
总体评分:{self._calculate_score(complexity, deps, layers)}/100
"""
4.3 架构演进策略
# architecture_evolution.py - 架构演进策略
class ArchitectureEvolutionStrategy:
"""架构演进策略"""
# 演进模式
PATTERNS = {
"strangler_pattern": "绞杀者模式 - 逐渐替换旧系统",
"branch_by_abstraction": "抽象分支 - 通过抽象层过渡",
"parallel_run": "并行运行 - 新旧系统同时运行对比",
"seam_carving": "接缝拆分 - 在边界处切割"
}
@classmethod
def strangler_fig(cls, old_system, new_component, gateway):
"""
绞杀者模式实现
适用场景:逐步替换大型遗留系统
"""
# 步骤1:在入口处添加路由逻辑
gateway.add_route("/api/users", old_system, version="v1")
gateway.add_route("/api/users", new_component, version="v2")
# 步骤2:逐步迁移功能
features = ["用户登录", "用户注册", "用户信息", "用户设置"]
for feature in features:
# 灰度发布
gateway.add_feature_toggle(
feature_name=feature,
percentage=10, # 10%流量走新系统
target=new_component
)
# 监控指标
if new_component.error_rate < 0.01 and new_component.latency < old_system.latency:
# 逐渐放量
gateway.update_percentage(feature, 100)
elif new_component.error_rate > 0.05:
# 回滚
gateway.rollback(feature)
# 步骤3:完全替换后移除旧系统
if gateway.all_features_migrated():
old_system.decommission()
@classmethod
def branch_by_abstraction(cls, old_impl, new_impl):
"""
抽象分支模式实现
适用场景:重构核心模块
"""
class AbstractInterface:
"""抽象接口"""
def operation(self):
raise NotImplementedError()
class BranchImplementation(AbstractInterface):
def __init__(self, use_new: bool = False):
self.use_new = use_new
self.old = old_impl
self.new = new_impl
def operation(self):
if self.use_new:
return self.new.improved_operation()
else:
return self.new.legacy_operation()
# 使用配置控制
config = ConfigManager()
use_new = config.get_bool("use_new_implementation", False)
# 注入抽象分支
service = BranchImplementation(use_new=use_new)
return service