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

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,提问作者ララララ

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 13:53:09