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

Gmail API多进程解析邮件返回空列表问题求助

Gmail邮件批量解析的多进程并发问题解决

问题背景

使用Google API及相关库成功提取Gmail账户中邮件的主题、发件人、日期、正文信息,但解析效率随邮件数量增长急剧下降——解析34封邮件耗时近15秒,无法扩展到千级规模。尝试用ProcessPoolExecutor对parse_message()函数做并发/多进程处理后,始终得到空的combined列表,需要实现批量解析邮件并汇总到列表的功能。

原代码及运行输出:

from __future__ import print_function
import os.path
from google.auth.transport.requests import Request
from google.oauth2.credentials import Credentials
from google_auth_oauthlib.flow import InstalledAppFlow
from googleapiclient.discovery import build
from concurrent.futures import ProcessPoolExecutor
import base64
import re

combined = []

def authenticate():
    # If modifying these scopes, delete the file token.json.
    SCOPES = ['https://www.googleapis.com/auth/gmail.readonly']

    creds = None

    if os.path.exists('token.json'):
        creds = Credentials.from_authorized_user_file('token.json', SCOPES)

    if not creds or not creds.valid:
        if creds and creds.expired and creds.refresh_token:
            creds.refresh(Request())
        else:
            flow = InstalledAppFlow.from_client_secrets_file(
                'creds.json', SCOPES)
            creds = flow.run_local_server(port=0)

        with open('token.json', 'w') as token:
            token.write(creds.to_json())
    return creds

def get_messages(creds):
    # Get the messages
    days = 31
    service = build('gmail', 'v1', credentials=creds)
    results = service.users().messages().list(userId='me', q=f'newer_than:{days}d, in:inbox').execute()
    messages = results.get('messages', [])
    message_count = len(messages)
    print(f"You've received {message_count} email(s) in the last {days} days")
    if not messages:
        print(f'No Emails found in the last {days} days.')
    return messages


def parse_message(msg):
    # Call the Gmail API
    service = build('gmail', 'v1', credentials=creds)
    txt = service.users().messages().get(userId='me', id=msg['id']).execute()
    payload = txt['payload']
    headers = payload['headers']

    #Grab the Subject Line, From and Date from the Email
    for d in headers:
        if d['name'] == 'Subject':
            subject = d['value']
        if d['name'] == 'From':
            sender = d['value']
            try:
                match = re.search(r'<(.*)>', sender).group(1)
            except:
                match = sender
        if d['name'] == "Date":
            date_received = d['value']

    def get_body(payload):
        if 'body' in payload and 'data' in payload['body']:
            return payload['body']['data']
        elif 'parts' in payload:
            for part in payload['parts']:
                data = get_body(part)
                if data:
                    return data
        else:
            return None

    data = get_body(payload)

    data = data.replace("-","+").replace("_","/")
    decoded_data = base64.b64decode(data).decode("UTF-8")
    decoded_data = (decoded_data.encode('ascii', 'ignore')).decode("UTF-8")
    decoded_data = decoded_data.replace('\n','').replace('\r','').replace('\t', '')

    # Append parsed message to shared list
    return combined.append([date_received, subject, match, decoded_data])

if __name__ == '__main__':
    creds = authenticate()
    messages = get_messages(creds)
    # Create a process pool with 4 worker processes
    with ProcessPoolExecutor(max_workers=4) as executor:
        # Submit the parse_message function for each message in the messages variable
        executor.map(parse_message, messages)
   
    print(f"Combined: {combined}")

运行输出:

You've received 34 email(s) in the last 31 days
Combined: []

问题原因

  1. 多进程内存隔离:Python多进程中,每个子进程会复制主进程的内存空间,主进程的combined列表在子进程中是独立副本,子进程对副本的修改不会同步到主进程的combined。
  2. 全局变量跨进程无效:parse_message中直接使用主进程的全局变量creds,子进程无法正确继承该变量,可能导致API调用失败。
  3. 返回值错误:parse_message返回的是combined.append(...)的结果(append方法返回None),没有实际返回解析后的邮件数据。

修复方案

调整后的代码

from __future__ import print_function
import os.path
from google.auth.transport.requests import Request
from google.oauth2.credentials import Credentials
from google_auth_oauthlib.flow import InstalledAppFlow
from googleapiclient.discovery import build
from concurrent.futures import ProcessPoolExecutor
import base64
import re
import pickle

