21天驯化AI打工仔 - 数据库的优化
五一小长假之前,我在搭档 007 的帮助下已经成功实现了从 Tushare 获取 OHLC 数据,并通过 Redis 消息队列将数据存储到 ClickHouse 数据库。为了进一步完善量化交易系统的数据支持,今天我们将聚焦于数据库的优化设计,主要涉及获取日线复权因子、获取分钟线数据以及计算分钟线数据(例如 15 分钟)三个方面。五一结束之后,我们计划通过这些优化,能够为量化交易策略提供更丰富、更精细的数据,从而提升系统的性能和决策能力。
1. 获取日线的复权因子
我:007!帮我在
data_producer.py和data_consumer.py中添加获取日线因子的代码。007:🫡收到请求,正在处理...
为了在 data_consumer.py 和 data_producer.py 中补充因子的数据,我们需要对原代码进行一些修改。主要思路是在生产者代码中获取因子数据并封装到数据结构中,然后在消费者代码中处理这些因子数据。因为后续我们要收集分钟级的数据,所以我新建了一张表,命名为 daily_data,用于存储日线因子数据。

然后,我需要针对这个数据表,对 data_producer.py 和 data_consumer.py 进行修改。

007 回答我的完整代码如下:
"""data_producer.py"""
import redis
import tushare as ts
import json
from datetime import datetime
# Tushare和Redis配置
TUSHARE_TOKEN = "YOUR TOKEN"
REDIS_HOST = "localhost"
REDIS_PORT = 6379
REDIS_QUEUE_NAME = "tushare_data_queue"
# 初始化连接
pro = ts.pro_api(TUSHARE_TOKEN)
redis_client = redis.StrictRedis(host=REDIS_HOST, port=REDIS_PORT, decode_responses=True)
def fetch_ohlc_daily_data(ts_code: str, start_date: str, end_date: str) -> list:
try:
df = pro.daily(ts_code=ts_code, start_date=start_date, end_date=end_date)
return df[['ts_code', 'trade_date', 'open', 'high', 'low', 'close', 'vol']].to_dict('records')
except Exception as e:
print(f"OHLC数据获取失败: {str(e)}")
return []
def fetch_adj_factor(ts_code: str, start_date: str, end_date: str) -> list:
try:
df = pro.adj_factor(ts_code=ts_code, start_date=start_date, end_date=end_date)
return df[['ts_code', 'trade_date', 'adj_factor']].to_dict('records')
except Exception as e:
print(f"复权因子获取失败: {str(e)}")
return []
def produce_data(ts_code_list: list, date_range: tuple):
start_date, end_date = date_range
for ts_code in ts_code_list:
data_package = {
"timestamp": datetime.now().isoformat(),
"ts_code": ts_code,
"ohlc_data": fetch_ohlc_daily_data(ts_code, start_date, end_date),
"adj_factor": fetch_adj_factor(ts_code, start_date, end_date)
}
redis_client.lpush(REDIS_QUEUE_NAME, json.dumps(data_package))
print(f"已推送数据: {ts_code} - {start_date}至{end_date}")
if __name__ == "__main__":
STOCK_CODES = ["000001.SZ", "600519.SH"]
DATE_RANGE = ("20230101", "20231231")
produce_data(STOCK_CODES, DATE_RANGE)
"""data_consumer.py"""
import redis
import json
from clickhouse_driver import Client
from datetime import datetime
# 配置参数
REDIS_HOST = "localhost"
REDIS_PORT = 6379
REDIS_QUEUE_NAME = "tushare_data_queue"
CLICKHOUSE_HOST = "localhost"
CLICKHOUSE_PORT = 9000
CLICKHOUSE_DB = "default"
# 初始化 Redis 和 ClickHouse 客户端
redis_client = redis.StrictRedis(host=REDIS_HOST, port=REDIS_PORT, decode_responses=True)
clickhouse_client = Client(host=CLICKHOUSE_HOST, port=CLICKHOUSE_PORT, database=CLICKHOUSE_DB)
def insert_to_clickhouse(data):
"""将数据插入到 ClickHouse"""
query = """
INSERT INTO daily_data (ts_code, trade_date, open, high, low, close, vol, adj_factor)
VALUES
"""
values = []
for record in data["ohlc_data"]:
adj_factor_record = next((adj for adj in data["adj_factor"] if adj["trade_date"] == record["trade_date"]), None)
adj_factor = adj_factor_record["adj_factor"] if adj_factor_record else None
values.append((
record["ts_code"],
datetime.strptime(record["trade_date"], "%Y%m%d").date(),
record["open"],
record["high"],
record["low"],
record["close"],
record["vol"],
adj_factor
))
clickhouse_client.execute(query, values)
def consume_data():
"""数据消费主函数"""
print("启动数据消费者,等待队列数据...")
while True:
try:
# 阻塞式获取队列数据
result = redis_client.brpop(REDIS_QUEUE_NAME, timeout=1)
if result is None:
# 如果没有获取到数据,说明队列为空,退出循环
print("Redis 队列为空,停止消费数据。")
break
_, json_data = result
data_package = json.loads(json_data)
insert_to_clickhouse(data_package)
print(f"成功插入数据: {len(data_package['ohlc_data'])} 条")
except Exception as e:
print(f"数据处理异常: {str(e)}")
continue
if __name__ == "__main__":
consume_data()
2. 获取分钟线数据
同上述步骤,我需要获取分钟线数据,并添加到数据表中。

007 给我提供了以下建议:
- 首先创建分钟线数据表;
- 修改生产者代码,添加分钟线数据获取功能;
- 创建对应的消费者代码。
2.1. 创建分钟线数据表
CREATE TABLE IF NOT EXISTS minute_data (
ts_code String,
trade_time DateTime,
open Float32,
high Float32,
low Float32,
close Float32,
vol Float32,
amount Float32
) ENGINE = MergeTree()
ORDER BY (ts_code, trade_time);
2.2. 修改生产者代码,添加分钟线数据获取功能
007 还是采用 tushare 来获取分钟线数据,并添加到数据表中。但是,这里获取分钟线数据我并不打算用 tushare 作为数据源,而是打算采用 qmt 提供的 API 接口来获取分钟级的数据。




