匡醍量化|大富翁量化

21 Days to AI Trader: Optimizing System Logic and Minute-Level Data Synthesis

中文 📅 2025-06-15 👁 views this month —

When tick data floods in like a tide, how can the system intelligently synthesize it into valuable minute-level data? This article takes you deep into the core of quantitative trading systems—the world of data synthesis and system architecture optimization!

"007, our real-time tick data subscription system is basically complete, but now I'm facing a new challenge." I said to my AI assistant while looking at the mountain of tick data piling up in Redis.

"What's the challenge?" 007 replied immediately.

"We now have massive amounts of tick data, but quantitative strategies require minute-level data. Moreover, I want the system to intelligently handle both intraday and historical data, allowing multiple clients to query seamlessly." I pointed to the dense data on the screen.

This is Day 9 of our quantitative trading system development. In the previous days, we successfully built the infrastructure to fetch data from Tushare and implemented real-time tick data subscription via QMT. However, in actual usage, I identified a critical issue: although tick data is precise, minute-level data is what most quantitative strategies actually need.

More importantly, I needed an intelligent system architecture that could:

  • Real-time synthesize tick data into multi-period minute-level bars
  • Intelligently distinguish between the storage and querying of intraday and historical data
  • Support simultaneous queries from multiple clients without impacting system performance
  • Ensure data integrity and consistency

🎯 Requirements Analysis: Building an Intelligent Data Synthesis System

"Before we start coding, we need to clarify the core requirements of the system." I said to 007.

After deep consideration, I outlined the following key system architecture:

Note:

  • All minute-level data (whether synthesized intraday or subscribed historical minute data) must be during trading hours. Otherwise, it is meaningless.
  • Intraday synthesized minute data is stored in Redis, not ClickHouse.
  • Only historical minute data subscribed from QMT is stored in ClickHouse via Redis. Please distinguish this from intraday synthesized minute data.

References: ClickHouse Redis Engine and ClickHouse Materialized Views.

"This architecture looks complex, especially the data synthesis part." I said with some concern.

"No worries! We can implement it step by step. First, build the basic architecture, then gradually optimize the data synthesis algorithms." 007 answered confidently.

🏗️ System Architecture Design: Three-Tier Separated Intelligent Architecture

"We need to design a truly intelligent three-tier separated architecture." 007 began the architecture design.

Overall Architecture Diagram

┌─────────────────┐    ┌─────────────────┐    ┌─────────────────┐
│   Windows端     │    │    远程Redis     │    │    Mac端        │
│   数据生产者     │───▶│   消息队列+缓存   │───▶│   数据消费者     │
│                 │    │                 │    │                 │
│ • QMT分笔订阅    │    │ • 分笔数据队列   │    │ • 历史数据存储   │
│ • 分钟线合成     │    │ • 当日分钟线缓存 │    │ • ClickHouse管理 │
│ • 交易时间验证   │    │ • 数据路由       │    │ • 数据清理       │
└─────────────────┘    └─────────────────┘    └─────────────────┘
                                │
                                ▼
                       ┌─────────────────┐
                       │   多Client端    │
                       │   数据查询者     │
                       │                 │
                       │ • Web查询界面   │
                       │ • 智能数据路由   │
                       │ • 24小时制时间   │
                       └─────────────────┘

Data Flow Design

"The design of the data flow is the core of the entire system." I emphasized to 007.

Data Flow:

  1. Tick Data Flow: QMT → Windows Client → Redis Queue
  2. Intraday Minute Data Flow: Windows Client Synthesis → Redis Cache → Client Query
  3. Historical Minute Data Flow: Mac Client Processing → ClickHouse Storage → Client Query
  4. Hybrid Query Flow: Client → Redis + ClickHouse → Data Merge → Return Results

🔧 Windows Client Implementation

"The Windows client is the heart of the entire system, responsible for real-time data synthesis." 007 began the Windows client design.

How to Synthesize Tick Data into Minute Bars?

007 designed an ingenious data synthesis engine:

