如何创建支持MySQL查询、SQL校验及结果邮件推送的自动化流水线
实现MySQL只读查询自动化流水线方案
1 核心设计原则
优先做双层校验,避免单一层漏判导致安全问题:
- 上层SQL语法校验拦截非允许操作
- 底层数据库账号做权限兜底,完全杜绝写入风险
2 安全校验模块实现
可以用成熟的SQL解析工具实现校验,不需要自行写规则匹配:
- 校验逻辑:
- 第一步:解析SQL语法,确认语句根类型为
SELECT,排除所有INSERT/UPDATE/DELETE/CREATE/DROP/CALL等非查询类语句 - 第二步:提取SQL中所有涉及的表名,和预先配置的表白名单逐一比对,存在非白名单表直接拦截
- 第三步:额外拦截高危语法,比如
INTO OUTFILE、LOAD DATA等可能导出写入的查询后缀
- 第一步:解析SQL语法,确认语句根类型为
- 数据库兜底配置:单独创建MySQL账号,仅授予该账号白名单表的
SELECT权限,不给其他任何库表的操作权限,即使校验逻辑出现漏洞,数据库层面也会拦截非法操作 - 示例(Python环境用sqlparse实现校验):
import sqlparse ALLOWED_TABLES = {"user_info", "order_record", "product_detail"} # 替换为你的白名单表 def sql_security_check(sql: str) -> bool: parsed = sqlparse.parse(sql.strip())[0] # 校验是否为SELECT语句 if parsed.get_type() != "SELECT": return False # 提取所有关联表名 for token in parsed.flatten(): if token.value.upper() in ("FROM", "JOIN", "INNER JOIN", "LEFT JOIN", "RIGHT JOIN"): # 定位表名位置 next_token = token.next while next_token.is_whitespace: next_token = next_token.next table_name = next_token.value.strip("`").strip("'").strip('"') if table_name not in ALLOWED_TABLES: return False return True
3 查询执行模块实现
执行查询时添加异常处理和超时限制,避免慢查询阻塞流水线:
- 配置MySQL连接超时、查询超时参数,超过设定时间直接中断查询
- 对返回结果集大小做限制,超过阈值直接提示缩小查询范围,避免内存溢出
- 示例(Python用pymysql+pandas执行查询):
import pymysql import pandas as pd def execute_query(sql: str) -> pd.DataFrame: conn = pymysql.connect( host="你的MySQL服务地址", user="仅读权限账号", password="账号密码", database="指定查询库名", connect_timeout=10, read_timeout=30 # 配置查询30秒超时 ) try: df = pd.read_sql(sql, conn) return df except Exception as e: raise RuntimeError(f"查询执行失败:{str(e)}") finally: conn.close()
4 结果导出与邮件发送模块实现
查询结果可以按需求导出为CSV或Excel格式,再通过SMTP协议发送邮件:
- 小结果集直接作为邮件附件发送,大结果集可以压缩后发送
- 示例(Python发邮件带附件):
import smtplib from email.mime.multipart import MIMEMultipart from email.mime.text import MIMEText from email.mime.application import MIMEApplication def send_result_email(receiver_email: str, result_df: pd.DataFrame): # 导出结果为Excel result_df.to_excel("查询结果.xlsx", index=False) # 构造邮件内容 msg = MIMEMultipart() msg["From"] = "发件人邮箱地址" msg["To"] = receiver_email msg["Subject"] = "MySQL自动查询结果通知" # 添加邮件正文 msg.attach(MIMEText("本次查询执行完成,结果详见附件。", "plain", "utf-8")) # 添加附件 with open("查询结果.xlsx", "rb") as f: part = MIMEApplication(f.read()) part.add_header("Content-Disposition", "attachment", filename="查询结果.xlsx") msg.attach(part) # 发送邮件 smtp = smtplib.SMTP_SSL("你的邮箱SMTP服务器地址", 465) smtp.login("发件人邮箱地址", "邮箱授权码") smtp.sendmail("发件人邮箱地址", receiver_email, msg.as_string()) smtp.quit()
5 流水线编排
把上述三个模块串起来即可组成完整流水线:
- 接收输入的SQL语句
- 调用
sql_security_check做安全校验,校验不通过直接返回错误 - 校验通过后调用
execute_query执行查询 - 查询成功后调用
send_result_email发送结果邮件
- 定时执行需求可以用Linux crontab、Windows任务计划程序配置定时触发,或者用Airflow等调度工具实现更复杂的流水线依赖配置。
内容的提问来源于stack exchange,提问作者Kani
相关产品推荐
相关产品推荐

