如何暂停数据流循环,每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:
- Thread Safety with Locks: The
data_lockensures that we never try to append ticks todfwhile we're copying/resetting it for processing. This prevents messy race conditions that could corrupt your data. - Background Processing Thread: The
periodic_processorruns in the background, checking every 10 seconds if it's time to process data. It doesn't block your WebSocket handler at all. - Snapshot Processing: We make a copy of the current
dffor processing, then reset the maindfimmediately. This means you start collecting new ticks right away, with no downtime or data loss. - Why
sleepFailed: Puttingsleepinsideon_tickshalts 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
相关产品推荐
相关产品推荐

