批量处理1分钟时间序列:确保起止为5分钟倍数的可扩展方案
Solution for 5-Minute Aligned Time Series Data Processing
1. Single DataFrame Processing Logic
First, ensure your timestamp column is parsed as a datetime type, then filter rows to start from the first timestamp that falls on a 5-minute interval (e.g., 0:00, 0:05, 0:10, etc.).
Here's the code snippet for a single DataFrame:
import pandas as pd # Sample DataFrame (replace with your data loading step) data = { 'timestamp': ['2020-10-10 0:02', '2020-10-10 0:03', '2020-10-10 0:04', '2020-10-10 0:05', '2020-10-10 0:06'], 'col': [10, 20, 12, 30, 25] } df = pd.DataFrame(data) # Convert timestamp column to datetime type df['timestamp'] = pd.to_datetime(df['timestamp']) # Find the first timestamp that's a multiple of 5 minutes first_valid_ts = df[df['timestamp'].dt.minute % 5 == 0]['timestamp'].iloc[0] # Filter the DataFrame to keep rows from the first valid timestamp onwards processed_df = df[df['timestamp'] >= first_valid_ts] print(processed_df)
This will output:
timestamp col 3 2020-10-10 00:05:00 30 4 2020-10-10 00:06:00 25
2. Scalable Batch Processing for Multiple Files
To handle a folder of files, wrap the single-file logic into reusable functions, then loop through all files in the input directory. Save processed files to a separate output folder to avoid overwriting original data.
Full Batch Processing Code
import pandas as pd import os from multiprocessing import Pool # Optional, for parallel processing def process_single_file(file_info): input_path, output_folder = file_info try: # Load the file (adjust read method if using parquet/excel) df = pd.read_csv(input_path, parse_dates=['timestamp']) df['timestamp'] = pd.to_datetime(df['timestamp']) # Skip if no valid 5-minute timestamps exist valid_timestamps = df[df['timestamp'].dt.minute % 5 == 0] if valid_timestamps.empty: print(f"Skipping {input_path}: No 5-minute interval timestamps found") return first_valid_ts = valid_timestamps['timestamp'].iloc[0] processed_df = df[df['timestamp'] >= first_valid_ts] # Save processed file os.makedirs(output_folder, exist_ok=True) output_path = os.path.join(output_folder, os.path.basename(input_path)) processed_df.to_csv(output_path, index=False) print(f"Processed: {input_path} -> {output_path}") except Exception as e: print(f"Error processing {input_path}: {str(e)}") def batch_process_files(input_folder, output_folder, use_parallel=False): # Collect all file paths in input folder file_paths = [os.path.join(input_folder, f) for f in os.listdir(input_folder) if os.path.isfile(os.path.join(input_folder, f))] # Prepare file info tuples for processing file_info_list = [(path, output_folder) for path in file_paths] if use_parallel: # Use multiprocessing for faster processing of large file sets with Pool() as pool: pool.map(process_single_file, file_info_list) else: # Sequential processing for smaller file sets for info in file_info_list: process_single_file(info) # Example usage if __name__ == "__main__": INPUT_FOLDER = "/path/to/your/input/files" OUTPUT_FOLDER = "/path/to/processed/output" batch_process_files(INPUT_FOLDER, OUTPUT_FOLDER, use_parallel=True)
Key Notes:
- File Format Flexibility: If your files are in Parquet or Excel format, replace
pd.read_csv/to_csvwithpd.read_parquet/to_parquetorpd.read_excel/to_excel. - Parallel Processing: Enable
use_parallel=Trueto speed up processing for hundreds/thousands of files (utilizes multiple CPU cores). - Error Handling: The code includes basic error catching to skip problematic files and log issues.
- Edge Case Handling: Checks for files with no valid 5-minute timestamps and skips them with a message.
内容的提问来源于stack exchange,提问作者prof32
相关产品推荐
相关产品推荐