def authenticate():
    SCOPES = ['https://www.googleapis.com/auth/gmail.readonly']
    creds = None
    if os.path.exists('token.json'):
        creds = Credentials.from_authorized_user_file('token.json', SCOPES)
    if not creds or not creds.valid:
        if creds and creds.expired and creds.refresh_token:
            creds.refresh(Request())
        else:
            flow = InstalledAppFlow.from_client_secrets_file(
                'creds.json', SCOPES)
            creds = flow.run_local_server(port=0)
        with open('token.json', 'w') as token:
            token.write(creds.to_json())
    return creds

def get_messages(creds):
    days = 31
    service = build('gmail', 'v1', credentials=creds)
    results = service.users().messages().list(userId='me', q=f'newer_than:{days}d, in:inbox').execute()
    messages = results.get('messages', [])
    message_count = len(messages)
    print(f"最近{days}天收到{message_count}封邮件")
    if not messages:
        print(f"最近{days}天未找到邮件")
    return messages

def parse_message(args):
    msg, creds_data = args
    # 从序列化数据重建Credentials
    creds = Credentials.from_authorized_user_info(pickle.loads(creds_data))
    service = build('gmail', 'v1', credentials=creds)
    txt = service.users().messages().get(userId='me', id=msg['id']).execute()
    payload = txt['payload']
    headers = payload['headers']

    subject = ""
    match = ""
    date_received = ""
    for d in headers:
        if d['name'] == 'Subject':
            subject = d['value']
        if d['name'] == 'From':
            sender = d['value']
            try:
                match = re.search(r'<(.*)>', sender).group(1)
            except:
                match = sender
        if d['name'] == "Date":
            date_received = d['value']

    def get_body(payload):
        if 'body' in payload and 'data' in payload['body']:
            return payload['body']['data']
        elif 'parts' in payload:
            for part in payload['parts']:
                data = get_body(part)
                if data:
                    return data
        else:
            return None

    data = get_body(payload)
    if not data:
        decoded_data = ""
    else:
        data = data.replace("-","+").replace("_","/")
        decoded_data = base64.b64decode(data).decode("UTF-8")
        decoded_data = (decoded_data.encode('ascii', 'ignore')).decode("UTF-8")
        decoded_data = decoded_data.replace('\n','').replace('\r','').replace('\t', '')

    # 返回解析后的邮件数据,而非修改全局列表
    return [date_received, subject, match, decoded_data]

if __name__ == '__main__':
    creds = authenticate()
    messages = get_messages(creds)
    # 序列化Credentials,以便在子进程中重建
    creds_data = pickle.dumps(creds.to_json())
    
    combined = []
    with ProcessPoolExecutor(max_workers=4) as executor:
        # 将msg和creds数据作为参数传入,收集每个进程的返回值
        results = executor.map(parse_message, [(msg, creds_data) for msg in messages])
        # 将结果汇总到combined列表
        combined.extend(results)
    
    print(f"汇总结果长度:{len(combined)}")
    # 可打印部分结果验证
    if combined:
        print("第一条邮件数据:", combined[0])

关键修改点

  1. 传递Credentials:将主进程的creds序列化为数据,作为参数传递给子进程的parse_message,子进程内部重建Credentials实例,避免全局变量跨进程问题。
  2. 返回解析数据:parse_message不再修改全局列表,而是直接返回解析后的邮件条目,主进程通过executor.map收集所有返回值并汇总到combined。
  3. 处理空body情况:增加对data为空的判断,避免解码时抛出异常。
  4. 内存隔离适配:利用多进程返回值的特性,绕过子进程与主进程的内存隔离问题,不再依赖共享全局变量。

进一步优化建议

  • 批量API请求:使用Google API的批量请求(Batch Requests),一次性获取多封邮件的详情,减少HTTP请求次数,大幅提升效率。
  • 限制进程池大小:Gmail API有请求频率限制,进程池大小建议设置为4-8,避免触发限流。
  • 异步请求:改用ThreadPoolExecutor(IO密集型任务更适合线程池),减少进程创建的开销,同时避免多进程的内存复制成本。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 18:50:38