匡醍量化|大富翁量化

Optimizing Quant Data Pipelines: DB Design for Factors & Intraday Bars

中文 📅 2025-05-13 👁 views this month —

Before the May Day holiday, with the help of my partner 007, I successfully implemented fetching OHLC data from Tushare and storing it in a ClickHouse database via a Redis message queue. To further enhance the data support for our quantitative trading system, this article focuses on database optimization design, specifically covering the acquisition of daily adjusted factors, the retrieval of minute-level data, and the computation of minute-level data (e.g., 15-minute intervals). After the holiday, we plan to leverage these optimizations to provide richer, more granular data for quantitative trading strategies, thereby improving system performance and decision-making capabilities.


1. Acquiring Daily Adjusted Factors

Me: 007! Please help me add code to fetch daily factors in data_producer.py and data_consumer.py.

007: 🫡 Request received, processing...

To incorporate factor data into data_consumer.py and data_producer.py, we need to modify the original code. The core idea is to fetch factor data in the producer code, encapsulate it into a data structure, and then process this factor data in the consumer code. Since we will be collecting minute-level data later, I created a new table named daily_data to store daily factor data.

Next, I need to modify data_producer.py and data_consumer.py to handle this data table.

Here is the complete code provided by 007:

"""data_producer.py"""
import redis
import tushare as ts
import json
from datetime import datetime

# Tushare and Redis configuration
TUSHARE_TOKEN = "YOUR TOKEN"
REDIS_HOST = "localhost"
REDIS_PORT = 6379
REDIS_QUEUE_NAME = "tushare_data_queue"