class BarDataSynthesizer:
    """分钟线数据合成器"""

    def __init__(self):
        # 存储各个股票的分笔数据缓存
        self.tick_cache: Dict[str, List[TickData]] = defaultdict(list)
        # 存储各个周期的分钟线缓存
        self.bar_cache: Dict[int, Dict[str, List[BarData]]] = {
            1: defaultdict(list),
            5: defaultdict(list),
            15: defaultdict(list),
            30: defaultdict(list)
        }
        # 交易时间验证器
        self.trading_validator = TradingTimeValidator()

    def add_tick_data(self, tick_data: TickData):
        """添加分笔数据并触发合成"""
        # 验证交易时间
        if not self.trading_validator.validate_tick_data(tick_dict):
            return

        # 缓存分笔数据
        self.tick_cache[tick_data.symbol].append(tick_data)

        # 合成1分钟线
        bar_1min = self._synthesize_1min_bar(tick_data.symbol)
        if bar_1min:
            self.bar_cache[1][tick_data.symbol].append(bar_1min)

            # 基于1分钟线合成其他周期
            for period in [5, 15, 30]:
                bar = self._synthesize_multi_min_bar(tick_data.symbol, period)
                if bar:
                    self.bar_cache[period][tick_data.symbol].append(bar)

1-Minute Bar Synthesis:

def _synthesize_1min_bar(self, symbol: str) -> BarData:
    """合成1分钟线"""
    ticks = self.tick_cache[symbol]
    if not ticks:
        return None

    # 获取当前分钟的开始时间
    current_time = ticks[-1].time
    minute_start = current_time.replace(second=0, microsecond=0)
    minute_end = minute_start + timedelta(minutes=1)

    # 筛选当前分钟的分笔数据
    minute_ticks = [
        tick for tick in ticks
        if minute_start <= tick.time < minute_end
    ]

    if not minute_ticks:
        return None

    # 计算OHLCV
    prices = [tick.price for tick in minute_ticks]
    volumes = [tick.volume for tick in minute_ticks]
    amounts = [tick.amount for tick in minute_ticks]

    return BarData(
        symbol=symbol,
        frame=minute_start,
        open=prices[0],
        high=max(prices),
        low=min(prices),
        close=prices[-1],
        vol=sum(volumes),
        amount=sum(amounts)
    )

Multi-Period Synthesis:

def _synthesize_multi_min_bar(self, symbol: str, period: int) -> BarData:
    """合成多分钟线(5分钟、15分钟、30分钟)"""
    bars_1min = self.bar_cache[1][symbol]
    if not bars_1min:
        return None

    # 获取当前周期的开始时间
    current_time = bars_1min[-1].frame
    period_start = self._get_period_start(current_time, period)
    period_end = period_start + timedelta(minutes=period)

    # 筛选当前周期的1分钟线数据
    period_bars = [
        bar for bar in bars_1min
        if period_start <= bar.frame < period_end
    ]

    if len(period_bars) == 0:
        return None

    # 合成多分钟线
    return BarData(
        symbol=symbol,
        frame=period_start,
        open=period_bars[0].open,
        high=max(bar.high for bar in period_bars),
        low=min(bar.low for bar in period_bars),
        close=period_bars[-1].close,
        vol=sum(bar.vol for bar in period_bars),
        amount=sum(bar.amount for bar in period_bars)
    )

Our synthesis method is to first synthesize 1-minute bars, and then synthesize other periods based on the 1-minute bars, ensuring data consistency.

Windows Client Monitoring Page

However, I felt that the layout of this monitoring page was too simplistic, so I asked 007 to design a new frontend page to better present the monitoring status.

🍎 Mac Client Implementation

Data Consumption and Storage

The core responsibility of the Mac client is to transfer historical data from Redis to ClickHouse:

class MacDataService:
    """Mac端数据服务"""

    def __init__(self):
        self.redis_manager = RedisManager()
        self.clickhouse_manager = ClickHouseManager()

    def consume_historical_data(self):
        """消费历史数据"""
        while self.is_running:
            try:
                # 从Redis获取历史分钟线数据
                for period in [1, 5, 15, 30]:
                    queue_name = f"historical_bar_data_{period}min"
                    data = self.redis_manager.client.brpop(queue_name, timeout=1)

                    if data:
                        bar_data = BarData(**json.loads(data[1]))

                        # 检查是否已存在
                        if not self.clickhouse_manager.data_exists(bar_data, period):
                            # 插入ClickHouse
                            self.clickhouse_manager.insert_bar_data(bar_data, period)
                        else:
                            # 数据已存在,直接删除Redis中的数据
                            self.logger.info(f"数据已存在,跳过: {bar_data.symbol} {bar_data.frame}")

            except Exception as e:
                self.logger.error(f"数据消费错误: {e}")
                time.sleep(5)

Special Handling at 2 AM

"The system needs special data cleanup at 2 AM." I explained the requirements to 007.

def handle_cleanup_time(self):
    """处理凌晨2点的数据清理"""
    try:
        # 1. 处理前一天的历史数据
        self.process_previous_day_data()

        # 2. 清理Redis的订阅消息队列
        self.cleanup_redis_queues()

        # 3. 数据完整性检查
        self.verify_data_integrity()

        self.logger.info("凌晨2点数据清理完成")

    except Exception as e:
        self.logger.error(f"数据清理错误: {e}")

