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: []
问题原因
- 多进程内存隔离:Python多进程中,每个子进程会复制主进程的内存空间,主进程的
combined列表在子进程中是独立副本,子进程对副本的修改不会同步到主进程的combined。 - 全局变量跨进程无效:
parse_message中直接使用主进程的全局变量creds,子进程无法正确继承该变量,可能导致API调用失败。 - 返回值错误:
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])
关键修改点
- 传递Credentials:将主进程的
creds序列化为数据,作为参数传递给子进程的parse_message,子进程内部重建Credentials实例,避免全局变量跨进程问题。 - 返回解析数据:
parse_message不再修改全局列表,而是直接返回解析后的邮件条目,主进程通过executor.map收集所有返回值并汇总到combined。 - 处理空body情况:增加对
data为空的判断,避免解码时抛出异常。 - 内存隔离适配:利用多进程返回值的特性,绕过子进程与主进程的内存隔离问题,不再依赖共享全局变量。
进一步优化建议
- 批量API请求:使用Google API的批量请求(Batch Requests),一次性获取多封邮件的详情,减少HTTP请求次数,大幅提升效率。
- 限制进程池大小:Gmail API有请求频率限制,进程池大小建议设置为4-8,避免触发限流。
- 异步请求:改用
ThreadPoolExecutor(IO密集型任务更适合线程池),减少进程创建的开销,同时避免多进程的内存复制成本。
内容的提问来源于stack exchange,提问作者Jered
相关产品推荐
相关产品推荐

