快递批量查询工具的技术实现:从API对接到并发查询的完整方案
一、为什么需要批量查询?场景与需求分析
在日常电商运营和物流管理工作中,快递单号查询是一个高频且繁琐的操作。对于日均发货几百甚至上千单的商家来说,逐一手动复制单号、打开快递查询页面、粘贴查询、记录状态——这一连串操作每天要重复上百次。
核心痛点:
- 效率低下:单量越大,重复劳动的时间成本越高
- 容易出错:手动复制粘贴容易遗漏或错位
- 异常难以及时发现:物流滞留、派送失败等问题无法第一时间掌握
功能需求清单:
| 功能模块 | 需求描述 |
|---|---|
| 单号批量导入 | 支持批量粘贴、文件导入(TXT/CSV/Excel) |
| 快递公司自动识别 | 根据单号格式自动判断所属快递公司 |
| 批量并发查询 | 同时查询数百上千个单号,而非逐个串行 |
| 结果展示与筛选 | 清晰展示物流状态,支持按状态/快递公司筛选 |
| 数据导出 | 支持CSV、Excel等格式导出,用于对账和分析 |
目前市面上已有成熟的解决方案,例如卢米快递查询助手,它将上述所有功能封装为开箱即用的桌面工具,覆盖国内外千余家快递公司,支持不限单量的批量查询与一键导出。如果不想从头开发,直接使用这类工具是最高效的选择。但如果想深入理解技术原理,或者需要将物流查询能力集成到自有系统中,下面的内容会带你从零开始搭建一套完整的批量查询系统。
二、整体技术架构
在动手写代码之前,先明确系统的整体架构。
技术选型:
| 组件 | 选型 | 理由 |
|---|---|---|
| 编程语言 | Python 3.8+ | 生态丰富,API调用方便,适合快速开发 |
| HTTP客户端 | requests / aiohttp | 同步/异步两种场景都支持 |
| 并发框架 | asyncio + aiohttp | 实现高并发查询的核心 |
| 数据处理 | pandas | 方便数据清洗、筛选和导出 |
| 数据存储 | SQLite(可选) | 轻量级本地存储,无需额外服务 |
| 桌面界面 | PyQt / Tkinter | 如需GUI可选用 |
系统逻辑架构:
[单号输入] → [快递识别引擎] → [批量查询调度器] → [API并发调用] → [结果聚合] → [展示/导出]
↓ ↓
[规则库匹配] [限流与重试]
[缓存查询] [超时控制]
三、快递公司自动识别:从规则到算法
批量查询的第一步,是判断每个单号属于哪家快递公司。这是整个系统的基础。
3.1 单号编码规则分析
不同快递公司的单号有各自的编码规律:
| 快递公司 | 单号特征 | 示例 |
|---|---|---|
| 顺丰速运 | 12/15位,常以SF开头 | SF123456789012 |
| 中通快递 | 12位纯数字 | 751234567890 |
| 圆通速递 | 10/12位,常以YT开头 | YT1234567890 |
| 韵达快递 | 13位数字 | 1201234567890 |
| EMS | 字母数字混合 | EE123456789CN |
| 京东快递 | 以JD开头 | JD1234567890 |
| 极兔速递 | 12-14位,常以JT开头 | JT123456789012 |
| 申通快递 | 12位纯数字 | 888123456789 |
3.2 基于正则表达式的规则引擎
最直接的方式是用正则表达式构建规则库:
import re
# 快递公司识别规则库
EXPRESS_RULES = [
("顺丰速运", r'^(SF|SFL)\d{12,15}$'),
("中通快递", r'^\d{12}$'), # 注意:与申通格式重叠
("圆通速递", r'^(YT|YTO)\d{10,12}$'),
("韵达快递", r'^[1-9]\d{12}$'),
("EMS", r'^[A-Z]{2}\d{9,11}[A-Z]{2}$'),
("京东快递", r'^JD[A-Z0-9]{10,12}$'),
("极兔速递", r'^JT\d{12,14}$'),
("申通快递", r'^\d{12}$'), # 与中通格式重叠
("德邦快递", r'^DB\d{10,12}$'),
("百世快递", r'^\d{12}$'),
]
def identify_express_company(tracking_number: str) -> str:
"""根据单号识别快递公司"""
tracking_number = tracking_number.strip().upper()
for company, pattern in EXPRESS_RULES:
if re.match(pattern, tracking_number):
return company
return "未知"
3.3 处理规则冲突与歧义
纯规则匹配存在一个典型问题:格式重叠。中通和申通的单号都是12位纯数字,规则上无法区分。
解决方案是多级识别策略:
def identify_with_confidence(tracking_number: str) -> dict:
"""
带置信度的多级识别
1. 规则匹配 → 2. 前缀特征 → 3. 历史数据辅助
"""
tracking_number = tracking_number.strip().upper()
# 第一级:精确规则匹配
candidates = []
for company, pattern in EXPRESS_RULES:
if re.match(pattern, tracking_number):
confidence = 0.85
# 前缀特征加分
prefix_boost = {
"SF": ("顺丰速运", 0.98),
"JT": ("极兔速递", 0.95),
"JD": ("京东快递", 0.95),
"YT": ("圆通速递", 0.92),
"DB": ("德邦快递", 0.92),
}
for prefix, (name, score) in prefix_boost.items():
if tracking_number.startswith(prefix):
if name == company:
confidence = score
break
candidates.append({"company": company, "confidence": confidence})
if not candidates:
return {"company": "未知", "confidence": 0}
# 按置信度排序
candidates.sort(key=lambda x: x["confidence"], reverse=True)
# 如果最高置信度 > 阈值,返回
if candidates[0]["confidence"] > 0.8:
return candidates[0]
# 置信度不足时,可查询历史数据辅助判断
# history_result = query_user_history(tracking_number)
# if history_result:
# return {"company": history_result, "confidence": 0.7}
return {"company": "需人工确认", "confidence": candidates[0]["confidence"]}
3.4 输入规范化与纠错
实际使用中,用户输入的单号可能包含空格、连字符等特殊字符,甚至OCR识别结果可能有字符混淆:
def normalize_tracking_number(raw: str) -> str:
"""
规范化单号输入
- 去除空格、连字符、下划线
- 统一大小写
- 纠正常见OCR错误(7/B, 0/O等)
"""
# 去除特殊字符
cleaned = re.sub(r'[\s\-_\.]', '', raw.strip())
# 字母统一大写
cleaned = cleaned.upper()
# OCR常见混淆纠错
ocr_fixes = {
'B': '8', # B可能被识别为8
'O': '0', # O可能被识别为0
'I': '1', # I可能被识别为1
'Z': '2', # Z可能被识别为2
}
# 注意:仅对纯数字段做纠错,保留字母段
# 实际实现需要更精细的规则
return cleaned
四、物流API对接:核心查询能力
识别出快递公司之后,下一步是通过API获取物流信息。
4.1 API服务商选择
目前国内主流的物流API聚合服务商主要有快递鸟和快递100两家。它们已经集成了数百家快递公司的接口,一次对接即可查询绝大多数快递,无需逐一对接各家快递公司。
快递鸟API特点:
- 支持2500+家快递公司
- 提供单条查询、批量查询(单次最多100个单号)
- 支持在途监控、轨迹订阅推送
- 请求需要数据签名认证
快递100 API特点:
- 覆盖2100+家快递公司
- 支持自动识别快递公司
- 提供实时查询接口
4.2 API调用基础:签名与请求
以快递鸟API为例,展示完整的单号查询流程:
import requests
import json
import hashlib
import base64
class ExpressAPI:
"""快递鸟API封装类"""
def __init__(self, ebid: str, api_key: str):
self.ebid = ebid
self.api_key = api_key
self.api_url = "https://api.kdniao.com/Ebusiness/EbusinessOrderHandle.aspx"
def _generate_sign(self, request_data: str) -> str:
"""生成数据签名(防篡改)"""
raw_sign = request_data + self.api_key
md5_sign = hashlib.md5(raw_sign.encode('utf-8')).hexdigest()
return base64.b64encode(md5_sign.encode('utf-8')).decode()
def query_single(self, shipper_code: str, logistic_code: str) -> dict:
"""
查询单个快递
shipper_code: 快递公司编码(如SF=顺丰)
logistic_code: 快递单号
"""
# 1. 构造请求数据
request_data = {
"ShipperCode": shipper_code,
"LogisticCode": logistic_code
}
json_data = json.dumps(request_data, ensure_ascii=False)
# 2. 构造完整请求参数
post_data = {
"RequestData": json_data,
"EBusinessID": self.ebid,
"RequestType": "1002", # 1002=即时查询
"DataSign": self._generate_sign(json_data),
"DataType": "2" # 返回JSON格式
}
# 3. 发送请求
try:
response = requests.post(self.api_url, data=post_data, timeout=10)
result = response.json()
if result.get("Success"):
return {
"success": True,
"number": logistic_code,
"status_code": result.get("State"),
"status_text": self._status_code_to_text(result.get("State")),
"traces": result.get("Traces", []),
"company": result.get("ShipperCode")
}
else:
return {
"success": False,
"number": logistic_code,
"error": result.get("Reason", "查询失败")
}
except requests.exceptions.Timeout:
return {
"success": False,
"number": logistic_code,
"error": "请求超时"
}
except Exception as e:
return {
"success": False,
"number": logistic_code,
"error": str(e)
}
def _status_code_to_text(self, code) -> str:
"""状态码转文字"""
status_map = {
0: "无轨迹",
1: "已揽收",
2: "运输中",
3: "已签收",
4: "问题件",
5: "已退件"
}
return status_map.get(code, "未知状态")
4.3 快递公司编码映射
各API服务商的快递公司编码可能不同,需要建立映射表:
# 快递公司名称 → 快递鸟编码
COMPANY_CODE_MAP = {
"顺丰速运": "SF",
"中通快递": "ZTO",
"圆通速递": "YTO",
"韵达快递": "YD",
"申通快递": "STO",
"EMS": "EMS",
"京东快递": "JD",
"极兔速递": "JT",
"德邦快递": "DB",
"百世快递": "BST",
}
def get_company_code(company_name: str) -> str:
"""获取快递公司API编码"""
return COMPANY_CODE_MAP.get(company_name, "")
五、批量查询引擎:从串行到并发
单号查询的API调好了,现在要实现核心功能——批量查询。
5.1 同步批量查询(基础版)
最简单的实现是循环调用,但这种方式串行执行,速度慢:
def batch_query_sync(api: ExpressAPI, tracking_list: list) -> list:
"""
同步批量查询(串行执行)
适用场景:单量少(<50个),对速度要求不高
"""
results = []
for item in tracking_list:
result = api.query_single(
shipper_code=item["company_code"],
logistic_code=item["number"]
)
results.append({
"number": item["number"],
"company": item["company_name"],
"result": result
})
return results
问题:查询100个单号,如果每个耗时0.5秒,总共需要50秒。对于电商大促场景动辄数千单的情况,这种方式完全不可用。
5.2 异步并发批量查询(进阶版)
用asyncio和aiohttp实现并发查询,将总时间压缩到单次请求的时间量级:
import asyncio
import aiohttp
import json
import hashlib
import base64
class AsyncExpressAPI:
"""异步快递鸟API封装"""
def __init__(self, ebid: str, api_key: str):
self.ebid = ebid
self.api_key = api_key
self.api_url = "https://api.kdniao.com/Ebusiness/EbusinessOrderHandle.aspx"
def _generate_sign(self, request_data: str) -> str:
raw_sign = request_data + self.api_key
md5_sign = hashlib.md5(raw_sign.encode('utf-8')).hexdigest()
return base64.b64encode(md5_sign.encode('utf-8')).decode()
async def query_one(self, session: aiohttp.ClientSession,
shipper_code: str, logistic_code: str) -> dict:
"""异步查询单个快递"""
request_data = {
"ShipperCode": shipper_code,
"LogisticCode": logistic_code
}
json_data = json.dumps(request_data, ensure_ascii=False)
post_data = {
"RequestData": json_data,
"EBusinessID": self.ebid,
"RequestType": "1002",
"DataSign": self._generate_sign(json_data),
"DataType": "2"
}
try:
async with session.post(self.api_url, data=post_data, timeout=10) as resp:
result = await resp.json()
if result.get("Success"):
return {
"success": True,
"number": logistic_code,
"status_code": result.get("State"),
"traces": result.get("Traces", [])
}
else:
return {
"success": False,
"number": logistic_code,
"error": result.get("Reason", "查询失败")
}
except asyncio.TimeoutError:
return {
"success": False,
"number": logistic_code,
"error": "请求超时"
}
except Exception as e:
return {
"success": False,
"number": logistic_code,
"error": str(e)
}
async def batch_query_async(api: AsyncExpressAPI,
tracking_list: list,
concurrency: int = 10) -> list:
"""
异步批量查询(并发执行)
tracking_list: [{"company_code": "SF", "number": "123"}]
concurrency: 同时并发数(避免API限流)
"""
# 使用Semaphore控制并发量
semaphore = asyncio.Semaphore(concurrency)
async def query_with_limit(session, item):
async with semaphore:
return await api.query_one(
session,
item["company_code"],
item["number"]
)
async with aiohttp.ClientSession() as session:
tasks = [query_with_limit(session, item) for item in tracking_list]
return await asyncio.gather(*tasks)
# 使用示例
async def main():
api = AsyncExpressAPI("your_ebid", "your_api_key")
tracking_list = [
{"company_code": "SF", "number": "SF123456789012"},
{"company_code": "ZTO", "number": "751234567890"},
# ... 更多单号
]
# 并发数设为10,避免触发API限流
results = await batch_query_async(api, tracking_list, concurrency=10)
for result in results:
print(result)
# asyncio.run(main())
5.3 并发控制与限流策略
API服务商通常有QPS限制(每秒请求数),需要合理控制并发量:
import time
from collections import deque
class RateLimiter:
"""滑动窗口限流器"""
def __init__(self, max_requests: int, window_seconds: int):
self.max_requests = max_requests
self.window_seconds = window_seconds
self.requests = deque()
def wait_if_needed(self):
"""如果达到限流阈值,等待"""
now = time.time()
# 清理窗口外的记录
while self.requests and now - self.requests[0] > self.window_seconds:
self.requests.popleft()
# 如果达到上限,等待
if len(self.requests) >= self.max_requests:
sleep_time = self.window_seconds - (now - self.requests[0])
if sleep_time > 0:
time.sleep(sleep_time + 0.1)
self.requests.append(now)
# 在同步API调用中使用
limiter = RateLimiter(max_requests=20, window_seconds=1) # 每秒最多20次
def query_with_rate_limit(api, shipper_code, logistic_code):
limiter.wait_if_needed()
return api.query_single(shipper_code, logistic_code)
5.4 结果缓存:减少重复API调用
对于短时间内重复查询的单号,增加缓存机制可以大幅减少API调用量:
import hashlib
import json
import sqlite3
from datetime import datetime, timedelta
class CacheManager:
"""SQLite本地缓存管理器"""
def __init__(self, db_path: str = "express_cache.db", ttl_seconds: int = 300):
self.db_path = db_path
self.ttl = ttl_seconds
self._init_db()
def _init_db(self):
"""初始化数据库表"""
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
cursor.execute('''
CREATE TABLE IF NOT EXISTS cache (
tracking_number TEXT PRIMARY KEY,
result TEXT,
updated_at TIMESTAMP
)
''')
conn.commit()
conn.close()
def get(self, tracking_number: str) -> dict:
"""获取缓存结果"""
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
cursor.execute(
"SELECT result, updated_at FROM cache WHERE tracking_number = ?",
(tracking_number,)
)
row = cursor.fetchone()
conn.close()
if row:
result_json, updated_at = row
updated_time = datetime.fromisoformat(updated_at)
if datetime.now() - updated_time < timedelta(seconds=self.ttl):
return json.loads(result_json)
return None
def set(self, tracking_number: str, result: dict):
"""存入缓存"""
conn = sqlite3.connect(self.db_path)
cursor = conn.cursor()
cursor.execute(
"INSERT OR REPLACE INTO cache (tracking_number, result, updated_at) VALUES (?, ?, ?)",
(tracking_number, json.dumps(result), datetime.now().isoformat())
)
conn.commit()
conn.close()
六、数据处理与导出
6.1 使用Pandas处理批量数据
import pandas as pd
def convert_to_dataframe(raw_results: list) -> pd.DataFrame:
"""将查询结果转换为DataFrame"""
rows = []
for result in raw_results:
# 提取最新轨迹
traces = result.get("traces", [])
latest_trace = traces[-1] if traces else {}
row = {
"快递单号": result.get("number", ""),
"快递公司": result.get("company", ""),
"物流状态": result.get("status_text", ""),
"最新轨迹": latest_trace.get("AcceptStation", ""),
"更新时间": latest_trace.get("AcceptTime", ""),
"查询状态": "成功" if result.get("success") else "失败"
}
rows.append(row)
df = pd.DataFrame(rows)
# 添加辅助列:是否异常
df["是否异常"] = df["物流状态"].apply(
lambda x: x in ["问题件", "已退件"]
)
return df
def filter_abnormal(df: pd.DataFrame) -> pd.DataFrame:
"""筛选异常订单"""
return df[df["是否异常"] == True]
def group_by_company(df: pd.DataFrame) -> pd.DataFrame:
"""按快递公司分组统计"""
return df.groupby("快递公司").agg({
"快递单号": "count",
"是否异常": "sum"
}).rename(columns={
"快递单号": "总单量",
"是否异常": "异常件数"
})
6.2 多格式导出
import csv
import json
def export_to_csv(data: list, filepath: str):
"""导出为CSV"""
if not data:
return
fieldnames = data[0].keys()
with open(filepath, 'w', newline='', encoding='utf-8-sig') as f:
writer = csv.DictWriter(f, fieldnames=fieldnames)
writer.writeheader()
writer.writerows(data)
def export_to_excel(data: list, filepath: str):
"""导出为Excel(需安装openpyxl)"""
df = pd.DataFrame(data)
df.to_excel(filepath, index=False, engine='openpyxl')
def export_to_json(data: list, filepath: str):
"""导出为JSON"""
with open(filepath, 'w', encoding='utf-8') as f:
json.dump(data, f, ensure_ascii=False, indent=2)
七、完整命令行工具
将以上模块整合,实现一个可用的命令行批量查询工具:
#!/usr/bin/env python
# -*- coding: utf-8 -*-
import argparse
import asyncio
import sys
from typing import List, Dict
class BatchQueryTool:
"""批量查询工具主类"""
def __init__(self, ebid: str, api_key: str):
self.api = AsyncExpressAPI(ebid, api_key)
self.cache = CacheManager()
def load_numbers(self, input_file: str) -> List[Dict]:
"""从文件加载单号并自动识别快递公司"""
import pandas as pd
if input_file.endswith('.csv'):
df = pd.read_csv(input_file)
numbers = df.iloc[:, 0].dropna().tolist()
elif input_file.endswith(('.xlsx', '.xls')):
df = pd.read_excel(input_file)
numbers = df.iloc[:, 0].dropna().tolist()
else:
with open(input_file, 'r', encoding='utf-8') as f:
numbers = [line.strip() for line in f if line.strip()]
tracking_list = []
for num in numbers:
num = normalize_tracking_number(str(num))
company = identify_express_company(num)
tracking_list.append({
"number": num,
"company_name": company,
"company_code": get_company_code(company)
})
return tracking_list
async def run(self, input_file: str, output_file: str,
concurrency: int = 10, export_format: str = "csv"):
"""运行批量查询"""
print(f"📦 加载单号:{input_file}")
tracking_list = self.load_numbers(input_file)
print(f"✅ 共加载 {len(tracking_list)} 个单号")
# 过滤掉未识别快递公司的单号
unknown = [t for t in tracking_list if not t["company_code"]]
if unknown:
print(f"⚠️ {len(unknown)} 个单号未能识别快递公司,将跳过")
tracking_list = [t for t in tracking_list if t["company_code"]]
print(f"🔄 开始批量查询(并发数:{concurrency})...")
results = await batch_query_async(self.api, tracking_list, concurrency)
success_count = sum(1 for r in results if r.get("success"))
print(f"✅ 查询完成:成功 {success_count},失败 {len(results) - success_count}")
# 转换为DataFrame并导出
df = convert_to_dataframe(results)
if export_format == "csv":
df.to_csv(output_file, index=False, encoding='utf-8-sig')
elif export_format == "excel":
df.to_excel(output_file, index=False, engine='openpyxl')
elif export_format == "json":
df.to_json(output_file, orient='records', force_ascii=False, indent=2)
print(f"💾 结果已导出:{output_file}")
print(f"📊 异常件数:{len(filter_abnormal(df))}")
def main():
parser = argparse.ArgumentParser(description="快递批量查询工具")
parser.add_argument("input", help="输入文件路径(CSV/Excel/TXT)")
parser.add_argument("-o", "--output", default="result.csv", help="输出文件路径")
parser.add_argument("-c", "--concurrency", type=int, default=10, help="并发数")
parser.add_argument("-f", "--format", choices=["csv", "excel", "json"],
default="csv", help="导出格式")
parser.add_argument("--ebid", required=True, help="快递鸟商户ID")
parser.add_argument("--api-key", required=True, help="快递鸟API密钥")
args = parser.parse_args()
tool = BatchQueryTool(args.ebid, args.api_key)
asyncio.run(tool.run(args.input, args.output, args.concurrency, args.format))
if __name__ == "__main__":
main()
八、从自建到成熟方案:技术实现的边界
通过上述代码,我们完成了一套完整的批量查询系统原型。从快递公司自动识别、API签名认证、异步并发查询,到结果缓存、数据筛选与多格式导出,核心功能均已覆盖。
这套方案适合哪些场景?
- 需要将物流查询能力集成到自有ERP/OMS系统
- 有技术团队维护,需要定制化功能
- 单量极大,需要深度优化查询性能
自建方案需要考虑的成本:
- API服务商有调用次数限制,大促期间可能触发限流
- 快递公司单号规则会变化,规则库需要持续维护
- 不同API服务商的接口格式不同,切换成本高
- 需要自行处理异常重试、超时、数据一致性等问题
如果不想维护这套复杂的技术栈,市面上已有成熟的解决方案。例如卢米快递查询助手,它将本文讲述的所有技术——快递公司自动识别、批量并发查询、数据筛选与多格式导出——全部封装为开箱即用的桌面工具,覆盖国内外千余家快递公司,且支持不限单量的批量查询,无需关心API限流、签名算法、规则库维护等技术细节。
选择自建还是使用成熟工具,取决于你的技术团队配置和业务需求。理解技术原理后,无论走哪条路都能做出更明智的决策。
九、总结
本文从零开始搭建了一套快递批量查询系统,涵盖了以下核心技术模块:
- 快递公司自动识别:基于正则表达式规则库,结合前缀特征和置信度评分
- 物流API对接:以快递鸟为例,实现了请求签名、状态码映射和异常处理
- 批量并发查询:使用asyncio + aiohttp实现高并发,配合Semaphore控制限流
- 结果缓存:基于SQLite的本地缓存,减少重复API调用
- 数据筛选与导出:支持CSV、Excel、JSON多格式导出
这套代码可以直接作为原型使用,也可以作为集成到自有系统的参考实现。
更多推荐




所有评论(0)