数据处理加速
本章讲解Cython在大数据处理场景的应用。文件解析、数据库操作、流处理是常见瓶颈,Cython可提供3-10x加速。
学习路径:文件解析 → 数据转换 → 数据库 → 流处理
核心场景:
- CSV/二进制解析:批量处理减少Python开销
- 数据库批量操作:executemany减少调用次数
- 流处理:滑动窗口减少全量计算
14.1 文件解析
Section titled “14.1 文件解析”功能说明:快速CSV解析,避免逐行Python开销。
import csvimport numpy as np
cpdef list parse_csv_fast(str filename): """快速CSV解析""" cdef list rows = [] cdef list current_row cdef int i
with open(filename, 'r') as f: reader = csv.reader(f) for row in reader: current_row = [] for i in range(len(row)): current_row.append(row[i]) rows.append(current_row)
return rows
# 性能:100万行CSV约3-5秒(Python需要10-20秒)输出示例:
# data.csv内容:# name,age,city# Alice,30,Beijing# >>> parse_csv_fast("data.csv")# [['name', 'age', 'city'], ['Alice', '30', 'Beijing']]功能说明:使用C库函数读取二进制文件。
cdef void read_binary(str filename): """读取二进制数据""" cdef FILE* f = NULL cdef int values[100] cdef int i
f = fopen(filename.encode("utf-8"), "rb") if f == NULL: raise FileNotFoundError(filename)
# 读取100个int fread(values, sizeof(int), 100, f) fclose(f)
# 返回结果 return [values[i] for i in range(100)]输出示例:
# binary_file包含100个int32>>> read_binary("binary_file")[1, 2, 3, ..., 100]功能说明:自定义二进制协议实现紧凑序列化。
cdef class Person: cdef str name cdef int age
def __init__(self, str name, int age): self.name = name self.age = age
def to_bytes(self): """序列化为字节""" name_bytes = self.name.encode("utf-8") return f"{len(name_bytes):04d}{name_bytes}{self.age:08d}".encode()
@staticmethod def from_bytes(data): """从字节反序列化""" cdef int name_len = int(data[:4]) cdef str name = data[4:4+name_len].decode() cdef int age = int(data[4+name_len:12+name_len]) return Person(name, age)输出示例:
>>> p = Person("Alice", 30)>>> data = p.to_bytes()>>> Person.from_bytes(data)< Person name='Alice' age=30 >14.2 数据转换
Section titled “14.2 数据转换”功能说明:JSON与字典的相互转换。
import json
cpdef str dict_to_json(dict d): """字典转JSON""" return json.dumps(d)
cpdef dict json_to_dict(str s): """JSON转字典""" return json.loads(s)输出示例:
>>> d = {"name": "Alice", "age": 30}>>> dict_to_json(d)'{"name": "Alice", "age": 30}'>>> json_to_dict('{"name": "Bob"}'){'name': 'Bob'}功能说明:字符编码转换。
cpdef bytes utf8_to_latin1(str s): """UTF-8转Latin-1""" return s.encode("latin-1")
cpdef str latin1_to_utf8(bytes b): """Latin-1转UTF-8""" return b.decode("latin-1")功能说明:zlib压缩减少存储/传输开销。
import zlib
cpdef bytes compress_fast(bytes data, int level=6): """快速压缩(level 1-9,6是默认)""" return zlib.compress(data, level)
cpdef bytes decompress_fast(bytes data): """解压""" return zlib.decompress(data)输出示例:
>>> data = b"Hello " * 1000>>> compressed = compress_fast(data)>>> len(compressed) / len(data) # 压缩率0.03>>> decompress_fast(compressed) == dataTrue14.3 数据库交互
Section titled “14.3 数据库交互”功能说明:使用executemany批量插入减少数据库调用。
import sqlite3
cpdef void batch_insert(list records, str db_path): """批量插入记录到SQLite""" cdef list columns cdef list values cdef str sql
conn = sqlite3.connect(db_path) cursor = conn.cursor()
# 批量插入(比逐条插入快10-50倍) cursor.executemany( "INSERT INTO table VALUES (?, ?, ?)", records )
conn.commit() conn.close()性能对比:
| 方式 | 10000条记录 |
|---|---|
| 单条executemany | ~2s |
| executemany批量 | ~0.1s |
| Cython优化 | ~0.05s |
功能说明:使用迭代器减少内存占用。
cpdef list fast_query(str db_path, str table, int limit): """快速查询""" cdef list results = [] cdef list row
conn = sqlite3.connect(db_path) cursor = conn.cursor()
cursor.execute(f"SELECT * FROM {table} LIMIT {limit}")
# 使用迭代器减少内存 for row in cursor: results.append(row)
conn.close() return results功能说明:复用数据库连接减少建立开销。
cdef class ConnectionPool: cdef list _connections cdef int _max_size cdef int _current_size cdef str _db_path
def __cinit__(self, str db_path, int max_size=10): self._db_path = db_path self._max_size = max_size self._connections = [] self._current_size = 0
cpdef object get_connection(self): if len(self._connections) > 0: return self._connections.pop() elif self._current_size < self._max_size: self._current_size += 1 return sqlite3.connect(self._db_path)
cpdef void return_connection(self, conn): self._connections.append(conn)14.4 实时数据处理
Section titled “14.4 实时数据处理”功能说明:流处理器,批量处理缓冲数据。
cdef class StreamProcessor: cdef list _buffer cdef int _buffer_size
def __init__(self, int buffer_size=1000): self._buffer = [] self._buffer_size = buffer_size
cpdef void add(self, object item): self._buffer.append(item) if len(self._buffer) > self._buffer_size: self._process_buffer()
cdef void _process_buffer(self): """批量处理缓冲区数据""" cdef list batch = self._buffer[:self._buffer_size] self._buffer = self._buffer[self._buffer_size:] # 处理逻辑:发送到下游、保存到DB等
cpdef list get_buffer(self): return self._buffer.copy()功能说明:滑动窗口,动态聚合窗口内数据。
cdef class SlidingWindow: cdef list _window cdef int _size cdef object _reducer
def __init__(self, int size, object reducer): self._window = [] self._size = size self._reducer = reducer
cpdef object add(self, object value): self._window.append(value) if len(self._window) > self._size: self._window.pop(0)
return self._reducer(self._window)输出示例:
>>> window = SlidingWindow(3, sum)>>> window.add(1)1>>> window.add(2)3>>> window.add(3)6>>> window.add(4) # 窗口变为[2,3,4]9功能说明:带状态的数据流处理器。
cdef class StatefulProcessor: cdef dict _state cdef int _counter
def __init__(self): self._state = {} self._counter = 0
cpdef object process(self, object key, object value): if key not in self._state: self._state[key] = [] self._state[key].append(value) self._counter += 1
return self._state[key]
cpdef dict get_state(self): return self._state.copy()
cpdef void reset(self): self._state = {} self._counter = 0| 操作 | Python | Cython加速 |
|---|---|---|
| CSV解析 | 慢 | 5-10x |
| 压缩解压 | 中 | 2-3x |
| 数据库插入 | 慢 | 3-5x |
| 流处理 | 中 | 2-5x |
- CSV解析用csv.reader避免split
- 数据库用executemany批量操作
- 连接池复用连接减少开销
- 流处理用窗口函数减少全量计算
- 实现快速CSV解析器,对比100万行处理时间
- 创建二进制文件读写函数(读写int32数组)
- 实现数据库连接池,测试连接复用
- 创建滑动窗口处理器(支持sum、avg、max)
- 实现带状态的数据流处理(按key聚合)
- 测试JSON压缩解压的性能和压缩率