共计 1973 个字符,预计需要花费 5 分钟才能阅读完成。
背景与痛点
天气数据服务看似简单,实则面临三大技术挑战:

- 实时性要求:气象数据每分钟都在变化,过期的数据毫无价值
- 数据源异构性:不同供应商返回 JSON/XML 格式不一,单位制式各异(华氏 / 摄氏)
- 高并发访问:极端天气时请求量可能暴增 100 倍
传统直接调用第三方 API 的方案存在明显缺陷:
- API 调用次数限制(通常 500 次 / 天)
- 响应延迟不可控(依赖外网质量)
- 无法定制数据清洗逻辑
技术选型对比
| 方案 | 优点 | 缺点 |
|---|---|---|
| 直接调用第三方 API | 开发简单 | 受配额限制,不可靠 |
| 自建数据管道 | 数据可控,性能优化空间大 | 运维成本较高 |
我们选择自建方案,核心优势在于:
- 数据本地缓存避免重复调用
- 自定义清洗逻辑保证数据一致性
- 分布式扩展能力
核心架构设计
数据采集 Agent
采用生产者 - 消费者模式,关键设计点:
- 智能重试机制:
- 指数退避重试(1s, 2s, 4s…)
-
网络异常自动切换备用数据源
-
数据校验:
- 字段完整性检查(温度 / 湿度必填)
- 数值范围校验(湿度 0 -100%)
@retry(stop_max_attempt_number=3, wait_exponential_multiplier=1000)
def fetch_weather_data(api_url: str) -> dict:
try:
resp = requests.get(api_url, timeout=5)
resp.raise_for_status()
if not validate_data(resp.json()):
raise ValueError("Invalid data structure")
return resp.json()
except Exception as e:
log.error(f"Fetch failed: {str(e)}")
raise
数据清洗流水线
使用 Pandas 构建标准化处理流程:
- 异常值处理:
- 3σ 原则剔除离群点
-
线性插值补全缺失值
-
单位统一化:
- 温度统一转摄氏
- 风速转 m / s 单位
def clean_data(raw_df: pd.DataFrame) -> pd.DataFrame:
# 华氏转摄氏
if 'temp_f' in raw_df.columns:
raw_df['temp_c'] = (raw_df['temp_f'] - 32) * 5/9
# 处理风速单位
if 'wind_mph' in raw_df.columns:
raw_df['wind_mps'] = raw_df['wind_mph'] * 0.44704
return raw_df.dropna()
存储方案
选用 TimescaleDB 时序数据库,配置策略:
- 按城市分片(北京 / 上海等单独分片)
- 自动过期策略(保留最近 30 天数据)
- 压缩比设置(节约 80% 存储空间)
完整代码实现
FastAPI 接口层
内置两级缓存:
- 内存缓存(LRU 策略)
- Redis 分布式缓存
@app.get("/weather/{city}")
@cache(expire=300) # 5 分钟缓存
async def get_weather(city: str):
# 先查本地缓存
if data := local_cache.get(city):
return data
# 查数据库
query = "SELECT * FROM weather_data WHERE city = %s ORDER BY ts DESC LIMIT 1"
record = await database.fetch_one(query, (city,))
if not record:
raise HTTPException(404, "Data not found")
# 写入缓存
local_cache[city] = record
return record
性能优化实战
压测对比(JMeter 测试)
| 场景 | QPS | 平均延迟 | 错误率 |
|---|---|---|---|
| 直接调用 API | 58 | 210ms | 12% |
| 本地缓存方案 | 4200 | 8ms | 0% |
内存优化技巧
- 使用 Pandas 的
category类型存储城市名(减少 75% 内存) - 禁用 PyArrow 的线程池(降低上下文切换开销)
- 分片加载历史数据
避坑指南
- API 配额管理:
- 使用
token-bucket算法限流 -
关键监控指标:
- 每日调用量
- 错误码 429 出现频率
-
时区陷阱:
- 统一存储 UTC 时间
-
前端按用户时区转换显示
-
缓存雪崩预防:
- 过期时间添加随机抖动(±10%)
- 永不过期的基准数据
扩展为分布式架构
未来可升级为数据网格架构:
- 每个城市部署独立 Agent
- 通过 gRPC 实现节点间通信
- 使用 Consul 实现服务发现
graph TD
A[用户请求] --> B(负载均衡器)
B --> C[北京节点]
B --> D[上海节点]
C --> E[本地数据库]
D --> F[本地数据库]
总结心得
构建天气服务就像做一道红烧肉:
- 好的数据源是新鲜五花肉(基础)
- 清洗流程像焯水去腥(必要步骤)
- 缓存机制如同小火慢炖(提升口感)
建议从小规模单节点开始,逐步验证各项技术假设,再扩展为分布式系统。完整代码已上传 GitHub,包含 Docker 部署脚本和性能测试工具链。
正文完
