匡醍量化|大富翁量化

21-Day AI Quant Build: Data Pipeline with Redis & Tushare

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

IDEA: Can my collaboration with my "AI black slave" develop a quantitative trading system in just 21 days?

What a fascinating challenge. Let’s start today! "Data is the foundation of everything." I’ve decided to have my AI assistant begin by handling data acquisition. (I felt "AI black slave" wasn't very polite, so I named it "Code Machine 007." I hope "007 the Cow" can "code with full power.")

Alright! Let’s officially start the development work!


Redis Startup!

I aim to develop a quantitative trading framework and use Redis to build a message queue, fetching OHLC data and adjustment factors from Tushare.

First, let’s consider the code structure. It needs to be divided into two main parts: a data-fetching producer and a consumer that processes the data. The producer fetches data from Tushare and pushes it into the Redis queue; the consumer retrieves data from the queue for processing. We need to handle Tushare API calls, including fetching OHLC and adjustment factors. This likely requires using Tushare’s pro_api, so users must provide a token.

Next, for the Redis part, we’ll use the redis-py library. The producer uses lpush to push data into the queue, and the consumer uses brpop to fetch data in a blocking manner. We must handle exceptions, such as Tushare API call failures or Redis connection issues, by adding try-except blocks.

Then, consider the independence and executability of the code. We need to import necessary libraries like redis, tushare, and json. In the producer function, data is converted to JSON format before being pushed into the queue. In the consumer function, data is retrieved, parsed, and printed; this can be extended to actual processing logic, such as storing data in a database or performing further calculations.

We also need to address configuration. Redis connection parameters and the Tushare token should be configurable items for easy user modification. The code includes comments indicating where users need to replace values, such as tushare_token and Redis host/port.

Finally, ensure the doubaocanvas format is correct. The identifier is quant-data-redis-queue, the type is general_code, and the title is descriptive. The two code sections serve as scripts for the producer and consumer, with paths producer.py and consumer.py, and accurate descriptions. Check for undefined variables to ensure the code is complete and independently runnable; users only need to replace the configuration to use it.

My "007" indicated this was simple and instantly provided two code files: data_producer.py and data_consumer.py.

"""data_consumer.py"""
import redis
import json
from typing import Dict

# Configuration parameters (must match the producer)
REDIS_HOST = "localhost"
REDIS_PORT = 6379
REDIS_QUEUE_NAME = "tushare_data_queue"

redis_client = redis.StrictRedis(host=REDIS_HOST, port=REDIS_PORT, decode_responses=True)

def process_data(data_package: Dict):
    """Data processing logic (can be extended for actual business processing)"""
    print(f"Data received at: {data_package['timestamp']}")
    print(f"Processing stock code: {data_package['ts_code']}")
    print(f"Number of OHLC data entries: {len(data_package['ohlc_data'])}")
    print(f"Number of adjustment factor entries: {len(data_package['adj_factor'])}")
    # Add actual business logic here, such as data storage/calculation/distribution

def consume_data():
    """Main data consumption function (blocking fetch)"""
    print("Starting data consumer, waiting for queue data...")
    while True:
        try:
            # Blocking fetch from queue (timeout 30 seconds)
            _, json_data = redis_client.brpop(REDIS_QUEUE_NAME, timeout=30)
            data_package = json.loads(json_data)
            process_data(data_package)
            
        except Exception as e:
            print(f"Data processing exception: {str(e)}")
            continue

if __name__ == "__main__":
    consume_data()
"""data_producer.py"""
import redis
import tushare as ts
import json
from datetime import datetime

# Configuration parameters
TUSHARE_TOKEN = "87a5a442931d15e749b4453b4f802c3b87794c3a5968b7b40a153edd"  # Replace with your Tushare 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_data(ts_code: str, start_date: str, end_date: str) -> list:
    """Fetch OHLC data"""
    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:
    """Fetch adjustment factors"""
    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"Adjustment factor fetch failed: {str(e)}")
        return []

def produce_data(ts_code_list: list, date_range: tuple):
    """Main data production function"""
    start_date, end_date = date_range
    for ts_code in ts_code_list:
        # Encapsulate data structure
        data_package = {
            "timestamp": datetime.now().isoformat(),
            "ts_code": ts_code,
            "ohlc_data": fetch_ohlc_data(ts_code, start_date, end_date),
            "adj_factor": fetch_adj_factor(ts_code, start_date, end_date)
        }
        
        # Push to Redis queue
        redis_client.lpush(REDIS_QUEUE_NAME, json.dumps(data_package))
        print(f"Data pushed: {ts_code} - {start_date} to {end_date}")

if __name__ == "__main__":
    # Example parameters
    STOCK_CODES = ["000001.SZ", "600519.SH"]  # Example CSI 300 component stocks
    DATE_RANGE = ("20230101", "20231231")   # Annual data range
    
    produce_data(STOCK_CODES, DATE_RANGE)

I ran both files, but there was no response. It turned out that this computer hadn’t downloaded and started Redis. I quickly downloaded and started it manually, and praised 007 for thoughtfully providing a script to check if Redis was running.

At this point, running 007’s "thoughtful code" showed a successful connection to Redis:

To view and manage data in Redis intuitively, I installed RedisInsight.

Fetching Data from Tushare

With the Redis connection established, I directly ran the data_producer.py provided by 007, and successfully saw the data in RedisInsight:

Then, running data_consumer.py successfully retrieved data from the message queue. 007 even explained the relationship between the producer and consumer to me. Another day dependent on my AI assistant!