Python:从实时更新的DataFrame筛选子集并触发邮件通知
持续更新DataFrame并触发邮件通知的解决方案
原代码核心问题
scheduled_update函数中return data写在while True循环内部,导致线程仅执行一次就终止,无法实现持续更新- 数据筛选与去重顺序错误,先筛选再去重会遗漏有效数据,且未跟踪新增的符合条件的行,无法触发邮件
- 邮件发送逻辑未与数据更新流程关联,没有在新符合条件的行出现时自动调用
修正后的完整代码
import pandas as pd import yfinance as yf import datetime import pytz from threading import Thread from time import sleep from email.mime.application import MIMEApplication from email.mime.multipart import MIMEMultipart from email.mime.text import MIMEText import smtplib # 邮件发送函数 def send_tradeNotification(send_to, subject, df): send_from = 'xxxx1@gmail.com' password = 'password' # 建议使用Google App Password而非明文密码 message = """\ <p><strong>交易通知</strong></p> <p>发现符合条件的EURUSD行情数据,请查看附件。</p> <p><strong>此致</strong></p> """ for receiver in send_to: multipart = MIMEMultipart() multipart['From'] = send_from multipart['To'] = receiver multipart['Subject'] = subject attachment = MIMEApplication(df.to_csv()) attachment['Content-Disposition'] = f'attachment; filename="{subject}.csv"' multipart.attach(attachment) multipart.attach(MIMEText(message, 'html')) try: server = smtplib.SMTP('smtp.gmail.com', 587) server.starttls() server.login(multipart['From'], password) server.sendmail(multipart['From'], multipart['To'], multipart.as_string()) print(f"通知邮件已发送至 {receiver}") except Exception as e: print(f"邮件发送失败: {str(e)}") finally: server.quit() def scheduled_update(): # 时区设置 tz = pytz.timezone('Etc/GMT-5') # 初始化过去24小时数据 current_time = datetime.datetime.now(tz) prev_24hrs = current_time - datetime.timedelta(hours=25) # 获取初始数据并筛选符合条件的行 data = yf.download(tickers='EURUSD=X', start=prev_24hrs, end=current_time, interval='1m').iloc[:-1] # 移除最后一行可能不完整的数据 # 跟踪已处理过的索引,避免重复发送邮件 processed_indices = set(data.index) # 筛选初始符合条件的数据(如果需要首次发送) initial_matches = data[data['High'] < 0.98000] if not initial_matches.empty: send_tradeNotification(['xxxx2@gmail.com'], '初始符合条件的EURUSD行情', initial_matches) # 持续更新数据 while True: sleep(60) current_time = datetime.datetime.now(tz) # 获取最近5分钟的数据(覆盖可能的延迟) upd_data = yf.download(tickers='EURUSD=X', start=current_time - datetime.timedelta(minutes=5), end=current_time, interval='1m').iloc[:-1] if len(upd_data) == 0: print('无新数据') continue # 过滤掉已处理过的旧数据 new_rows = upd_data[~upd_data.index.isin(processed_indices)] if new_rows.empty: print('无新增数据行') continue # 筛选新增行中符合条件的数据 new_matches = new_rows[new_rows['High'] < 0.98000] # 更新主数据集 data = pd.concat([data, new_rows]) data = data[~data.index.duplicated(keep='first')] # 更新已处理索引 processed_indices.update(new_rows.index) # 如果有符合条件的新增行,发送邮件 if not new_matches.empty: print(f"发现{len(new_matches)}条符合条件的新数据,触发邮件通知") send_tradeNotification(['xxxx2@gmail.com'], 'EURUSD交易信号通知', new_matches) else: print('新增数据无符合条件的行') # 启动线程 thread = Thread(target=scheduled_update) thread.start()
关键修改说明
- 修复线程终止问题:移除
while True循环内的return data,让线程持续运行 - 新增重复处理防护:用
processed_indices集合跟踪已处理的行索引,确保只处理新增数据 - 优化数据流程:先获取新数据→过滤旧数据→筛选符合条件的新增行→更新主数据集→触发邮件,逻辑更清晰
- 邮件逻辑关联:在发现符合条件的新增行时自动调用
send_tradeNotification函数 - 错误处理优化:给邮件发送逻辑增加异常捕获,避免因邮件发送失败导致整个更新流程终止
- 时区与时间范围修正:统一使用指定时区,避免时间偏移导致的数据重复或遗漏
内容的提问来源于stack exchange,提问作者jonboy
相关产品推荐
相关产品推荐

