Skip to content

数据处理加速

本章讲解Cython在大数据处理场景的应用。文件解析、数据库操作、流处理是常见瓶颈,Cython可提供3-10x加速。

学习路径:文件解析 → 数据转换 → 数据库 → 流处理

核心场景:

  • CSV/二进制解析:批量处理减少Python开销
  • 数据库批量操作:executemany减少调用次数
  • 流处理:滑动窗口减少全量计算

功能说明:快速CSV解析,避免逐行Python开销。

import csv
import 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 >

功能说明: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) == data
True

功能说明:使用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)

功能说明:流处理器,批量处理缓冲数据。

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

操作PythonCython加速
CSV解析慢5-10x
压缩解压中2-3x
数据库插入慢3-5x
流处理中2-5x
  1. CSV解析用csv.reader避免split
  2. 数据库用executemany批量操作
  3. 连接池复用连接减少开销
  4. 流处理用窗口函数减少全量计算

  1. 实现快速CSV解析器,对比100万行处理时间
  2. 创建二进制文件读写函数(读写int32数组)
  3. 实现数据库连接池,测试连接复用
  4. 创建滑动窗口处理器(支持sum、avg、max)
  5. 实现带状态的数据流处理(按key聚合)
  6. 测试JSON压缩解压的性能和压缩率