Python 高阶特性
Python 高阶特性是提升代码效率、可读性和工程化能力的核心,尤其在 AI Agent、后端开发、数据处理等场景中不可或缺。本文针对有 Java 基础但刚接触 Python 的开发者,通过「概念解释 + 核心语法 + Java 对比 + 实战代码」的形式,系统讲解 Python 高阶特性,帮你快速理解并落地。
1. 生成器 (Generator) 与 yield
1.1 什么是生成器
生成器是 Python 特有的惰性迭代器:
- 普通函数用
return一次性返回所有结果,生成器用yield关键字按需生成单个值; - 核心优势:极致省内存 —— 无需像列表一样把所有数据加载到内存,而是 “用一个生成一个”,适合处理超大数据集(如百万级日志、大数据流)。
Java 对比:
Python 生成器 ≈ Java 的
Iterator接口,但实现成本极低:
- Java 需手动实现
hasNext()和next()方法,维护迭代状态;- Python 仅需一个
yield关键字,解释器自动维护迭代暂停 / 继续状态。
1.2 核心语法
| 语法 | 作用 |
|---|---|
yield |
函数执行到此处暂停,返回当前值;下次调用(next()/ 遍历)从暂停处继续 |
next() |
手动获取生成器的下一个值,无值时抛出 StopIteration(遍历自动捕获) |
(i for i in range(n)) |
生成器表达式(极简创建生成器) |
1.3 实战案例
1.3.1 列表 vs 生成器(内存对比)
import sys
# 1. 列表推导式:立即生成所有数据,占用大量内存(百万级数据约8MB)
list_data = [i for i in range(1000000)]
print(f"列表占用内存: {sys.getsizeof(list_data)} 字节") # 输出:8000056 字节
# 2. 生成器表达式:仅保存生成规则,几乎不占内存(约100字节)
gen_data = (i for i in range(1000000))
print(f"生成器占用内存: {sys.getsizeof(gen_data)} 字节") # 输出:112 字节
# 验证:生成器按需生成值(遍历才会计算)
print(next(gen_data)) # 0
print(next(gen_data)) # 1
1.3.2 自定义生成器函数(斐波那契数列)
def fibonacci(n: int):
"""生成前n个斐波那契数的生成器"""
a, b = 0, 1 # 初始值
count = 0
while count < n:
yield a # 暂停并返回当前值,下次从这里继续
a, b = b, a + b # 更新值
count += 1
# 使用生成器(三种方式)
f = fibonacci(10)
# 方式1:手动next()
print(next(f)) # 0
print(next(f)) # 1
# 方式2:遍历(自动处理StopIteration)
for num in f:
print(num, end=" ") # 输出:1 2 3 5 8 13 21 34
# 方式3:转为列表(一次性消费所有值)
print("\n完整数列:", list(fibonacci(5))) # [0, 1, 1, 2, 3]
2. 装饰器 (Decorator)
2.1 什么是装饰器
装饰器本质是接收函数作为参数、返回新函数的闭包,核心价值是:
- 在不修改原函数代码的前提下,为函数增加额外功能(日志、计时、权限校验、缓存等);
- 是 Python 实现面向切面编程 (AOP) 的核心方式。
Java 对比:
Python 装饰器 ≈ Java 的「注解 + AOP」(如 Spring 的
@Transactional):
- Java 需通过 AspectJ/Spring AOP 实现切面逻辑,Python 仅需几行代码;
- Java 动态代理也能实现类似功能,但装饰器语法更直观、轻量化。
2.2 核心语法
| 语法 | 作用 |
|---|---|
@装饰器名 |
语法糖,等价于 原函数 = 装饰器(原函数) |
functools.wraps |
保留原函数的元数据(函数名、文档注释、参数说明),避免被装饰后丢失 |
*args/**kwargs |
接收任意位置参数 / 关键字参数,让装饰器适配所有函数 |
2.3 实战案例
2.3.1 基础:计时装饰器(通用版)
import time
import functools
def timer(func):
"""通用计时装饰器:计算函数执行耗时"""
# 保留原函数元数据(关键!否则 func.__name__ 会变成 wrapper)
@functools.wraps(func)
def wrapper(*args, **kwargs):
start_time = time.time()
# 执行原函数并获取返回值
result = func(*args, **kwargs)
end_time = time.time()
# 打印耗时(保留4位小数)
print(f"[{func.__name__}] 执行耗时: {end_time - start_time:.4f} 秒")
return result # 必须返回原函数结果,否则调用者拿不到返回值
return wrapper
# 使用装饰器
@timer # 等价于 heavy_task = timer(heavy_task)
def heavy_task(n: float):
"""模拟耗时IO/计算任务"""
time.sleep(n) # 模拟等待(IO密集型)
return f"任务完成(休眠{n}秒)"
@timer # 装饰器适配任意函数
def calculate_sum(num: int) -> int:
"""模拟计算任务"""
return sum(range(num))
# 调用测试
print(heavy_task(1.5)) # 输出:[heavy_task] 执行耗时: 1.5001 秒 + 任务完成
print(calculate_sum(1000000))# 输出:[calculate_sum] 执行耗时: 0.0200 秒 + 499999500000
2.3.2 进阶:带参数的装饰器(日志等级可控)
# 临时导入日志
import logging
import functools
# 配置日志
logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s")
def logger(level: str = "info"):
"""带参数的日志装饰器:自定义日志等级"""
def decorator(func):
@functools.wraps(func)
def wrapper(*args, **kwargs):
# 根据参数选择日志等级
log_func = getattr(logging, level.lower())
log_func(f"开始执行 {func.__name__},参数:args={args}, kwargs={kwargs}")
result = func(*args, **kwargs)
log_func(f"{func.__name__} 执行完成,返回值:{result}")
return result
return wrapper
return decorator
# 使用带参数的装饰器
@logger(level="warning") # 日志等级为WARNING
def delete_user(user_id: int) -> bool:
"""模拟删除用户"""
print(f"删除用户 {user_id} 成功")
return True
@logger() # 使用默认等级INFO
def add_user(name: str, age: int) -> dict:
"""模拟新增用户"""
return {"id": 1001, "name": name, "age": age}
# 调用测试
delete_user(100) # 输出 WARNING 级日志
add_user("Alice", 25) # 输出 INFO 级日志
3. 类型注解 (Type Hints)
3.1 为什么需要类型注解
Python 是动态类型语言(变量类型可随时改变),优点是灵活,缺点是:
- 大型项目中易因类型错误导致 Bug(如把字符串传给需要整数的函数);
- IDE 无法提供智能提示,开发效率低;
- 团队协作时,代码可读性差(不清楚函数参数 / 返回值类型)。
类型注解的核心价值:静态提示,动态不校验 —— 既保留 Python 灵活性,又能享受静态语言的类型安全。
Java 对比:
- Python 类型注解是「提示性」的(运行时不校验,写错也能跑,IDE 会警告);
- Java 是「强类型」语言(编译时强制校验,类型错误直接编译失败);
- Python 配合
mypy工具可实现静态类型检查,接近 Java 的编译校验。
3.2 核心语法(Python 3.8+)
| 场景 | 语法示例 | 说明 | |
|---|---|---|---|
| 基础变量 | age: int = 25 |
变量名:类型 = 值 | |
| 函数参数 / 返回值 | def func(a: int) -> bool: |
参数标注类型,-> 标注返回值类型 | |
| 容器类型 | List[int] / Dict[str, int] |
列表(元素为 int)/ 字典(键 str 值 int) | |
| 可选类型 | Optional[str] |
等价于 `str | None`(Python 3.10+) |
| 联合类型 | Union[int, float] |
等价于 `int | float`(Python 3.10+) |
| 任意类型 | Any |
不限制类型(等同于无注解) |
3.3 实战案例
# Python 3.9+ 可直接用 list/dict,无需导入 typing;3.8及以下需 from typing import List, Dict
from typing import List, Dict, Optional, Union, Any
# 1. 基础变量注解
name: str = "Alice"
age: int = 25
height: float = 1.68
is_student: bool = True
# 2. 容器类型注解
scores: List[int] = [80, 90, 95] # 整数列表
user_info: Dict[str, Union[str, int]] = {"name": "Bob", "age": 30} # 混合类型字典
optional_name: Optional[str] = None # 可能为字符串或None
mixed_type: Union[int, str] = 100 # 可以是int或str
# 3. 函数注解(核心:提升可读性和IDE提示)
def greet(name: str, times: int = 1) -> str:
"""
生成问候语
Args:
name (str): 要问候的人名(必填)
times (int): 重复次数(默认1)
Returns:
str: 拼接后的问候语
"""
return f"Hello, {name}! " * times
def find_user(user_id: int) -> Optional[str]:
"""
根据ID查找用户
Args:
user_id (int): 用户ID
Returns:
Optional[str]: 找到返回用户名,未找到返回None
"""
if user_id == 1:
return "Admin"
return None
# 4. 函数调用(IDE会提示参数类型,写错会警告)
print(greet("Bob", 2)) # 正确:Hello, Bob! Hello, Bob!
# greet(123) # IDE警告:Expected type 'str', got 'int' instead(运行时仍能执行)
# 5. 类型别名(简化复杂类型)
UserList = List[Dict[str, Union[str, int]]] # 定义类型别名
def get_users() -> UserList:
"""返回用户列表"""
return [{"name": "Charlie", "age": 35}, {"name": "David", "age": 40}]
4. Pydantic(数据验证神器)
4.1 什么是 Pydantic
Pydantic 是 Python 最流行的数据验证与序列化库,核心能力:
- 定义数据模型(继承
BaseModel),自动校验字段类型和规则; - 自动类型转换(如字符串 “123” 转整数 123);
- 支持 JSON 序列化 / 反序列化,完美适配 API/LLM 数据处理;
- 是 LangChain、FastAPI、AI Agent 开发的核心依赖。
Java 对比:
Pydantic ≈ Java 的「JavaBean + Hibernate Validator (JSR-380)」:
- Java 需手动写 getter/setter + 注解(如
@NotBlank@Min);- Pydantic 自动生成校验逻辑,且支持更灵活的字段规则(如
min_lengthge);- 特别适合处理 LLM 返回的结构化 JSON 数据,避免手动解析出错。
4.2 核心语法
| 语法 / 类 | 作用 |
|---|---|
BaseModel |
所有数据模型的基类,提供校验、序列化 / 反序列化能力 |
Field(...) |
定义字段规则(必填 / 默认值 / 长度 / 范围),... 表示必填 |
ValidationError |
校验失败时抛出的异常,包含详细的错误信息 |
model_dump() |
将模型转为字典,model_dump_json() 转为 JSON 字符串 |
4.3 实战案例
from pydantic import BaseModel, Field, ValidationError
from typing import List, Optional
# 1. 定义数据模型(核心)
class User(BaseModel):
"""用户数据模型(AI Agent 接收的用户信息)"""
id: int # 必填字段,类型int
name: str = Field(..., min_length=2, max_length=10, description="用户名,2-10个字符")
# ... 表示必填;min_length/max_length 限制长度
tags: List[str] = Field(default=[], description="用户标签列表,默认空")
# 默认值为空列表
email: Optional[str] = Field(None, pattern=r'^[\w\.-]+@[\w\.-]+\.\w+$', description="邮箱(可选)")
# 可选字段,正则校验邮箱格式;None 表示默认值
score: float = Field(0.0, ge=0, le=100, description="评分,0-100")
# ge=大于等于,le=小于等于
# 2. 数据校验(自动类型转换 + 规则校验)
try:
# 原始数据(模拟API/LLM返回的JSON转字典)
raw_data = {
"id": "123", # 字符串自动转为int
"name": "Alice", # 符合长度规则
"tags": ["AI", "Python"], # 列表元素为str
"score": 95.5 # 符合0-100范围
}
# 实例化模型(自动校验)
user = User(**raw_data)
print("校验成功:", user)
# 输出:id=123 name='Alice' tags=['AI', 'Python'] email=None score=95.5
# 序列化(转为JSON)
print("JSON格式:", user.model_dump_json(indent=2, ensure_ascii=False))
except ValidationError as e:
print("数据校验失败:", e)
# 3. 错误案例(触发校验异常)
try:
# 错误数据:name太短、id非数字、score超出范围
User(id="abc", name="A", score=101)
except ValidationError as e:
# 打印详细错误信息
print("\n错误详情:")
for error in e.errors():
print(f"字段 {error['loc'][0]}: {error['msg']}")
# 输出:
# 字段 id: Input should be a valid integer
# 字段 name: String should have at least 2 characters
# 字段 score: Input should be less than or equal to 100
5. 异步编程 (Async/Await)
5.1 为什么需要异步
Python 同步代码在处理 IO 密集型任务(爬虫、API 调用、数据库查询)时,会 “傻等” 响应(如等待网络请求返回),导致 CPU 闲置。
异步编程的核心:在等待 IO 时切换到其他任务,极大提升并发效率(如同时处理 1000 个 HTTP 请求,耗时与处理 1 个几乎相同)。
Java 对比:
Python
async/await≈ Java 的CompletableFuture/ JavaScriptPromise:
- Python 基于「单线程事件循环 (Event Loop)」实现并发,无线程切换开销;
- Java 传统上用线程池处理并发,JDK 21+ 引入虚拟线程(Virtual Threads),接近 Python 协程的轻量性;
- Python 异步适合高并发 IO(如 Web 服务、聊天机器人),Java 多线程适合 CPU 密集型 + IO 密集型混合场景。
5.2 核心语法
| 语法 / 函数 | 作用 |
|---|---|
async def |
定义协程函数(不能直接调用,需用 await 或 asyncio.run()) |
await |
挂起当前协程,等待耗时操作完成(仅能在协程函数内使用) |
asyncio.run() |
运行顶层协程函数(程序入口) |
asyncio.gather() |
并发执行多个协程,等待所有完成并返回结果列表 |
asyncio.sleep() |
模拟异步耗时操作(非阻塞,区别于 time.sleep() 的阻塞) |
5.3 实战案例(同步 vs 异步对比)
import asyncio
import time
# 模拟耗时IO操作(如HTTP请求、数据库查询)
async def download_url(url: str):
"""异步下载URL(非阻塞)"""
print(f"开始下载: {url}")
await asyncio.sleep(1) # 模拟网络延迟(非阻塞,事件循环可处理其他任务)
print(f"下载完成: {url}")
return f"{url} 数据"
def sync_download_url(url: str):
"""同步下载URL(阻塞)"""
print(f"开始下载: {url}")
time.sleep(1) # 阻塞,CPU闲置
print(f"下载完成: {url}")
return f"{url} 数据"
# 异步主函数
async def async_main():
start = time.time()
# 创建协程任务列表
tasks = [
download_url("http://baidu.com"),
download_url("http://google.com"),
download_url("http://python.org")
]
# 并发执行所有任务(总耗时≈1秒)
results = await asyncio.gather(*tasks)
end = time.time()
print(f"\n异步总耗时: {end - start:.2f} 秒")
print("异步结果:", results)
# 同步主函数
def sync_main():
start = time.time()
# 串行执行(总耗时≈3秒)
results = [
sync_download_url("http://baidu.com"),
sync_download_url("http://google.com"),
sync_download_url("http://python.org")
]
end = time.time()
print(f"\n同步总耗时: {end - start:.2f} 秒")
print("同步结果:", results)
# 运行测试
if __name__ == "__main__":
print("=== 异步执行 ===")
asyncio.run(async_main()) # 运行异步程序
print("\n=== 同步执行 ===")
sync_main()
# 输出结果:
# === 异步执行 ===
# 开始下载: http://baidu.com
# 开始下载: http://google.com
# 开始下载: http://python.org
# 下载完成: http://baidu.com
# 下载完成: http://google.com
# 下载完成: http://python.org
# 异步总耗时: 1.00 秒
#
# === 同步执行 ===
# 开始下载: http://baidu.com
# 下载完成: http://baidu.com
# 开始下载: http://google.com
# 下载完成: http://google.com
# 开始下载: http://python.org
# 下载完成: http://python.org
# 同步总耗时: 3.00 秒
6. 正则表达式 (re 模块)
6.1 核心功能
正则表达式是处理字符串的 “瑞士军刀”,核心用于:
- 匹配:检查字符串是否符合规则(如手机号、邮箱);
- 查找:提取字符串中的目标内容(如所有手机号);
- 替换:批量修改字符串(如脱敏、格式化);
- 分割:按复杂规则分割字符串(如按任意空白符分割)。
Java 对比:
- 转义字符:Java 字符串中反斜杠需双重转义(如匹配数字需写
"\\d"),Python 可使用r"..."原生字符串(只需写r"\d"),极大简化正则编写;- API 差异:Java 使用
Pattern和Matcher类,Python 直接使用re模块函数(如re.findall)更便捷。
6.2 常用方法与语法
| 方法 | 作用 |
|---|---|
re.search() |
查找第一个匹配项,返回 Match 对象(无则 None) |
re.findall() |
查找所有匹配项,返回列表 |
re.sub() |
替换匹配的内容(支持正则分组) |
re.compile() |
预编译正则表达式,提升多次使用的效率 |
| 正则元字符 | 作用 | 示例 |
|---|---|---|
\d |
匹配数字 | \d{11} 匹配 11 位手机号 |
\w |
匹配字母 / 数字 / 下划线 | \w+ 匹配单词 |
[] |
匹配字符集 | [a-z] 匹配小写字母 |
+ |
匹配 1 次或多次 | \d+ 匹配连续数字 |
* |
匹配 0 次或多次 | \w* 匹配任意单词(包括空) |
{n} |
匹配恰好 n 次 | \d{4} 匹配 4 位数字 |
6.3 实战案例
import re
# 原始文本(模拟爬虫/日志内容)
text = """
联系邮箱: support@python.org, 备用邮箱: admin@example.com
客服电话: 13800138000, 技术电话: 13912345678
用户评论: "Python 高阶特性很实用!123456"
"""
# 1. 匹配单个内容(邮箱)
email_pattern = r'[\w\.-]+@[\w\.-]+\.\w+' # 简易邮箱正则
email_match = re.search(email_pattern, text)
if email_match:
print("找到第一个邮箱:", email_match.group()) # support@python.org
# 2. 查找所有匹配项(手机号)
phone_pattern = r'\d{11}' # 11位数字(手机号)
phones = re.findall(phone_pattern, text)
print("所有手机号:", phones) # ['13800138000', '13912345678']
# 3. 替换内容(手机号脱敏)
desensitized_text = re.sub(r'(\d{3})\d{4}(\d{4})', r'\1****\2', text)
# \1 表示第一个分组,\2 表示第二个分组
print("脱敏后文本:\n", desensitized_text)
# 4. 预编译正则(多次使用时提升效率)
comment_pattern = re.compile(r'[^\d\s]+') # 匹配非数字、非空白的字符
comments = comment_pattern.findall(text)
## 7. 数据类 (Dataclasses)
### 7.1 什么是 Dataclasses
Python 3.7+ 引入的 `@dataclass` 装饰器,旨在**自动生成类的常用方法**(如 `__init__`, `__repr__`, `__eq__`),极大地简化了 “仅用于存储数据” 的类定义。
> **Java 对比**:
>
> - Python `@dataclass` ≈ Java 14+ 的 **Record (记录类)** 或 Lombok 的 **`@Data`** 注解;
> - 核心目的相同:消除样板代码(Boilerplate Code),让开发者专注于数据定义。
### 7.2 核心语法
| 语法 | 作用 |
| :---------------: | :----------------------------------------------------------: |
| `@dataclass` | 类装饰器,自动生成 `__init__`, `__repr__` 等方法 |
| `frozen=True` | 参数,创建**不可变对象**(类似 Java Record),属性只读 |
| `field(default=)` | `dataclasses.field`,用于定义复杂默认值(如列表、工厂函数) |
### 7.3 实战案例
```python
from dataclasses import dataclass, field
from typing import List
# 1. 传统类写法(繁琐)
class UserOld:
def __init__(self, name: str, age: int):
self.name = name
self.age = age
def __repr__(self):
return f"UserOld(name={self.name}, age={self.age})"
def __eq__(self, other):
if not isinstance(other, UserOld):
return False
return self.name == other.name and self.age == other.age
# 2. Dataclass 写法(简洁)
@dataclass
class User:
name: str
age: int
# 复杂默认值需用 field(default_factory=...)
tags: List[str] = field(default_factory=list)
# 验证
u1 = User("Alice", 30)
u2 = User("Alice", 30)
u3 = User("Bob", 25)
# 自动生成的 __repr__
print(u1) # User(name='Alice', age=30, tags=[])
# 自动生成的 __eq__ (按属性值比较)
print(u1 == u2) # True (内容相同即相等,类似 Java Record)
print(u1 == u3) # False
# 3. 不可变数据类 (frozen=True)
@dataclass(frozen=True)
class Point:
x: int
y: int
p = Point(10, 20)
print(p)
# p.x = 30 # 报错:FrozenInstanceError (类似 Java Record 的 final 属性)
print(“提取非数字文本:”, comments) # 包含邮箱、评论等非数字内容
7. 日志模块 (logging)
7.1 为什么不用 print
print 是调试工具,不适合生产环境:
- 无法分级(调试 / 信息 / 警告 / 错误);
- 无法输出到文件,仅能打印到控制台;
- 无法自定义格式(如时间、日志等级);
logging是 Python 标准库,满足生产环境所有日志需求。
Java 对比:
Python
logging≈ Java 的Log4j/SLF4J:
- 核心概念一致:Logger(记录器)、Handler(处理器)、Formatter(格式化器)、Level(级别);
- Python 配置更简单,无需额外依赖;
- 日志等级完全对应:DEBUG < INFO < WARNING < ERROR < CRITICAL。
7.2 核心语法
| 日志等级 | 作用 | 使用场景 |
|---|---|---|
DEBUG |
调试信息(开发阶段) | 打印变量、函数入参 |
INFO |
普通信息(运行状态) | 系统启动、任务完成 |
WARNING |
警告信息(非致命问题) | 磁盘空间不足、配置过期 |
ERROR |
错误信息(功能异常) | 数据库连接失败、API 调用失败 |
CRITICAL |
严重错误(系统崩溃) | 磁盘满、核心服务挂掉 |
7.3 实战案例
import logging
# 1. 基础配置(一次性配置,全局生效)
logging.basicConfig(
level=logging.DEBUG, # 日志等级(低于该等级的日志不输出)
format='%(asctime)s - %(name)s - %(levelname)s - %(filename)s:%(lineno)d - %(message)s',
# 格式:时间 - 日志器名 - 等级 - 文件名:行号 - 消息
handlers=[
logging.FileHandler('app.log', encoding='utf-8'), # 输出到文件
logging.StreamHandler() # 输出到控制台
]
)
# 2. 获取日志器(推荐按模块名命名)
logger = logging.getLogger(__name__)
# 3. 输出不同等级的日志
logger.debug("这是调试信息(变量值:%s)", "test123") # 支持格式化参数
logger.info("系统启动成功,开始处理任务")
logger.warning("磁盘空间不足 20%,请及时清理")
try:
1 / 0 # 模拟错误
except ZeroDivisionError as e:
logger.error("计算出错:%s", e, exc_info=True) # exc_info=True 打印堆栈信息
logger.critical("核心数据库连接失败,系统即将退出")
# 4. 日志文件内容示例:
# 2024-05-20 10:00:00,123 - __main__ - DEBUG - demo.py:15 - 这是调试信息(变量值:test123)
# 2024-05-20 10:00:00,124 - __main__ - INFO - demo.py:16 - 系统启动成功,开始处理任务
# ...
8. 综合案例:Agent 数据清洗管道
本案例整合所有高阶特性,模拟 AI Agent 的数据处理流程:
- 生成器:流式读取海量原始数据;
- Pydantic:校验数据格式,过滤脏数据;
- 装饰器:记录处理耗时;
- 日志模块:记录运行状态;
- 类型注解:提升代码可读性。
8.1 完整代码
import time
import logging
import random
from typing import Iterator, Tuple
from pydantic import BaseModel, ValidationError, Field
# ===================== 1. 基础配置 =====================
# 配置日志(生产环境级)
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(module)s:%(lineno)d - %(message)s',
handlers=[logging.FileHandler('agent_pipeline.log', encoding='utf-8'), logging.StreamHandler()]
)
logger = logging.getLogger(__name__)
# ===================== 2. 数据模型定义(Pydantic) =====================
class UserData(BaseModel):
"""AI Agent 处理的用户评论数据模型"""
user_id: int = Field(..., ge=1000, description="用户ID,至少1000")
content: str = Field(..., min_length=5, description="评论内容,至少5个字符")
score: float = Field(default=0.0, ge=0, le=100, description="评论评分,0-100")
# ===================== 3. 装饰器定义 =====================
def process_timer(func):
"""记录函数处理耗时的装饰器"""
def wrapper(*args, **kwargs) -> Tuple[int, int]:
start_time = time.time()
valid_count, error_count = func(*args, **kwargs)
end_time = time.time()
logger.info(f"[{func.__name__}] 总耗时: {end_time - start_time:.4f} 秒")
return valid_count, error_count
return wrapper
# ===================== 4. 生成器定义(流式数据) =====================
def data_stream_generator(batch_size: int) -> Iterator[dict]:
"""
模拟从数据库/消息队列流式读取原始数据
:param batch_size: 读取批次大小
:return: 原始数据迭代器(dict)
"""
logger.info(f"开始生成 {batch_size} 条原始数据(含脏数据)")
for i in range(batch_size):
# 模拟随机数据,每5条生成1条脏数据(短内容)
yield {
"user_id": 1000 + i,
"content": f"User comment {i + 1}" if i % 5 != 0 else "Hi", # 脏数据:仅2字符
"score": round(random.uniform(0, 100), 1)
}
time.sleep(0.05) # 模拟IO延迟(50ms/条)
logger.info("原始数据生成完成")
# ===================== 5. 核心处理管道 =====================
@process_timer
def run_agent_pipeline(batch_size: int) -> Tuple[int, int]:
"""
AI Agent 数据清洗管道主函数
:param batch_size: 处理批次大小
:return: (有效数据数, 错误数据数)
"""
logger.info(">>> 启动 AI Agent 数据清洗管道 <<<")
valid_count = 0
error_count = 0
# 流式处理数据(生成器)
for raw_data in data_stream_generator(batch_size):
try:
# 数据校验(Pydantic)
user_data = UserData(**raw_data)
# 模拟业务处理(如关键词提取、情感分析)
logger.debug(f"处理有效数据: ID={user_data.user_id}, 内容={user_data.content}")
valid_count += 1
except ValidationError as e:
# 捕获校验错误,记录详情
logger.warning(f"脏数据过滤: {raw_data} | 错误原因: {e.errors()[0]['msg']}")
error_count += 1
# 统计结果
logger.info(f"<<< 管道执行结束 >>> 有效数据: {valid_count} 条, 错误数据: {error_count} 条")
return valid_count, error_count
# ===================== 6. 程序入口 =====================
if __name__ == "__main__":
# 运行管道(处理10条数据)
valid, error = run_agent_pipeline(10)
# 输出最终统计
print(f"\n最终结果:有效数据 {valid} 条,错误数据 {error} 条")
8.2 案例流程图
8.3 运行效果
2024-05-20 11:00:00,000 - INFO - demo:58 - >>> 启动 AI Agent 数据清洗管道 <<<
2024-05-20 11:00:00,001 - INFO - demo:42 - 开始生成 10 条原始数据(含脏数据)
2024-05-20 11:00:00,052 - WARNING - demo:64 - 脏数据过滤: {'user_id': 1000, 'content': 'Hi', 'score': 88.5} | 错误原因: String should have at least 5 characters
2024-05-20 11:00:00,103 - DEBUG - demo:61 - 处理有效数据: ID=1001, 内容=User comment 2
2024-05-20 11:00:00,154 - DEBUG - demo:61 - 处理有效数据: ID=1002, 内容=User comment 3
2024-05-20 11:00:00,205 - DEBUG - demo:61 - 处理有效数据: ID=1003, 内容=User comment 4
2024-05-20 11:00:00,256 - DEBUG - demo:61 - 处理有效数据: ID=1004, 内容=User comment 5
2024-05-20 11:00:00,307 - WARNING - demo:64 - 脏数据过滤: {'user_id': 1005, 'content': 'Hi', 'score': 45.2} | 错误原因: String should have at least 5 characters
2024-05-20 11:00:00,358 - DEBUG - demo:61 - 处理有效数据: ID=1006, 内容=User comment 7
2024-05-20 11:00:00,409 - DEBUG - demo:61 - 处理有效数据: ID=1007, 内容=User comment 8
2024-05-20 11:00:00,460 - DEBUG - demo:61 - 处理有效数据: ID=1008, 内容=User comment 9
2024-05-20 11:00:00,511 - DEBUG - demo:61 - 处理有效数据: ID=1009, 内容=User comment 10
2024-05-20 11:00:00,512 - INFO - demo:48 - 原始数据生成完成
2024-05-20 11:00:00,513 - INFO - demo:68 - <<< 管道执行结束 >>> 有效数据: 8 条, 错误数据: 2 条
2024-05-20 11:00:00,514 - INFO - demo:35 - [run_agent_pipeline] 总耗时: 0.5130 秒
最终结果:有效数据 8 条,错误数据 2 条
评论