Mac Client Monitoring Page

To ensure that data transmission between Redis and ClickHouse is normal, and to monitor the speed and progress of data transmission in real-time, I asked 007 to design a Mac client monitoring page:

💻 Client Implementation

"The Client is the interface users interact with directly; it must be perfect." I emphasized to 007.

However, during the development of the Client, we encountered a series of challenges...

First Attempt: Complex Debugging System

Initially, 007 designed a feature-rich debugging system for the Client:

# 复杂的调试逻辑
def query_bar_data(self, symbol: str, start_time: datetime, end_time: datetime, period: int):
    debug_info = []
    debug_info.append(f"查询参数: {symbol}, {period}分钟, {start_time} 到 {end_time}")
    debug_info.append(f"今天: {today}, 查询日期范围: {start_date} 到 {end_date}")

    # 大量的调试信息...
    if len(all_redis_data) > 0 and len(redis_data) == 0:
        debug_info.append(f"⚠️ 时间过滤导致数据为空")
        debug_info.append(f"数据样本时间: {sample_bar.frame}")
        # 更多调试信息...

"This debugging system is too complex!" I said, rubbing my temples looking at the screen full of debugging code.

Second Attempt: The Time Format Nightmare

Next, we encountered time format issues. Users complained that the Web interface displayed AM/PM format:

"I want 24-hour time queries. Why is there AM and PM on the frontend? I don't want AM/PM. I'm so annoyed!"

007 immediately fixed it, replacing the datetime-local input with separated time inputs:

<!-- 24小时制时间输入 -->
<div class="col-md-3">
    <label for="start_time" class="form-label">开始时间 (24小时制)</label>
    <div class="row g-1">
        <div class="col-6">
            <input type="date" class="form-control" id="start_date" required>
        </div>
        <div class="col-3">
            <input type="number" class="form-control" id="start_hour" min="0" max="23" placeholder="时" required>
        </div>
        <div class="col-3">
            <input type="number" class="form-control" id="start_minute" min="0" max="59" placeholder="分" required>
        </div>
    </div>
</div>

Third Attempt: JSON Serialization Traps

Then, we encountered JSON serialization errors:

查询失败: 请求处理错误: Object of type datetime is not JSON serializable

"This error is common; datetime objects cannot be serialized directly." 007 explained.

We tried various solutions and finally adopted manual serialization:

# 手动序列化,确保datetime正确转换
data_list = []
for bar in result.data:
    data_list.append({
        "symbol": bar.symbol,
        "frame": bar.frame.isoformat(),  # 手动转换datetime为字符串
        "open": float(bar.open),
        "high": float(bar.high),
        "low": float(bar.low),
        "close": float(bar.close),
        "vol": float(bar.vol),
        "amount": float(bar.amount)
    })

Final Refactoring: Concise and Powerful

"007! I think you've made the Client query code a mess. You should tear it all down and rewrite the queries according to the logic of Windows and Mac." I finally couldn't hold back.

007 immediately performed a thorough refactoring, rewriting the code from scratch according to the original system architecture design:

def query_bar_data(self, symbol: str, start_time: datetime, end_time: datetime, period: int) -> QueryResponse:
    """
    查询分钟线数据

    按照系统架构:
    1. 如果查询的分钟线数据是当日的,则直接从Redis中读取合成的分钟线数据
    2. 如果查询的分钟线数据是历史的,则直接从ClickHouse中读取
    3. 如果查询的分钟线数据是既有当日的又有历史的,则合并数据返回给Client
    """
    try:
        today = date.today()
        start_date = start_time.date()
        end_date = end_time.date()

        redis_data = []
        clickhouse_data = []

        # 1. 查询当日数据(从Redis读取)
        if end_date >= today:
            redis_data = self.redis_manager.get_current_bar_data(period, symbol)
            # 过滤时间范围
            redis_data = [bar for bar in redis_data if start_time <= bar.frame <= end_time]

        # 2. 查询历史数据(从ClickHouse读取)
        if start_date < today:
            # 避免与当日数据重复,历史数据查询到今天之前
            hist_end_time = min(end_time, datetime.combine(today, datetime.min.time()))
            if start_time < hist_end_time:
                clickhouse_data = self.clickhouse_manager.query_bar_data(
                    symbol, start_time, hist_end_time, period
                )

        # 3. 合并数据
        merged_data = self.data_merger.merge_bar_data(redis_data, clickhouse_data)

        return QueryResponse(
            success=True,
            message=f"当日数据: {len(redis_data)} 条,历史数据: {len(clickhouse_data)} 条,合并后: {len(merged_data)} 条",
            data=merged_data,
            total_count=len(merged_data)
        )

    except Exception as e:
        return QueryResponse(
            success=False,
            message=f"查询失败: {str(e)}",
            data=[],
            total_count=0
        )