# Initialize connections
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 data fetch failed: {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"Adjusted factor fetch failed: {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"Data pushed: {ts_code} - {start_date} to {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

# Configuration parameters
REDIS_HOST = "localhost"
REDIS_PORT = 6379
REDIS_QUEUE_NAME = "tushare_data_queue"

CLICKHOUSE_HOST = "localhost"
CLICKHOUSE_PORT = 9000
CLICKHOUSE_DB = "default"

# Initialize Redis and ClickHouse clients
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):
    """Insert data into 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():
    """Main data consumption function"""
    print("Starting data consumer, waiting for queue data...")
    while True:
        try:
            # Blocking fetch from queue
            result = redis_client.brpop(REDIS_QUEUE_NAME, timeout=1)
            if result is None:
                # If no data is fetched, the queue is empty, exit loop
                print("Redis queue is empty, stopping data consumption.")
                break
            _, json_data = result
            data_package = json.loads(json_data)
            insert_to_clickhouse(data_package)
            print(f"Successfully inserted data: {len(data_package['ohlc_data'])} records")
        except Exception as e:
            print(f"Data processing exception: {str(e)}")
            continue

if __name__ == "__main__":
    consume_data()

2. Acquiring Minute-Level Data

Following the steps above, I need to acquire minute-level data and add it to the data table.

007 provided the following suggestions:

  1. First, create the minute-level data table;
  2. Modify the producer code to add minute-level data acquisition functionality;
  3. Create the corresponding consumer code.

2.1. Creating the Minute-Level Data Table

![](https://cdn.jsdelivr.net/gh/zillionare/images@main/images/2025/05/3_04.png)
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);

Bridging OS Gaps: QMT Minute Data via Redis to ClickHouse

Implementing a cross-platform data pipeline using QMT, Redis, and ClickHouse to ingest high-frequency minute-level data from Windows-based brokerage APIs into a Mac-based quantitative research environment.

2.2. Modifying the Producer Code to Ingest Minute-Level Data

While the original author (007) relied on Tushare for minute-level data, this implementation switches to the QMT API for data acquisition.

I added a new producer file, minute_producer.py, modifying the original data_producer.py with the following key changes:

However, I encountered a significant hurdle: QMT currently supports only Windows, while my development environment is macOS.

To resolve this, I implemented a solution using Redis as a middleware to transfer data from the Windows machine to the macOS program, ultimately storing it in ClickHouse.

Following 007’s recommendation, I developed the following code.

2.2.3. Windows Data Producer Code

The original code provided by 007 is as follows:

import redis
import json
from datetime import datetime
from xtquant.xtdata import (
    init,
    download_history_data,
    get_local_data,
    get_trading_dates,
    close
)

# Redis Configuration - Use Mac's IP Address
REDIS_HOST = "Replace with Mac's IP Address"
REDIS_PORT = 6379
REDIS_QUEUE_NAME = "qmt_minute_queue"
REDIS_PASSWORD = None  # Set if password is required

def setup_redis_client():
    """Initialize Redis Client"""
    return redis.StrictRedis(
        host=REDIS_HOST, 
        port=REDIS_PORT, 
        password=REDIS_PASSWORD,
        decode_responses=True
    )

def fetch_minute_data(stock_code, date_str):
    """Fetch minute-level data for a specific date"""
    try:
        # Fetch minute-level data
        df = get_local_data(stock_code, 'min1', date_str, date_str)
        if df is None or len(df) == 0:
            return []
        
        # Convert data format
        records = []
        for time, row in df.iterrows():
            records.append({
                "ts_code": stock_code,
                "trade_time": time.strftime("%Y-%m-%d %H:%M:%S"),
                "open": float(row['open']),
                "high": float(row['high']),
                "low": float(row['low']),
                "close": float(row['close']),
                "vol": float(row['volume']),
                "amount": float(row['amount'])
            })
        return records
    except Exception as e:
        print(f"Failed to fetch minute data for {stock_code} {date_str}: {str(e)}")
        return []

def main():
    """Main Function"""
    # Initialize QMT Interface
    init()
        
    # Initialize Redis Client
    redis_client = setup_redis_client()
        
    try:
        # Configuration Parameters
        stock_list = ["000001.SZ", "600519.SH"]
        start_date = "20230101"
        end_date = "20230131"
            
        # Download Historical Data
        print(f"Starting historical data download: {start_date} to {end_date}")
        download_history_data(stock_list, 'min1', start_date, end_date)
        print("Historical data download complete")
            
        # Get Trading Date List
        trading_dates = get_trading_dates(start_date, end_date)
            
        # Fetch and Push Minute Data to Redis by Date and Stock Code
        for trade_date in trading_dates:
            date_str = trade_date.strftime("%Y%m%d")
            print(f"Processing Date: {date_str}")
                
            for stock_code in stock_list:
                minute_data = fetch_minute_data(stock_code, date_str)
                
                if minute_data:
                    # Package Data
                    data_package = {
                        "timestamp": datetime.now().isoformat(),
                        "ts_code": stock_code,
                        "trade_date": date_str,
                        "minute_data": minute_data
                    }
                        
                    # Push to Redis
                    redis_client.lpush(REDIS_QUEUE_NAME, json.dump(data_package))
                    print(f"Pushed minute data: {stock_code} - {date_str} ({len(minute_data)} records)")
        
    except Exception as e:
        print(f"Program execution error: {str(e)}")
        
    finally:
        # Close QMT Interface
        close()
        print("Program execution complete")

if __name__ == "__main__":
    main()

This code is non-functional in its current state due to potential changes in the QMT library version, where some modules may have been removed or modified. Additionally, connecting from a Windows machine to a Redis server running on macOS involves network connectivity, firewall settings, Redis configuration, and permission issues.

To address the Redis connection issues, I utilized 007’s "helpful code" to test the Redis connection:

# Redis Configuration - Use Mac's IP Address
REDIS_HOST = "Replace with Mac's IP Address"
REDIS_PORT = 6379
REDIS_QUEUE_NAME = "qmt_minute_queue"
REDIS_PASSWORD = None  # Set if password is required


# Test Redis Connection
import redis
import time

try:
    # Create Redis Client
    redis_client = redis.Redis(
        host=REDIS_HOST,
        port=REDIS_PORT,
        password=REDIS_PASSWORD,
        socket_timeout=5,
        decode_responses=True
    )
    
    # Test Connection - PING Command
    response = redis_client.ping()
    print(f"Redis Connection Test (PING): {'Success' if response else 'Failed'}")
    
    # Test Basic Operations - Write and Read
    test_key = "test_connection_key"
    test_value = f"test_value_{time.time()}"
    
    # Write Test
    redis_client.set(test_key, test_value)
    print(f"Redis Write Test: Successfully wrote key '{test_key}'")
    
    # Read Test
    read_value = redis_client.get(test_key)
    print(f"Redis Read Test: {'Success' if read_value == test_value else 'Failed'}")
    print(f"Written Value: {test_value}")
    print(f"Read Value: {read_value}")
    
    # Test Queue Operations
    redis_client.lpush(REDIS_QUEUE_NAME, "Test Message")
    queue_length = redis_client.llen(REDIS_QUEUE_NAME)
    print(f"Redis Queue Test: Successfully wrote to queue '{REDIS_QUEUE_NAME}', current queue length: {queue_length}")
    
    # Clean Up Test Data
    redis_client.delete(test_key)
    
    print("Redis connection and basic operation tests completed. Connection is normal.")
    
except redis.exceptions.ConnectionError as e:
    print(f"Redis Connection Error: {str(e)}")
    print("Please check the following:")
    print("1. Is the Redis server running?")
    print("2. Is the IP address correct?")
    print("3. Is the port correct?")
    print("4. Does the firewall allow the connection?")
    print("5. Is Redis configured to allow remote connections?")
    
except Exception as e:
    print(f"Other errors occurred during Redis testing: {str(e)}")

With the test code in place, we must address Redis’s default configuration. By default, Redis only accepts local connections, binding to 127.0.0.1. To allow remote connections, the Redis configuration file must be modified to bind to 0.0.0.0 or the Mac’s local network IP. This requires editing the redis.conf file to adjust the bind parameter.

  1. Modify Redis Configuration File By default, Redis listens only on the local IP (127.0.0.1). Adjust to allow remote connections:

    • Open the configuration file:
      sudo nano /usr/local/etc/redis.conf
    • Modify the following parameters:
      • Bind IP: Change bind 127.0.0.1 to bind 0.0.0.0 (allow all IPs) or replace it with the Mac’s local IP (e.g., 192.168.1.100).

      • Disable Protected Mode: Change protected-mode yes to protected-mode no.

      • Set Password (Optional but Recommended): Uncomment requirepass and set a password:

        requirepass your_password

    • Save and Exit: Press Ctrl+O to save, then Ctrl+X to exit.
  2. Restart Redis Service (Choose One)

    brew services restart redis  # For Homebrew installation
    redis-server /usr/local/etc/redis.conf  # Manual restart
  3. Open Mac Firewall Port

    • GUI Operation:
      • Navigate to System Preferences → Security & Privacy → Firewall.
      • Click the lock icon to unlock, then select Firewall Options.
      • Click + to add redis-server to the allowed list.
    • Command Line Operation (Requires Admin Privileges):
    sudo /usr/libexec/ApplicationFirewall/socketfilterfw --add /usr/local/bin/redis-server
    sudo /usr/libexec/ApplicationFirewall/socketfilterfw --unblockapp /usr/local/bin/redis-server

After completing the above steps, we configure the Windows side (installing the Redis client):

  1. Install Redis Client

    • Download Windows Redis: Download the stable version from the Redis Official Website and extract it to any directory (e.g., C:\redis).
    • Add to System Path: Add C:\redis\bin to the environment variable PATH to use redis-cli directly from the command line.
  2. Connect to Redis Server

    • Command Format:
    redis-cli -h <Mac's Local IP> -p 6379 -a <Password>
    • Example (Assuming Mac IP is 192.168.1.100 and password is your_redis_password):
      redis-cli -h 192.168.1.100 -p 6379 -a your_redis_password
    • Verify Connection:
      192.168.1.100:6379> PING
      PONG  # Connection Successful

The above steps do not guarantee that the Redis server on Windows will accept connections from the Mac. (We leave this as an open question; feel free to share your solutions in the comments.)

2.2.4. macOS Data Consumer Code

import redis
import json
from clickhouse_driver import Client
from datetime import datetime
import time

# Redis Configuration
REDIS_HOST = "localhost"  # Local Redis or Windows IP
REDIS_PORT = 6379
REDIS_QUEUE_NAME = "qmt_minute_queue"
REDIS_PASSWORD = None  # Set if password is required

# ClickHouse Configuration
CLICKHOUSE_HOST = "localhost"
CLICKHOUSE_PORT = 9000
CLICKHOUSE_DB = "default"
CLICKHOUSE_USER = "default"
CLICKHOUSE_PASSWORD = ""

def setup_redis_client():
    """Initialize Redis Client"""
    return redis.StrictRedis(
        host=REDIS_HOST, 
        port=REDIS_PORT, 
        password=REDIS_PASSWORD,
        decode_responses=True
    )

def setup_clickhouse_client():
    """Initialize ClickHouse Client"""
    return Client(
        host=CLICKHOUSE_HOST,
        port=CLICKHOUSE_PORT,
        database=CLICKHOUSE_DB,
        user=CLICKHOUSE_USER,
        password=CLICKHOUSE_PASSWORD
    )

def insert_to_clickhouse(client, data):
    """Insert minute-level data into ClickHouse"""
    if not data["minute_data"]:
        print("No data to insert")
        return 0
    
    query = """
    INSERT INTO minute_data 
    (ts_code, trade_time, open, high, low, close, vol, amount)
    VALUES
    """
    
    values = []
    for record in data["minute_data"]:
        values.append((
            record["ts_code"],
            datetime.strptime(record["trade_time"], "%Y-%m-%d %H:%M:%S"),
            record["open"],
            record["high"],
            record["low"],
            record["close"],
            record["vol"],
            record["amount"]
        ))
    
    if values:
        client.execute(query, values)
        return len(values)
    return 0

def main():
    """Main Function"""
    # Initialize Clients
    redis_client = setup_redis_client()
    clickhouse_client = setup_clickhouse_client()
    
    print("Starting minute-level data consumer, waiting for queue data...")
    
    try:
        while True:
            # Attempt to fetch data from Redis
            result = redis_client.brpop(REDIS_QUEUE_NAME, timeout=1)
            
            if result is None:
                print("Redis queue is empty, waiting for new data...")
                time.sleep(5)  # Wait 5 seconds before retrying
                continue
            
            # Parse Data
            _, json_data = result
            data_package = json.loads(json_data)
            
            # Insert into ClickHouse
            inserted_count = insert_to_clickhouse(clickhouse_client, data_package)
            print(f"Successfully inserted minute data: {data_package['ts_code']} - {data_package['trade_date']} ({inserted_count} records)")
    
    except KeyboardInterrupt:
        print("Program manually interrupted")
    
    except Exception as e:
        print(f"Program execution error: {str(e)}")
    
    finally:
        print("Program execution complete")

if __name__ == "__main__":
    main()

Since the Redis connection between Windows and macOS remains unresolved, I plan to first store the minute-level data obtained from QMT into 000001.SH_data.csv and 300750.SZ_data.csv, then load the data into ClickHouse, and finally process it using Python.

from xtquant import xtdata
import os
import pandas as pd

code_list = ['000001.SH', '300750.SZ']
period = '1h'
start_time = '20250101093000'
end_time = '20250201093000'

def on_data(datas):
    if datas:
        print(datas)
    else:
        print("Data download failed or is empty")

xtdata.download_history_data2(code_list, period, start_time, end_time, on_data)

# Create Directory (if it doesn't exist)
save_dir = 'C:\\wbq'
os.makedirs(save_dir, exist_ok=True)

for code in code_list:
    data = xtdata.get_market_data_ex([], [code], period, start_time, end_time)
    if code in data and not data[code].empty:
        # Create separate file for each stock
        file_path = os.path.join(save_dir, f'{code}_data.csv')
        # Ensure data is in DataFrame format
        df = data[code]
        # Save data with error handling
        try:
            df.to_csv(file_path)
            print(f'{code} data saved locally: {file_path}')
        except Exception as e:
            print(f'Error saving {code} data: {str(e)}')
    else:
        print(f'{code} did not retrieve any data')

print("Execution complete")