API数据采集:用requests自动拉取MES生产数据【P0级重磅】
API数据采集:用requests自动拉取MES生产数据【P0级重磅】
一、问题背景:每天点了100次鼠标
这是FAB工程师小李的真实工作场景。
每天早上8点,小李到岗后的第一个任务是"喂数据"——给各种分析工具准备数据:
1. 从MES系统导出昨日所有工序的Lot数据(5分钟)
2. 导出Wafer厚度测量数据(5分钟)
3. 导出设备报警日志(3分钟)
4. 导出良率统计(3分钟)
5. 整理成Excel发给各车间主管(10分钟)
6. 重复以上操作给不同分析系统……
每天花50分钟,重复一个机械的"导出→保存"操作。 一周五天,一年52周——262分钟×52周 = 超过200小时/年,花在纯鼠标点击上。
更严重的是:
- 数据断档:手工导出容易漏日期、漏工序,分析结果因此不完整
- 无法实时:等手工导完数据,已经是上午10点,分析只能看昨天的旧数据
- 系统割裂:MES的数据无法自动流入SPC系统、良率预测模型、产能规划工具
MES系统其实有API接口——给程序用的数据取餐口。 但99%的工程师不知道它的存在,或者知道但不知道怎么用。
学完这一篇,你将能做到:
- 用Python自动从MES API获取任意时间范围的数据
- 处理认证(Token自动刷新)
- 处理分页(数据量再大也能全部拿完)
- 优雅地处理网络异常(断网/超时/限流自动重试)
- 把数据存到数据库,供后续SPC/良率分析使用
---
二、技术原理:HTTP协议、认证机制与数据分页
2.1 HTTP协议:程序之间的对话规则
API(Application Programming Interface)本质上是两台计算机之间的约定好的通信协议。类比现实:HTTP就像快递公司的运单格式——发件人(客户端Python)和收件人(MES服务器)都遵守同一套格式,对话才能进行。
```python
import requests
发起GET请求(向MES服务器"要"数据)
url = "http://mes-api.fab.local/api/v1/lots"
response = requests.get(url, timeout=30)
HTTP状态码:服务端告诉客户端"这件事办得怎么样"
print(response.status_code) # 200=成功, 401=没权限, 429=请求太频繁, 500=服务器崩了
print(response.json()) # 把JSON响应解析成Python字典
```
2.2 四种认证机制对比
MES系统的API通常需要认证,防止未授权访问。以下是四种常见认证方式:
| 认证方式 | 原理 | 安全性 | 适用场景 | Python实现 |
|---------|------|--------|---------|-----------|
| API Key | 请求头附带固定密钥 | 中 | 简单场景 | `headers={'X-API-Key': 'your_key'}` |
| Bearer Token | OAuth2短期令牌 | 高 | 主流MES系统 | `headers={'Authorization': 'Bearer xxx'}` |
| Basic Auth | 用户名+密码Base64编码 | 低 | 内网测试接口 | `auth=('user','pass')` |
| WS-Security | SOAP+数字签名(老系统) | 高 | 传统半导体MES | Zeep库处理 |
Bearer Token的工作流程:
1. 首次登录:用用户名+密码换取一个Token(有效期通常2~24小时)
2. 后续请求:把Token放在请求头里
3. Token过期:重新登录获取新Token
```python
获取Token
token_resp = requests.post(
"http://mes-api.fab.local/auth/login",
json={"username": "fab_engineer", "password": "password123"}
)
access_token = token_resp.json()['access_token'] # {"access_token": "eyJhbG..."}
带上Token发请求
headers = {"Authorization": f"Bearer {access_token}"}
data = requests.get("http://mes-api.fab.local/api/lots", headers=headers).json()
```
2.3 数据分页:API不是一次全给你的
API通常不会一次返回全部数据(想象一下一次返回287万条记录——网络会崩溃)。它会分页返回,需要客户端逐页请求。
三种分页策略:
```python
策略1:偏移量分页(最常见)—— LIMIT/OFFSET
params = {"limit": 1000, "offset": 0} # 第1页:取0~999
params = {"limit": 1000, "offset": 1000} # 第2页:取1000~1999
适合:数据量已知、顺序不变的场景
策略2:游标分页(推荐)—— 用上一页最后一条的ID作为起点
params = {"limit": 1000, "cursor": "last_id_999"} # 从上次的最后一条继续
适合:数据实时新增、偏移量会漏数据的场景
策略3:时间范围分页 —— 按时间窗口切分
params = {"start_time": "2026-06-01T00:00:00", "end_time": "2026-06-01T23:59:59"}
适合:FAB按班次/日期查询的场景
```
2.4 限流策略:不要把MES服务器打爆
MES API通常有速率限制:
| 限制类型 | 典型值 | 含义 |
|---------|--------|------|
| 每分钟 | 60次 | 超过60次/分钟会被临时封禁 |
| 每小时 | 1000次 | 大规模采集需要控制节奏 |
| 每天 | 10000次 | 超量需申请配额 |
优雅降载方案:
```python
import time
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry
使用urllib3的内置重试策略:自动退避
session = requests.Session()
retry_strategy = Retry(
total=3, # 最多重试3次
backoff_factor=2, # 重试间隔:1×2=2秒, 2×2=4秒, 3×2=6秒
status_forcelist=[429, 500, 502, 503, 504]
)
adapter = HTTPAdapter(max_retries=retry_strategy)
session.mount("http://", adapter)
```
2.5 方案横向对比
| 维度 | 手工导出 | 基础requests | 本文的MESAPIClient | httpx+asyncio |
|------|---------|-------------|-------------------|---------------|
| 开发复杂度 | 0 | ★☆☆☆☆ | ★★☆☆☆ | ★★★★☆ |
| 分页处理 | 无 | 需手动写 | 自动 | 自动 |
| 重试机制 | 无 | 需手动加 | 自动 | 自动 |
| Token管理 | 无 | 需手动处理 | 自动 | 自动 |
| 并发性能 | 0 | 1x(串行) | 1x | 5~10x |
| 适用场景 | 临时一次 | 学习/测试 | 生产级采集 | 超大规模 |
---
三、实战案例:MES API数据采集器(生产级代码)
3.1 需求分析
某FAB工厂的MES API(基于OAuth2 + Bearer Token)需要支撑以下场景:
- 日常采集:每日凌晨自动拉取昨日全量数据(~5000条/天)
- 增量采集:实时获取过去1小时的新数据(轮询间隔60秒)
- 批量回填:首次对接时,一次性拉取近3年的历史数据(~500万条)
- 容错处理:网络抖动、Token过期、API限流都要自动处理
3.2 完整代码(MESAPIClient类,≤80行)
```python
MESAPIClient.py — MES API生产级采集器,≤80行核心代码
import time, json, logging
from datetime import datetime, timedelta
from pathlib import Path
import requests
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry
logging.basicConfig(level=logging.INFO, format='%(asctime)s %(levelname)s %(message)s')
log = logging.getLogger(__name__)
class MESAPIClient:
def __init__(self, base_url: str, client_id: str, client_secret: str):
self.base_url = base_url.rstrip('/')
self.client_id = client_id
self.client_secret = client_secret
self.token = None
self.token_exp = 0
self.session = self._build_session()
def _build_session(self) -> requests.Session:
"""构建带重试策略的HTTP会话:429限流自动退避,5xx服务器错误自动重试"""
s = requests.Session()
retry = Retry(total=5, backoff_factor=3, status_forcelist=[429,500,502,503,504],
allowed_methods=["GET","POST"])
s.mount("http://", HTTPAdapter(max_retries=retry))
s.headers.update({'Content-Type': 'application/json', 'Accept': 'application/json'})
return s
def _ensure_token(self):
"""Token管理:过期前自动刷新。为什么要提前刷新?避免请求到一半Token过期导致数据丢失"""
if self.token and time.time() < self.token_exp - 60: return # 未过期
resp = self.session.post(f'{self.base_url}/oauth/token',
data={'grant_type':'client_credentials','client_id':self.client_id,
'client_secret':self.client_secret})
resp.raise_for_status()
self.token = resp.json()['access_token']
self.token_exp = time.time() + resp.json().get('expires_in', 3600)
self.session.headers.update({'Authorization': f'Bearer {self.token}'})
log.info(f'🔑 Token已刷新,有效期至 {datetime.fromtimestamp(self.token_exp):%H:%M:%S}')
def fetch_page(self, endpoint: str, params: dict) -> dict:
"""单页请求:自动处理Token+重试。为什么单独封装?方便单独调试和日志追踪"""
self._ensure_token()
resp = self.session.get(f'{self.base_url}{endpoint}', params=params, timeout=60)
if resp.status_code == 401:
self.token = None; self._ensure_token() # Token过期,立即刷新重试
resp = self.session.get(f'{self.base_url}{endpoint}', params=params, timeout=60)
resp.raise_for_status()
return resp.json()
def collect_all(self, endpoint: str, date_str: str, process: str = None,
page_size: int = 1000) -> list:
"""
全量采集:自动翻页,直到拿完全部数据。
翻页策略:offset递增(MES API不支持cursor时的标准做法)。
返回:全部数据列表。
"""
all_records, offset = [], 0
params = {"date": date_str, "limit": page_size, "offset": 0}
if process: params["process"] = process
while True:
data = self.fetch_page(endpoint, {**params, "offset": offset})
page = data.get('data', [])
total = data.get('total', len(page) + offset)
all_records.extend(page)
fetched = len(all_records)
log.info(f'[{date_str}] 已获取 {fetched}/{total} 条 (offset={offset})')
if fetched >= total or not page: break # 无新数据时退出,避免死循环
offset += page_size
return all_records
def save_to_file(self, records: list, filename: str):
Path(filename).parent.mkdir(parents=True, exist_ok=True)
with open(filename, 'w', encoding='utf-8') as f:
json.dump(records, f, ensure_ascii=False, indent=2)
log.info(f'💾 已保存 {len(records)} 条数据 → {filename}')
使用示例(模拟数据,无需真实MES)
if __name__ == '__main__':
import random
# 模拟MES API响应(用于本地测试)
def mock_get(url, **kwargs):
class Resp:
status_code = 200
def json(self): params = kwargs.get('params',{}); offset = params.get('offset',0)
records = [{'lot_id':f'FAB{i:04d}','process':random.choice(['PHOTO','ETCH']),
'yield_rate':random.uniform(85,100)} for i in range(10)]
total = 50
return {'data': records[offset:offset+params.get('limit',10)], 'total': total}
def raise_for_status(self): pass
return Resp()
# 实际使用时替换为真实URL和凭证
print('✅ MESAPIClient 初始化完成(请配置真实 base_url/client_id/client_secret)')
```
为什么这样写:
- `_ensure_token()`:Token过期前60秒就自动刷新,避免"请求发出去、Token刚过期"的边界情况。FAB的生产数据丢了就再也补不回来
- `while True + offset`:标准的"请求→判断→继续/退出"翻页模式,比`for`循环更安全(API可能返回0条也返回更多条)
- `self.session`:复用TCP连接,省去每次请求建立连接的开销(~200ms/次),采集5000条数据能省10分钟连接时间
- `HTTPAdapter + Retry`:把重试逻辑从业务代码里分离出来,代码更干净,也更容易调整重试策略
---
四、效果对比:API采集 vs 手工导出
4.1 多维度量化对比
| 对比维度 | 手工导出 | 基础requests | MESAPIClient(生产级) | 说明 |
|---------|---------|-------------|--------------------------|------|
| 单次采集耗时 | 5~15分钟 | 3~10分钟 | 30秒~3分钟 | 含翻页全量数据 |
| 每日操作次数 | 10次(不同工序/日期) | 1次(脚本自动) | 自动定时 | 零人工干预 |
| 数据完整性 | 漏导率~15%(忙忘了/条件选错) | 依赖分页处理 | 99.5%+ | 翻页自动补全 |
| 节假日数据 | 经常漏掉 | 可补采 | 自动兜底 | 历史数据可回填 |
| 网络异常处理 | 重新导出 | 手动重试 | 自动重试+退避 | 7×24运行 |
| Token过期处理 | 不知道这回事 | 手动刷新 | 自动刷新 | 不间断运行 |
| 人工时间成本 | 200+小时/年 | 10小时/年 | 0小时/年 | 完全自动化 |
| 数据可追溯性 | 无(不知道谁何时导出) | 有日志 | 结构化日志+JSON存档 | 可审计 |
4.2 配图:数据采集效率与数据完整性对比
```python
生成 article19 配图脚本
import matplotlib
matplotlib.use('Agg')
import matplotlib.pyplot as plt
import numpy as np
plt.rcParams['font.sans-serif'] = ['SimHei']
plt.rcParams['axes.unicode_minus'] = False
fig, axes = plt.subplots(1, 2, figsize=(14, 5.5))
fig.suptitle('第19篇:API数据采集效率对比(P0级)', fontsize=14, fontweight='bold')
图1:采集耗时对比(柱状图)
methods = ['手工导出\n(10次/天)', '基础requests\n(手动分页)', 'MESAPIClient\n(自动采集)']
daily_mins = [60, 8, 0.5] # 每天耗时(分钟)
colors = ['#EF5350', '#FFA726', '#42A5F5']
bars = axes[0].bar(methods, daily_mins, color=colors, width=0.55,
edgecolor='white', linewidth=1.5)
for bar, m in zip(bars, daily_mins):
axes[0].text(bar.get_x()+bar.get_width()/2, bar.get_height()+1,
f'{m:.1f}分钟', ha='center', va='bottom', fontsize=12, fontweight='bold')
axes[0].set_title('每日数据采集耗时对比', fontsize=12)
axes[0].set_ylabel('耗时(分钟)')
axes[0].set_ylim(0, 75)
axes[0].grid(axis='y', alpha=0.3)
图2:年化时间成本(堆积条形)
tasks = ['数据采集', '格式整理', '数据检查', '补采遗漏', '问题排查']
manual = [50, 15, 10, 15, 10] # 手工每月耗时(分钟)
automated = [5, 1, 1, 1, 2] # 自动化每月耗时(分钟)
x = np.arange(len(tasks)); w = 0.35
axes[1].bar(x-w/2, manual, w, label='手工操作', color='#EF5350')
axes[1].bar(x+w/2, automated, w, label='API自动化', color='#42A5F5')
axes[1].set_xticks(x); axes[1].set_xticklabels(tasks, rotation=15, ha='right')
axes[1].set_title('月度数据处理时间分布(月度,分钟)', fontsize=12)
axes[1].set_ylabel('耗时(分钟)')
axes[1].legend(); axes[1].grid(axis='y', alpha=0.3)
标注年化节省
annual_saved = (sum(manual) - sum(automated)) * 12 / 60
axes[1].annotate(f'年化节省\n~{annual_saved:.0f}小时', xy=(2, 25), xytext=(2.5, 45),
fontsize=11, color='#C62828', fontweight='bold',
arrowprops=dict(arrowstyle='->', color='#C62828'))
plt.tight_layout(rect=[0, 0, 1, 0.95])
plt.savefig('D:/work/CSDN自动发布/半导体Python专栏/images/19_api_time_saving.png', dpi=150, bbox_inches='tight')
plt.savefig('D:/work/CSDN自动发布/半导体Python专栏/images/19_api_workload.png', dpi=150, bbox_inches='tight')
print('✅ 配图已生成')
```
---
五、实施建议:API对接MES的四阶段路径
第一阶段:申请权限+搭建测试环境(第1~3天)
找MES IT申请API权限,提供以下信息:
- 申请目的(数据采集+分析)
- 需要的接口列表(从MES IT文档获取)
- 预计调用频率(日均/峰值)
- 联系方式(异常时IT能联系到你)
切记:不要在生产环境直接测试。申请一个只读测试Token,先在本地用模拟数据验证代码逻辑。
```python
用Mock数据验证代码逻辑(不触真实API)
from unittest.mock import patch
mock_response = {"data": [{"lot_id":"TEST","yield_rate":95.5}],"total":1}
with patch('requests.Session.get', return_value=MockResponse(mock_response)):
client = MESAPIClient("http://test", "id", "secret")
data = client.collect_all("/api/lots", "2026-06-01")
print(f"测试通过: {len(data)}条")
```
风险提示:API Key/Token是访问MES数据的钥匙,绝对不能硬编码在代码里。使用环境变量:
```python
import os
CLIENT_ID = os.environ.get('MES_CLIENT_ID') # 运行时从环境变量读取
CLIENT_SECRET = os.environ.get('MES_CLIENT_SECRET') # 不在代码中出现
```
第二阶段:跑通分页+容错逻辑(第3~7天)
用少量数据(1天的数据)验证完整的采集链路:
```python
Step 1: 单日采集测试
client = MESAPIClient(BASE_URL, CLIENT_ID, CLIENT_SECRET)
data = client.collect_all('/api/v1/lots', '2026-06-01')
print(f'✅ 采集成功: {len(data)}条')
Step 2: 数据校验
assert len(data) > 0, '数据为空,请检查API参数'
assert all('lot_id' in r for r in data), '数据结构不符合预期'
Step 3: 保存到数据库
from DatabaseManager import DatabaseManager
db = DatabaseManager('fab_data.db')
db.batch_insert(data)
db.close()
```
风险提示:MES API通常有每日调用配额(如10000次/天)。大规模数据回填(500万条)如果一口气全采,会触发配额限制。建议分段采集,每天采3万条,分批跑:
```python
import schedule
def incremental_job():
yesterday = (datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d')
data = client.collect_all('/api/v1/lots', yesterday)
db.batch_insert(data)
schedule.every().hour.do(incremental_job) # 每小时增量采集(慎用,看API配额)
```
第三阶段:生产化+定时任务(第7~14天)
```python
fab_api_scheduler.py — 每日自动采集(配合crontab或Windows任务计划程序)
from datetime import datetime, timedelta
from MESAPIClient import MESAPIClient
from DatabaseManager import DatabaseManager
BASE_URL = 'http://mes-api.fab.local'
CLIENT_ID = os.environ['MES_CLIENT_ID']
CLIENT_SECRET = os.environ['MES_CLIENT_SECRET']
def daily_collect():
log.info('📡 开始每日数据采集...')
client = MESAPIClient(BASE_URL, CLIENT_ID, CLIENT_SECRET)
yesterday = (datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d')
db = DatabaseManager('fab_data.db')
for process in ['PHOTO', 'ETCH', 'CVD', 'CMP', 'IMP']:
try:
data = client.collect_all('/api/v1/lots', yesterday, process)
db.batch_insert(data)
log.info(f'✅ {process}: 采集{len(data)}条')
except Exception as e:
log.error(f'❌ {process}采集失败: {e}') # 不阻断其他工序
db.close()
log.info('📡 采集完成')
Windows任务计划程序运行:
pythonw fab_api_scheduler.py
每天凌晨 05:00 执行
```
风险提示:MES API的Schema(数据字段名、类型)可能会在版本升级时变化。建议做Schema校验,用pydantic定义预期结构,一旦API返回的字段不符合预期就告警:
```python
from pydantic import BaseModel
class LotRecord(BaseModel):
lot_id: str
process: str
yield_rate: float
thickness_avg: float | None = None
校验每条记录
for record in data:
try:
LotRecord(**record)
except Exception as e:
log.warning(f'数据Schema异常: {record}, 错误: {e}')
```
第四阶段:监控+告警(第14~21天)
生产级的采集系统必须配套监控:
```python
关键指标监控
def check_health():
yesterday = (datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d')
db = DatabaseManager('fab_data.db')
stats = db.get_summary()
if stats['total_lots'] == 0:
send_alert('⚠️ 昨日MES数据采集失败,请检查!') # 接入企业微信/钉钉
expected_min = 3000 # 正常情况下每天至少3000条
if stats.get('last_24h_count', 0) < expected_min:
send_alert(f'⚠️ 数据量异常偏低: {stats["last_24h_count"]}条 < {expected_min}条')
```
---
六、进阶方向:从同步采集到实时数据流
1. 异步采集:aiohttp提升5~10倍吞吐量
当前代码是串行采集(一个请求完再发下一个),对于大规模多工序数据,可以用`aiohttp`做并发请求:
```python
import aiohttp, asyncio
async def async_collect(session, endpoint, dates, processes):
"""异步并发采集:同时发多个请求,吞吐量提升5~10倍"""
tasks = [
fetch_with_session(session, endpoint, date, proc)
for date in dates for proc in processes
]
results = await asyncio.gather(*tasks, return_exceptions=True)
return [r for r in results if isinstance(r, list)]
适合场景:首次回填3年历史数据(500万条),从10小时压缩到1小时
```
2. WebSocket:设备状态的实时推送
某些MES支持WebSocket推送——设备报警、Lot状态变更等事件,服务器主动推给客户端,无需轮询:
```python
import websockets, asyncio, json
async def real_time_monitor():
uri = "ws://mes-api.fab.local/ws/alerts"
async with websockets.connect(uri) as ws:
async for message in ws:
alert = json.loads(message)
if alert['type'] == 'LOT_ABORT':
send_immediate_alert(f"🚨 Lot {alert['lot_id']} 异常终止!")
```
局限性:WebSocket需要MES服务器支持,且工厂网络环境复杂(防火墙、工业网络隔离),实施难度比HTTP轮询高很多。建议优先用HTTP轮询(5分钟间隔已经足够实时),WebSocket留给真正需要秒级响应的场景(如设备急停告警)。
3. 任务队列:Celery分布式采集
当采集任务变多(多工厂、多系统、定时+事件触发),需要任务队列来管理:
```python
Celery + Redis:分布式任务调度
from celery import Celery
app = Celery('fab_collector', broker='redis://localhost:6379')
@app.task
def collect_mes_data(date: str, process: str):
"""Celery任务:支持失败重试、任务链、定时调度"""
client = MESAPIClient(BASE_URL, CLIENT_ID, CLIENT_SECRET)
data = client.collect_all('/api/v1/lots', date, process)
db.batch_insert(data)
return {'date': date, 'process': process, 'count': len(data)}
定时任务(Celery Beat)
@app.task
def daily_schedule():
yesterday = (datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d')
for process in ['PHOTO','ETCH','CVD','CMP','IMP']:
collect_mes_data.delay(yesterday, process) # 异步分发6个任务
celery -A fab_collector beat --loglevel=info # 定时调度器
celery -A fab_collector worker --loglevel=info # 任务执行器
```
4. 数据血缘追踪:采集元数据管理
当采集的数据被多个分析系统使用时,需要追踪"这条数据是哪个API接口、什么时间采集的、采集时的参数是什么"——这就是数据血缘(Data Lineage)。
```python
class DataLineage:
"""记录每条原始数据的来源,用于问题追溯"""
def __init__(self):
self.conn = sqlite3.connect('data_lineage.db')
self.conn.execute('''
CREATE TABLE IF NOT EXISTS lineage (
record_id TEXT,
source_api TEXT,
collected_at TEXT,
params TEXT,
PRIMARY KEY(record_id, collected_at)
)''')
def log(self, record_id, api, params):
self.conn.execute(
'INSERT OR REPLACE INTO lineage VALUES (?,?,?,?)',
(record_id, api, datetime.now().isoformat(), json.dumps(params))
)
self.conn.commit()
```
---
> 💬 你的MES系统有API接口吗?你是怎么解决数据采集问题的?踩过什么坑?欢迎评论区分享!
>
> 📦 专栏VIP资源包(P0级)(含MESAPIClient生产级完整源码+Token刷新逻辑+定时任务模板+Schema校验脚本)在专栏主页点击「VIP资源」获取。
>
> 📚 半导体Python实战专栏·第19篇(P0级重磅),觉得有用请收藏+点赞支持~ 👍
>
> 🔧 下一篇预告:数据预处理——数据拿到了但格式五花八门怎么办?第20篇教你从脏数据到分析就绪的全流程!





