Python并行处理CSV时部分输出文件为空且未执行预期代码路径的问题排查求助
Python并行处理CSV时部分输出文件为空且未执行预期代码路径的问题排查求助
我现在碰到一个百思不得其解的问题:用Python并行处理一批CSV文件,每个输入CSV都确认有200行,所以我加了if num_videos == 200:的判断来执行后续的批量处理和保存逻辑。
正常情况下,处理完成的CSV会包含大量数据,同时控制台会打印process_videos_in_batches里的调试信息(比如Debug: {len(all_data)})和save_all_data_to_csv里的保存提示。但现在的状况很奇怪:
- 部分文件完全符合预期:调试输出正常,输出CSV内容完整
- 另一些文件却生成了空CSV,而且对应的调试打印完全没出现——看起来代码根本没走完那个判断分支,但空文件又确实被创建出来了
我已经做了这些排查:
- 确认所有输入CSV都确实有200行,
if条件应该都能满足 - 输出文件名都是唯一的,不存在文件覆盖的问题
- 用htop检查过系统状态,内存和CPU都没有耗尽的情况
我的运行环境:
- Python: 3.10.15
- Conda: 24.9.2
- VSCode: 1.96
- OS: Ubuntu 22.04.5 LTS
会不会是并行执行的机制导致某些进程提前退出或者静默失败了?明明if条件满足,却没走后续逻辑,还生成了空文件,实在想不通,求各位大佬给点思路或建议!
完整代码片段
process_videos_in_batches函数
def process_videos_in_batches(video_id_list, batch_size, API_KEY): all_data = [] all_video_ids = video_id_list[:] while all_video_ids: batch_video_ids = all_video_ids[:batch_size] all_video_ids = all_video_ids[batch_size:] videos_info = get_video_data(batch_video_ids, API_KEY) ### Debug ### print(f'Debug: {videos_info}, {batch_video_ids}') ############# # videos_infoがNoneの場合は次のバッチへ if videos_info is None: continue for video_id in batch_video_ids: video_data = videos_info.get(video_id) if videos_info else None if video_data: comments = get_all_video_comments(video_id, API_KEY) for comment in comments: # 親コメントの情報を展開 video_entry = video_data.copy() video_entry.update({ 'comment': comment.get('comment', pd.NA), 'comment_like_count': comment.get('comment_like_count', pd.NA), 'comment_author_channel_id': comment.get('comment_author_channel_id', pd.NA), 'comment_published_at': comment.get('comment_published_at', pd.NA), 'updated_at': comment.get('updated_at', pd.NA), 'has_reply': comment.get('has_reply', 'no'), 'reply_comment': comment.get('reply_comment', pd.NA), 'reply_like_count': comment.get('reply_like_count', pd.NA), 'reply_author_channel_id': comment.get('reply_author_channel_id', pd.NA), 'reply_published_at': comment.get('reply_published_at', pd.NA), 'reply_updated_at': comment.get('reply_updated_at', pd.NA) }) all_data.append(video_entry) ### Debug ### print(f'Debug: {len(all_data)}') print(f'Debug: {all_data[:1]}') ############# return all_data
save_all_data_to_csv函数
def save_all_data_to_csv(all_data, output_csv_path): # DataFrame に変換 final_df = pd.DataFrame(all_data) # CSV に保存 final_df.to_csv(output_csv_path, index=False, encoding='utf-8') ### Debug ### print("Debug: writing CSV to", output_csv_path) #############
process_single_file函数
def process_single_file(file_path, output_dir, batch_size, API_KEY): """ 1つのCSVファイルを処理し、必要に応じて結果をCSVで出力する。 """ df = pd.read_csv(file_path, encoding='utf-8') video_ids = list(df['related_video_id']) num_video = len(video_ids) if num_video == 200: all_data = process_videos_in_batches(video_ids, batch_size, API_KEY) output_csv_path = os.path.join(output_dir, os.path.basename(file_path)) save_all_data_to_csv(all_data, output_csv_path) # 終わったらファイルパスを返す(ログ代わり) return file_path
主进程逻辑
input_files = os.listdir(root_dir) file_paths = [os.path.join(root_dir, f) for f in input_files if os.path.isfile(os.path.join(root_dir, f))] # 並列実行する num_workers = 4 futures = [] with ProcessPoolExecutor(max_workers=num_workers) as executor: for file_path in file_paths: # process_single_file関数を並列に実行 future = executor.submit(process_single_file, file_path, output_dir, batch_size, API_KEY) futures.append(future) # 処理が終わったら順に結果を取得 # tqdmを使って「終わった数」を確認するときは、as_completedに対して進捗バーを回します for f in tqdm(as_completed(futures)): try: finished_file = f.result() except Exception as e: print("Error happened in worker:", e)
控制台输出示例
Debug: {'lAtasG8EVEg': {'video_id': 'lAtasG8EVEg'... ... Debug: 22058 Debug: [{'video_id': '8BtA6fO93_w'... Debug: writing CSV to /home/foo/mnt/vt/related_videos/data_info/C_dsxOR9JJw.csv
最小可复现示例
import os import time from concurrent.futures import ProcessPoolExecutor def example_task(x): time.sleep(1) # Simulate some work return x ** 2 def process_files(files, num_workers): with ProcessPoolExecutor(max_workers=num_workers) as executor: results = list(executor.map(example_task, files)) return results if __name__ == "__main__": # Input data to simulate files files = [1, 2, 3, 4, 5] # Test with multiprocessing num_workers = 4 # Adjust the number of workers try: output = process_files(files, num_workers) print("Output:", output) except Exception as e: print("Error:", e)
完整代码
import os import requests import pandas as pd import time from tqdm.notebook import tqdm from datetime import datetime, timedelta, timezone from concurrent.futures import ProcessPoolExecutor, as_completed # Retrieve API key and paths from environment variables API_KEY = os.getenv('YOUTUBE_API_KEY') batch_size = 50 mnt_path = os.getenv('MNT_PATH') root_dir = os.path.join(mnt_path, 'related_videos') # Create directories for thumbnails and output data thumbnail_dir = os.path.join(root_dir, 'thumbnails') os.makedirs(thumbnail_dir, exist_ok=True) output_dir = os.path.join(root_dir, 'data_info') os.makedirs(output_dir, exist_ok=True) # Check if the output directory exists if not os.path.exists(output_dir): print(f"Failed to create directory: {output_dir}") # Print the current time in JST (Japan Standard Time) def show_now_time(): jst = timezone(timedelta(hours=9)) jst_time = datetime.now(jst) print(jst_time) # Convert UTC time string to JST time string def convert_to_jst(utc_time_str): utc_time = datetime.strptime(utc_time_str, "%Y-%m-%dT%H:%M:%SZ") jst_time = utc_time + timedelta(hours=9) return jst_time.strftime("%Y-%m-%d %H:%M:%S") # Download a video thumbnail and save it locally def download_thumbnail(thumbnail_url, thumbnail_dir, video_id): try: response = requests.get(thumbnail_url) response.raise_for_status() # Raise an exception for HTTP errors except requests.exceptions.RequestException as e: print(f"Error occurred while downloading thumbnail: {e}") return None if response.ok: # Extract the file extension (e.g., jpg, png) file_extension = thumbnail_url.split('.')[-1] if len(file_extension) > 5: # Handle overly long extensions (e.g., URL parameters) file_extension = 'jpg' # Default to jpg # Construct the file path for saving file_path = os.path.join(thumbnail_dir, f"{video_id}.{file_extension}") try: # Save the image data to the file with open(file_path, 'wb') as file: file.write(response.content) return file_path except IOError as io_error: print(f"Error occurred while saving the file: {io_error}") return None else: print(f"Unexpected status code received: {response.status_code}") return None
备注:内容来源于stack exchange,提问作者ララララ
相关产品推荐
相关产品推荐