Client Query Page

For easy querying and visualized query results, I asked 007 to design a query page for the Client:

⚠️ Important: Data Merging and Trading Time Validation

Data Merging: Ensuring Query Accuracy

If the date a user needs to query contains both historical minute data and intraday minute data, we must merge the data from Redis and ClickHouse.

class DataMerger:
    """数据合并器 - 合并Redis当日数据和ClickHouse历史数据"""

    @staticmethod
    def merge_bar_data(redis_data: List[BarData], clickhouse_data: List[BarData]) -> List[BarData]:
        """合并分钟线数据"""
        # 合并数据并按时间排序
        all_data = redis_data + clickhouse_data

        # 去重(以frame和symbol为键)
        unique_data = {}
        for bar in all_data:
            key = (bar.symbol, bar.frame)
            unique_data[key] = bar

        # 按时间排序
        merged_data = list(unique_data.values())
        merged_data.sort(key=lambda x: x.frame)

        return merged_data

Trading Time Validation: Ensuring Data Quality

"All minute-level data must be within trading hours; otherwise, it is meaningless." I emphasized the importance of data quality to 007.

007 designed a specialized trading time validator:

class TradingTimeValidator:
    """交易时间验证器"""

    def __init__(self):
        self.trading_hours = {
            'morning_start': '09:30:00',
            'morning_end': '11:30:00',
            'afternoon_start': '13:00:00',
            'afternoon_end': '15:00:00'
        }

    def is_trading_time(self, dt: datetime) -> bool:
        """检查是否为交易时间"""
        time_obj = dt.time()

        morning_start = dt_time.fromisoformat(self.trading_hours['morning_start'])
        morning_end = dt_time.fromisoformat(self.trading_hours['morning_end'])
        afternoon_start = dt_time.fromisoformat(self.trading_hours['afternoon_start'])
        afternoon_end = dt_time.fromisoformat(self.trading_hours['afternoon_end'])

        return (
            (morning_start <= time_obj <= morning_end) or
            (afternoon_start <= time_obj <= afternoon_end)
        )

    def validate_tick_data(self, tick_dict: dict) -> bool:
        """验证分笔数据"""
        return self.is_trading_time(tick_dict['time'])

    def validate_bar_data(self, bar_dict: dict) -> bool:
        """验证分钟线数据"""
        return self.is_trading_time(bar_dict['frame'])

🎯 System Testing

"The system is built. Let's test its performance." I said with anticipation.

Fetching Historical Minute Data

Windows Client Performance:

Mac Client Performance:

You can notice that the elements in the Redis queue on the Mac client are decreasing, while the data stored in ClickHouse is increasing.

Let's look at ClickHouse again:

It can be seen that historical data has been stored in ClickHouse (and all of it meets the requirement of being within trading hours).

Fetching Intraday Minute Data

By starting the Windows and Mac clients during trading hours, you can see that Redis stores both tick data and synthesized intraday minute data. (This data will be cleared uniformly at 2 AM).

Redis:

We can query the intraday minute data via the Client (using 15-minute bars as an example):

📝 Summary

Through this system logic optimization and minute-level data synthesis development, 007 and I successfully built an intelligent, efficient, and stable quantitative trading data processing system. From the initial complex design to the final concise implementation, we not only solved technical challenges but also established an extensible and maintainable system architecture.

Nine days have passed in the 21-day challenge. Our quantitative trading system is becoming increasingly powerful. From simple data acquisition initially to intelligent data synthesis now, every step has been full of challenges and rewards.

007's performance has once again impressed me. From architecture design to code implementation, from problem diagnosis to system optimization, it has demonstrated professional technical capabilities. More importantly, it can quickly adjust direction under my guidance to ultimately achieve a concise and powerful system.

Especially during the Client refactoring process, 007 showed strong learning abilities. When I pointed out that the code was too complex, it immediately understood the problem and thoroughly refactored it according to the original system architecture design. This ability to respond quickly and self-correct is exactly the quality an excellent AI assistant should possess.

However, is our minute synthesis really reasonable? In the next chapter, we will use Tushare and AkShare to verify the accuracy of our synthesized real-time intraday minutes. For the multi-Client issue, we will also test the system using multiple Client machines. Stay tuned!