You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何暂停数据流循环,每5分钟导出DataFrame处理后再续收?

Solution: Periodically Process Streaming Data Without Blocking Collection

Got it, let's fix this. The main issue with using sleep directly is that it blocks your WebSocket tick handler—this means you'll miss incoming ticks or even get disconnected from the server while waiting. Instead, we need to separate data collection from periodic processing using threads and locks to keep everything safe and responsive.

Here's a modified version of your code that does exactly what you want:

import pandas as pd
import threading
import time

# Initialize global variables with thread safety
df = pd.DataFrame(data=None, columns=['name', 'ltp'])
last_process_time = time.time()
data_lock = threading.Lock()  # Prevents race conditions between threads
PROCESS_INTERVAL = 5 * 60  # 5 minutes in seconds

def process_snapshot(snapshot_df):
    """Your existing data processing logic, applied to a snapshot of the data"""
    snapshot_df['change'] = snapshot_df.groupby('name')['ltp'].pct_change() * 100
    g = snapshot_df.groupby('name')['change']
    counts = g.agg(
        pos_count=lambda s: s.gt(0).sum(),
        neg_count=lambda s: s.lt(0).sum(),
        net_count=lambda s: s.gt(0).sum() - s.lt(0).sum()
    ).astype(int)
    print(counts)

def periodic_processor():
    """Background thread that checks for processing intervals and runs the logic"""
    global df, last_process_time
    while True:
        current_time = time.time()
        # Check if 5 minutes have passed since last processing
        if current_time - last_process_time >= PROCESS_INTERVAL:
            with data_lock:
                # Take a snapshot of the current data and reset the main df
                data_snapshot = df.copy()
                df = pd.DataFrame(data=None, columns=['name', 'ltp'])
                last_process_time = current_time
            # Process the snapshot outside the lock to avoid blocking collection
            if not data_snapshot.empty:
                print(f"Processing {len(data_snapshot)} rows collected in the last 5 minutes...")
                process_snapshot(data_snapshot)
        # Check every 10 seconds to avoid unnecessary CPU usage
        time.sleep(10)

def on_ticks(ws, ticks):
    global df
    with data_lock:
        # Append new ticks to the df safely (no race conditions)
        for sc in ticks:
            token = sc['instrument_token']
            name = trd_portfolio[token]['name']
            ltp = sc['last_price']
            df1 = pd.DataFrame([[name, ltp]], columns=['name', 'ltp'])
            df = pd.concat([df, df1], ignore_index=True)
    # Optional: Print progress without blocking
    # print(f"Current df size: {len(df)} rows")

# Start the periodic processing thread before launching the WebSocket
processing_thread = threading.Thread(target=periodic_processor, daemon=True)
processing_thread.start()

# Your existing WebSocket setup code goes here (e.g., ws.run())
# ws = WebSocket(...)
# ws.on_ticks = on_ticks
# ws.run()

Key Explanations:

  1. Thread Safety with Locks: The data_lock ensures that we never try to append ticks to df while we're copying/resetting it for processing. This prevents messy race conditions that could corrupt your data.
  2. Background Processing Thread: The periodic_processor runs in the background, checking every 10 seconds if it's time to process data. It doesn't block your WebSocket handler at all.
  3. Snapshot Processing: We make a copy of the current df for processing, then reset the main df immediately. This means you start collecting new ticks right away, with no downtime or data loss.
  4. Why sleep Failed: Putting sleep inside on_ticks halts the callback function, so your WebSocket client can't handle new ticks until the sleep finishes. This leads to missed data or server disconnections—using a separate thread avoids this entirely.

内容的提问来源于stack exchange,提问作者push kala

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.07 19:27:54