前言

很多企业内部都有大量业务数据,例如订单数据、销售数据、库存数据、广告数据、客户数据、售后数据等。

但是在真实业务中,经常会出现这样的情况:

运营想知道:

上个月 Amazon 平台销售额是多少?

老板想知道:

哪个产品最近 30 天退款率最高?

客服主管想知道:

本周售后问题最多的产品是什么?

普通业务人员并不会写 SQL,每次都需要找技术同事帮忙查数据库、导报表、写统计逻辑。

这类需求非常适合用 AI 来解决。

如果我们能让用户直接用自然语言提问:

帮我查询 2026 年 5 月 Amazon 平台每个 SKU 的销售额

系统自动完成:

自然语言理解
  ↓
生成 SQL
  ↓
SQL 安全校验
  ↓
查询业务数据库
  ↓
生成分析总结
  ↓
返回表格结果

这就是企业级 AI 数据分析助手。

本文将用 FastAPI 实现一个简化版 AI 数据分析助手,包含:

  • 自然语言转 SQL;
  • SQL 安全校验;
  • SQLite 示例数据库;
  • 查询结果返回;
  • AI 自动生成数据分析结论;
  • 防止危险 SQL 执行;
  • 可扩展的企业级架构设计。

这篇文章适合:

  • AI 应用开发工程师;
  • 后端开发工程师;
  • 数据分析系统开发者;
  • 企业内部 BI 系统开发者;
  • 想把大模型接入业务数据库的技术团队。

系统整体流程图

下面是AI数据分析助手的完整工作流程:

校验失败

查询异常

用户输入自然语言问题

自然语言理解

生成SQL查询语句

SQL安全校验

查询业务数据库

生成分析总结

返回表格结果

返回错误提示

返回数据库错误

一、为什么 AI 数据分析助手值得做?

传统数据分析流程通常是:

业务人员提出问题
  ↓
技术或数据分析师理解需求
  ↓
编写 SQL
  ↓
查询数据库
  ↓
导出 Excel
  ↓
整理图表
  ↓
解释数据含义

这个流程效率很低。

而 AI 数据分析助手可以把流程变成:

业务人员用自然语言提问
  ↓
AI 生成 SQL
  ↓
系统校验 SQL 安全性
  ↓
自动查询数据库
  ↓
AI 总结分析结论

它的价值不是简单“替代 SQL”,而是降低业务人员使用数据的门槛。

但是这里有一个非常重要的问题:

不能让 AI 生成的 SQL 直接执行。

因为大模型可能生成:

DROP TABLE orders;

或者:

DELETE FROM users;

甚至可能查询敏感字段。

所以企业级 AI 数据分析助手的核心不是简单的 Text-to-SQL,而是:

Text-to-SQL + SQL 安全校验 + 权限控制 + 查询结果解释

二、本文项目目标

我们要实现一个接口:

POST /analysis/query

用户提交自然语言问题:

{
  "question": "查询 2026 年 5 月每个平台的销售额"
}

系统返回:

{
  "sql": "SELECT platform, SUM(amount) AS total_sales FROM orders WHERE order_date BETWEEN '2026-05-01' AND '2026-05-31' GROUP BY platform LIMIT 100;",
  "columns": ["platform", "total_sales"],
  "rows": [
    ["Amazon", 32890.5],
    ["Walmart", 15200.0],
    ["Shopify", 8600.0]
  ],
  "analysis": "2026 年 5 月 Amazon 平台销售额最高,占整体销售额的主要部分..."
}

这个项目会包含完整代码。


三、项目技术栈

本文使用:

Python 3.10+
FastAPI
SQLite
Pydantic
httpx
OpenAI-Compatible API

数据库使用 SQLite,方便本地快速运行。

真实生产环境可以替换为:

MySQL
PostgreSQL
SQL Server
ClickHouse
BigQuery
Snowflake
Amazon Redshift

四、项目目录结构

项目结构如下:

ai-data-analyst/
│
├── app/
│   ├── main.py
│   ├── config.py
│   ├── schemas.py
│   │
│   ├── db/
│   │   ├── database.py
│   │   └── init_db.py
│   │
│   ├── llm/
│   │   └── llm_client.py
│   │
│   ├── prompts/
│   │   └── sql_prompt.py
│   │
│   ├── services/
│   │   ├── sql_guard.py
│   │   └── analysis_service.py
│   │
│   └── metadata/
│       └── table_schema.py
│
├── requirements.txt
├── .env
└── README.md

这个目录结构体现了一个企业级 AI 应用的基本分层:

接口层
业务服务层
大模型调用层
数据库访问层
Prompt 管理层
安全校验层
元数据管理层

五、安装依赖

requirements.txt

fastapi==0.115.0
uvicorn==0.30.6
httpx==0.27.2
python-dotenv==1.0.1
pydantic==2.8.2

安装依赖:

pip install -r requirements.txt

六、配置环境变量

.env

LLM_API_KEY=你的API_KEY
LLM_BASE_URL=https://api.openai.com/v1
LLM_MODEL=gpt-4o-mini
DATABASE_PATH=./data.db

这里使用 OpenAI-Compatible API 格式。

如果你使用的是其他兼容 OpenAI 格式的大模型服务,只需要修改:

LLM_BASE_URL
LLM_MODEL
LLM_API_KEY

七、读取配置

app/config.py

import os
from dotenv import load_dotenv

load_dotenv()


class Settings:
    LLM_API_KEY: str = os.getenv("LLM_API_KEY", "")
    LLM_BASE_URL: str = os.getenv("LLM_BASE_URL", "https://api.openai.com/v1")
    LLM_MODEL: str = os.getenv("LLM_MODEL", "gpt-4o-mini")
    DATABASE_PATH: str = os.getenv("DATABASE_PATH", "./data.db")


settings = Settings()

八、创建示例数据库

为了方便演示,我们创建一个简单的订单表。

app/db/database.py

import sqlite3
from app.config import settings


def get_connection():
    return sqlite3.connect(settings.DATABASE_PATH)


def execute_query(sql: str):
    conn = get_connection()
    cursor = conn.cursor()

    try:
        cursor.execute(sql)
        rows = cursor.fetchall()
        columns = [desc[0] for desc in cursor.description] if cursor.description else []
        return columns, rows
    finally:
        cursor.close()
        conn.close()

九、初始化订单数据

app/db/init_db.py

import sqlite3
from app.config import settings


def init_database():
    conn = sqlite3.connect(settings.DATABASE_PATH)
    cursor = conn.cursor()

    cursor.execute("""
    DROP TABLE IF EXISTS orders;
    """)

    cursor.execute("""
    CREATE TABLE orders (
        id INTEGER PRIMARY KEY AUTOINCREMENT,
        order_id TEXT NOT NULL,
        platform TEXT NOT NULL,
        sku TEXT NOT NULL,
        product_name TEXT NOT NULL,
        quantity INTEGER NOT NULL,
        amount REAL NOT NULL,
        refund_amount REAL DEFAULT 0,
        order_date TEXT NOT NULL,
        country TEXT NOT NULL
    );
    """)

    sample_data = [
        ("10001", "Amazon", "TENT-001", "Screen House Tent", 2, 299.98, 0, "2026-05-01", "US"),
        ("10002", "Amazon", "TENT-002", "Bubble Tent", 1, 399.99, 0, "2026-05-03", "US"),
        ("10003", "Walmart", "CHAIR-001", "Camping Chair", 4, 239.96, 0, "2026-05-05", "US"),
        ("10004", "Shopify", "TENT-001", "Screen House Tent", 1, 149.99, 0, "2026-05-07", "US"),
        ("10005", "Amazon", "CHAIR-001", "Camping Chair", 3, 179.97, 59.99, "2026-05-10", "US"),
        ("10006", "Walmart", "TENT-002", "Bubble Tent", 2, 799.98, 0, "2026-05-12", "US"),
        ("10007", "Amazon", "PET-001", "Pet Tent", 5, 499.95, 0, "2026-05-15", "US"),
        ("10008", "Shopify", "PET-001", "Pet Tent", 2, 199.98, 0, "2026-05-18", "US"),
        ("10009", "Amazon", "TENT-001", "Screen House Tent", 3, 449.97, 149.99, "2026-05-20", "US"),
        ("10010", "Walmart", "CHAIR-001", "Camping Chair", 2, 119.98, 0, "2026-05-22", "US"),
        ("10011", "Amazon", "TENT-002", "Bubble Tent", 1, 399.99, 0, "2026-06-01", "US"),
        ("10012", "Shopify", "CHAIR-001", "Camping Chair", 6, 359.94, 0, "2026-06-03", "US")
    ]

    cursor.executemany("""
    INSERT INTO orders (
        order_id, platform, sku, product_name, quantity, amount, refund_amount, order_date, country
    ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?);
    """, sample_data)

    conn.commit()
    cursor.close()
    conn.close()

    return {
        "message": "database initialized",
        "rows": len(sample_data)
    }

这个订单表包含:

平台
SKU
产品名称
销售数量
销售金额
退款金额
订单日期
国家

后面 AI 会根据自然语言自动生成 SQL 查询这些数据。


十、定义数据库元数据

Text-to-SQL 最重要的一点是:必须告诉大模型有哪些表、哪些字段、字段含义是什么。

不要让模型猜数据库结构。

app/metadata/table_schema.py

ORDER_TABLE_SCHEMA = """
数据库中只有一张订单表:orders

表名:orders

字段说明:
- id: 主键,自增 ID
- order_id: 订单号
- platform: 销售平台,例如 Amazon、Walmart、Shopify
- sku: 商品 SKU
- product_name: 产品名称
- quantity: 销售数量
- amount: 销售金额
- refund_amount: 退款金额
- order_date: 订单日期,格式为 YYYY-MM-DD
- country: 国家,例如 US

注意:
1. 只能查询 orders 表;
2. 不允许查询不存在的表;
3. 不允许使用 INSERT、UPDATE、DELETE、DROP、ALTER;
4. 所有查询必须是 SELECT;
5. 如果用户问销售额,通常使用 SUM(amount);
6. 如果用户问退款率,可以使用 SUM(refund_amount) / SUM(amount);
7. 如果用户问销量,可以使用 SUM(quantity);
8. 日期字段使用 order_date。
"""

企业级 Text-to-SQL 应用必须维护数据字典。

否则模型很容易生成不存在的字段,比如:

SELECT sales_amount FROM order_table;

但是实际字段可能叫:

amount

十一、定义请求和响应结构

app/schemas.py

from pydantic import BaseModel, Field
from typing import List, Any


class AnalysisRequest(BaseModel):
    question: str = Field(..., description="用户自然语言问题")


class AnalysisResponse(BaseModel):
    question: str
    sql: str
    columns: List[str]
    rows: List[List[Any]]
    analysis: str

企业系统中,接口返回最好是结构化的。

这样前端可以直接渲染:

SQL
表格
分析结论
图表

十二、封装大模型客户端

app/llm/llm_client.py

import json
import httpx
from typing import Dict, Any
from app.config import settings


class LLMClient:
    def __init__(self):
        self.api_key = settings.LLM_API_KEY
        self.base_url = settings.LLM_BASE_URL.rstrip("/")
        self.model = settings.LLM_MODEL

    async def chat_text(self, system_prompt: str, user_prompt: str) -> str:
        url = f"{self.base_url}/chat/completions"

        headers = {
            "Authorization": f"Bearer {self.api_key}",
            "Content-Type": "application/json"
        }

        payload = {
            "model": self.model,
            "temperature": 0.2,
            "messages": [
                {
                    "role": "system",
                    "content": system_prompt
                },
                {
                    "role": "user",
                    "content": user_prompt
                }
            ]
        }

        async with httpx.AsyncClient(timeout=60) as client:
            response = await client.post(url, headers=headers, json=payload)
            response.raise_for_status()
            result = response.json()

        return result["choices"][0]["message"]["content"]

    async def chat_json(self, system_prompt: str, user_prompt: str) -> Dict[str, Any]:
        content = await self.chat_text(system_prompt, user_prompt)
        return self._parse_json(content)

    @staticmethod
    def _parse_json(content: str) -> Dict[str, Any]:
        content = content.strip()

        if content.startswith("```json"):
            content = content.replace("```json", "").replace("```", "").strip()

        if content.startswith("```"):
            content = content.replace("```", "").strip()

        try:
            return json.loads(content)
        except json.JSONDecodeError:
            raise ValueError(f"模型返回内容不是合法 JSON:{content}")


llm_client = LLMClient()

这里封装了两个方法:

chat_text()

用于普通文本生成。

chat_json()

用于结构化 JSON 输出。


十三、设计 Text-to-SQL Prompt

app/prompts/sql_prompt.py

from app.metadata.table_schema import ORDER_TABLE_SCHEMA


def build_sql_generation_prompt(question: str) -> str:
    return f"""
你是一个企业级数据分析助手,负责把用户的自然语言问题转换成安全的 SQLite SQL 查询语句。

【用户问题】
{question}

【数据库结构】
{ORDER_TABLE_SCHEMA}

请根据用户问题生成 SQL。

要求:
1. 只能生成 SELECT 查询;
2. 只能查询 orders 表;
3. 不允许生成 INSERT、UPDATE、DELETE、DROP、ALTER、TRUNCATE;
4. 不允许使用多条 SQL;
5. 不允许使用分号拼接多条语句;
6. 查询结果最多返回 100 行;
7. 如果用户没有明确要求明细数据,优先生成聚合查询;
8. 返回 JSON,不要输出 Markdown。

JSON 格式如下:

{{
  "sql": "SELECT platform, SUM(amount) AS total_sales FROM orders GROUP BY platform LIMIT 100;",
  "reason": "说明为什么这样生成 SQL"
}}
"""


def build_analysis_prompt(question: str, sql: str, columns: list, rows: list) -> str:
    return f"""
你是一个企业经营数据分析师,请根据用户问题、SQL 查询语句和查询结果,生成简洁清晰的数据分析结论。

【用户问题】
{question}

【SQL】
{sql}

【查询字段】
{columns}

【查询结果】
{rows}

要求:
1. 用中文回答;
2. 先总结核心结论;
3. 如果有明显最高值、最低值、占比或异常,请指出;
4. 不要编造查询结果中不存在的数据;
5. 如果数据量较少,请说明这是基于当前样例数据的分析;
6. 控制在 300 字以内。
"""

这里的 Prompt 有两个阶段:

第一阶段:

自然语言 -> SQL

第二阶段:

SQL 结果 -> 分析结论

这种分阶段设计比一次性让模型完成所有事情更稳定。


十四、SQL 安全校验

SQL安全校验流程图

下面是SQL安全校验的详细流程:

AI生成的SQL语句

SQL规范化处理

是否以SELECT开头?

抛出错误: 只允许SELECT查询

是否包含危险关键字?

抛出错误: 包含禁止关键字

是否包含多语句?

抛出错误: 不允许多语句

是否查询允许的表?

抛出错误: 不允许查询该表

是否包含LIMIT?

自动添加LIMIT 100

SQL校验通过

返回错误给用户

执行SQL查询

这是本文最重要的部分。

企业级 AI 数据分析助手绝对不能直接执行模型生成的 SQL。

必须先做安全校验。

app/services/sql_guard.py

import re


class SQLGuard:
    """
    SQL 安全校验器。
    目标:防止 AI 生成危险 SQL。
    """

    FORBIDDEN_KEYWORDS = [
        "insert",
        "update",
        "delete",
        "drop",
        "alter",
        "truncate",
        "create",
        "replace",
        "attach",
        "detach",
        "pragma",
        "vacuum"
    ]

    ALLOWED_TABLES = ["orders"]

    @classmethod
    def validate(cls, sql: str) -> str:
        if not sql or not sql.strip():
            raise ValueError("SQL 不能为空")

        normalized_sql = sql.strip()

        # 移除末尾分号
        if normalized_sql.endswith(";"):
            normalized_sql = normalized_sql[:-1].strip()

        lower_sql = normalized_sql.lower()

        # 1. 必须以 SELECT 开头
        if not lower_sql.startswith("select"):
            raise ValueError("只允许执行 SELECT 查询")

        # 2. 禁止危险关键字
        for keyword in cls.FORBIDDEN_KEYWORDS:
            pattern = r"\b" + keyword + r"\b"
            if re.search(pattern, lower_sql):
                raise ValueError(f"SQL 包含禁止关键字:{keyword}")

        # 3. 禁止多语句
        if ";" in normalized_sql:
            raise ValueError("不允许执行多条 SQL 语句")

        # 4. 只允许查询 orders 表
        table_names = cls._extract_table_names(lower_sql)
        for table in table_names:
            if table not in cls.ALLOWED_TABLES:
                raise ValueError(f"不允许查询表:{table}")

        # 5. 自动补 LIMIT
        if " limit " not in lower_sql:
            normalized_sql = normalized_sql + " LIMIT 100"

        return normalized_sql + ";"

    @staticmethod
    def _extract_table_names(sql: str):
        """
        简单提取 FROM 和 JOIN 后面的表名。
        生产环境建议使用 SQL Parser。
        """
        tables = []

        from_matches = re.findall(r"\bfrom\s+([a-zA-Z_][a-zA-Z0-9_]*)", sql)
        join_matches = re.findall(r"\bjoin\s+([a-zA-Z_][a-zA-Z0-9_]*)", sql)

        tables.extend(from_matches)
        tables.extend(join_matches)

        return tables

这段代码做了几个关键限制:

只能 SELECT
禁止 DELETE / DROP / UPDATE
禁止多语句
限制只能查 orders 表
自动添加 LIMIT 100

生产环境建议继续增强:

SQL Parser 解析
字段级权限控制
行级权限控制
最大扫描行数限制
查询超时控制
慢 SQL 阻断
敏感字段脱敏

错误处理与用户提示

在 AI 数据分析助手中,良好的错误处理机制至关重要。当 SQL 安全校验失败、数据库查询出错或大模型 API 调用异常时,应该向用户返回友好且安全的错误信息,避免暴露系统内部细节。

FastAPI 异常处理器示例

app/main.pyapp/exception_handlers.py 中添加以下异常处理器:

from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import JSONResponse
from sqlalchemy.exc import SQLAlchemyError
import logging

logger = logging.getLogger(__name__)

app = FastAPI()

# 自定义异常类
class SQLValidationError(Exception):
    """SQL 校验失败异常"""
    pass

class LLMServiceError(Exception):
    """大模型服务异常"""
    pass

class DatabaseQueryError(Exception):
    """数据库查询异常"""
    pass

# 全局异常处理器
@app.exception_handler(SQLValidationError)
async def sql_validation_exception_handler(request: Request, exc: SQLValidationError):
    """SQL 校验失败异常处理"""
    logger.warning(f"SQL 校验失败: {exc}", exc_info=True)
    return JSONResponse(
        status_code=400,
        content={
            "error": "SQL 校验失败",
            "message": str(exc),
            "suggestion": "请检查查询语句是否符合规范,或尝试简化查询条件"
        }
    )

@app.exception_handler(DatabaseQueryError)
async def database_query_exception_handler(request: Request, exc: DatabaseQueryError):
    """数据库查询异常处理"""
    logger.error(f"数据库查询异常: {exc}", exc_info=True)
    return JSONResponse(
        status_code=500,
        content={
            "error": "数据库查询异常",
            "message": "查询执行过程中出现错误",
            "suggestion": "请稍后重试,或联系管理员检查数据库状态"
        }
    )

@app.exception_handler(LLMServiceError)
async def llm_service_exception_handler(request: Request, exc: LLMServiceError):
    """大模型服务异常处理"""
    logger.error(f"大模型服务异常: {exc}", exc_info=True)
    return JSONResponse(
        status_code=503,
        content={
            "error": "AI 服务暂时不可用",
            "message": "大模型服务响应异常",
            "suggestion": "请稍后重试,或尝试简化您的问题描述"
        }
    )

@app.exception_handler(SQLAlchemyError)
async def sqlalchemy_exception_handler(request: Request, exc: SQLAlchemyError):
    """SQLAlchemy 数据库异常处理"""
    logger.error(f"数据库异常: {exc}", exc_info=True)
    
    # 根据异常类型提供更具体的错误信息
    error_msg = str(exc)
    if "permission denied" in error_msg.lower():
        return JSONResponse(
            status_code=403,
            content={
                "error": "权限不足",
                "message": "当前用户没有执行该查询的权限",
                "suggestion": "请联系管理员申请相应权限"
            }
        )
    elif "syntax error" in error_msg.lower():
        return JSONResponse(
            status_code=400,
            content={
                "error": "SQL 语法错误",
                "message": "生成的 SQL 语句存在语法问题",
                "suggestion": "请尝试重新描述您的问题,或联系技术支持"
            }
        )
    elif "timeout" in error_msg.lower() or "timed out" in error_msg.lower():
        return JSONResponse(
            status_code=504,
            content={
                "error": "查询超时",
                "message": "数据库查询执行时间过长",
                "suggestion": "请尝试缩小查询范围,或添加更具体的筛选条件"
            }
        )
    else:
        return JSONResponse(
            status_code=500,
            content={
                "error": "数据库错误",
                "message": "数据库操作过程中出现错误",
                "suggestion": "请稍后重试,或联系管理员检查数据库状态"
            }
        )

@app.exception_handler(Exception)
async def general_exception_handler(request: Request, exc: Exception):
    """通用异常处理"""
    logger.error(f"未处理的异常: {exc}", exc_info=True)
    return JSONResponse(
        status_code=500,
        content={
            "error": "服务器内部错误",
            "message": "系统处理您的请求时出现未知错误",
            "suggestion": "请稍后重试,或联系技术支持"
        }
    )
不同错误类型的最佳处理实践
  1. SQL 安全校验失败(权限不足、危险操作)

    • 状态码: 400 Bad Request
    • 用户提示: “查询语句不符合安全规范”
    • 建议: 明确告知违反的具体规则(如"只允许 SELECT 查询"),但不要暴露完整的校验逻辑
    • 日志记录: 记录详细的校验失败原因,便于安全审计
  2. SQL 语法错误(大模型生成错误 SQL)

    • 状态码: 400 Bad Request
    • 用户提示: “AI 生成的 SQL 存在语法问题”
    • 建议: 提示用户重新描述问题,避免技术术语
    • 处理: 可以尝试让大模型重新生成,但要有重试次数限制
  3. 数据库权限不足

    • 状态码: 403 Forbidden
    • 用户提示: “当前用户没有执行该查询的权限”
    • 建议: 引导用户联系管理员申请权限
    • 安全: 绝对不要透露具体的表名、字段名等敏感信息
  4. 网络超时(数据库或大模型 API)

    • 状态码: 504 Gateway Timeout 或 503 Service Unavailable
    • 用户提示: “服务响应超时,请稍后重试”
    • 建议: 提供重试按钮,设置合理的超时时间(如数据库查询 30 秒,大模型 60 秒)
    • 降级: 考虑返回缓存结果或简化查询
  5. 大模型 API 异常

    • 状态码: 503 Service Unavailable
    • 用户提示: “AI 服务暂时不可用”
    • 建议: 提示用户稍后重试,或尝试使用预设的常见查询模板
    • 容错: 实现服务降级,如切换到备用模型或返回预定义的 SQL 模板
  6. 查询结果过大

    • 状态码: 413 Payload Too Large
    • 用户提示: “查询结果数据量过大”
    • 建议: 提示用户添加更具体的筛选条件或时间范围
    • 防护: 强制添加 LIMIT,限制最大返回行数
错误处理的最佳实践
  1. 分层错误处理

    • SQL 校验层:业务逻辑错误(400)
    • 数据库层:执行错误(500/504)
    • 大模型层:服务错误(503)
    • 网络层:连接错误(502)
  2. 用户友好的消息

    • 避免技术术语和堆栈信息
    • 提供明确的下一步建议
    • 保持语气友好和专业
  3. 安全考虑

    • 不要暴露数据库结构、表名、字段名
    • 不要泄露 SQL 查询的具体内容
    • 记录完整的错误日志供内部排查
  4. 监控与告警

    • 对不同类型的错误设置不同级别的告警
    • 监控错误率,及时发现系统问题
    • 定期分析错误日志,优化系统

通过完善的错误处理机制,不仅能提升用户体验,还能增强系统的安全性和可维护性。

十五、实现分析服务

app/services/analysis_service.py

from app.llm.llm_client import llm_client
from app.prompts.sql_prompt import build_sql_generation_prompt, build_analysis_prompt
from app.services.sql_guard import SQLGuard
from app.db.database import execute_query


class AnalysisService:
    async def analyze(self, question: str):
        """
        主流程:
        1. 用户自然语言问题转 SQL
        2. SQL 安全校验
        3. 执行 SQL
        4. AI 生成分析结论
        """

        # 1. 生成 SQL
        sql_system_prompt = """
你是一个严谨的 Text-to-SQL 助手。
你只能根据提供的数据库结构生成 SQL。
你必须返回合法 JSON。
不要生成任何危险 SQL。
"""

        sql_user_prompt = build_sql_generation_prompt(question)

        sql_result = await llm_client.chat_json(
            system_prompt=sql_system_prompt,
            user_prompt=sql_user_prompt
        )

        raw_sql = sql_result["sql"]

        # 2. SQL 安全校验
        safe_sql = SQLGuard.validate(raw_sql)

        # 3. 执行 SQL
        columns, rows = execute_query(safe_sql)

        # 4. 生成分析结论
        analysis_system_prompt = """
你是一个企业数据分析师。
你只能基于查询结果进行分析,不允许编造数据。
"""

        analysis_user_prompt = build_analysis_prompt(
            question=question,
            sql=safe_sql,
            columns=columns,
            rows=rows
        )

        analysis = await llm_client.chat_text(
            system_prompt=analysis_system_prompt,
            user_prompt=analysis_user_prompt
        )

        return {
            "question": question,
            "sql": safe_sql,
            "columns": columns,
            "rows": [list(row) for row in rows],
            "analysis": analysis
        }


analysis_service = AnalysisService()

这个服务层就是 AI 数据分析助手的核心流程:

自然语言问题
  ↓
LLM 生成 SQL
  ↓
SQLGuard 安全校验
  ↓
数据库查询
  ↓
LLM 生成分析总结
  ↓
返回结果

注意这里有两个 LLM 调用:

第一次用于生成 SQL。

第二次用于解释查询结果。

这比让模型自己“猜数据”要可靠得多。


十六、实现 FastAPI 接口

app/main.py

from fastapi import FastAPI
from app.schemas import AnalysisRequest, AnalysisResponse
from app.services.analysis_service import analysis_service
from app.db.init_db import init_database

app = FastAPI(
    title="AI Data Analyst",
    description="企业级 AI 数据分析助手:自然语言转 SQL + 安全校验 + 分析总结",
    version="1.0.0"
)


@app.get("/health")
async def health_check():
    return {
        "status": "ok",
        "message": "AI Data Analyst is running"
    }


@app.post("/init-db")
async def init_db():
    return init_database()


@app.post("/analysis/query", response_model=AnalysisResponse)
async def analysis_query(request: AnalysisRequest):
    result = await analysis_service.analyze(request.question)
    return AnalysisResponse(**result)

启动服务:

uvicorn app.main:app --reload

访问接口文档:

http://127.0.0.1:8000/docs

十七、初始化数据库

启动项目后,先调用:

curl -X POST "http://127.0.0.1:8000/init-db"

返回:

{
  "message": "database initialized",
  "rows": 12
}

说明示例订单数据已经写入 SQLite。


十八、测试自然语言查询

示例 1:查询各平台销售额

请求:

curl -X POST "http://127.0.0.1:8000/analysis/query" \
-H "Content-Type: application/json" \
-d '{
  "question": "查询 2026 年 5 月每个平台的销售额"
}'

模型可能生成 SQL:

SELECT platform, SUM(amount) AS total_sales
FROM orders
WHERE order_date BETWEEN '2026-05-01' AND '2026-05-31'
GROUP BY platform
LIMIT 100;

返回结果可能是:

{
  "question": "查询 2026 年 5 月每个平台的销售额",
  "sql": "SELECT platform, SUM(amount) AS total_sales FROM orders WHERE order_date BETWEEN '2026-05-01' AND '2026-05-31' GROUP BY platform LIMIT 100;",
  "columns": ["platform", "total_sales"],
  "rows": [
    ["Amazon", 1829.86],
    ["Shopify", 349.97],
    ["Walmart", 1159.92]
  ],
  "analysis": "从当前样例数据看,2026 年 5 月 Amazon 平台销售额最高,为 1829.86;Walmart 排名第二,为 1159.92;Shopify 销售额最低,为 349.97。整体来看,Amazon 是主要销售来源。"
}

示例 2:查询退款率最高的 SKU

请求:

curl -X POST "http://127.0.0.1:8000/analysis/query" \
-H "Content-Type: application/json" \
-d '{
  "question": "2026 年 5 月退款率最高的 SKU 是哪个?"
}'

模型可能生成:

SELECT sku, SUM(refund_amount) / SUM(amount) AS refund_rate
FROM orders
WHERE order_date BETWEEN '2026-05-01' AND '2026-05-31'
GROUP BY sku
ORDER BY refund_rate DESC
LIMIT 100;

这个查询可以帮助运营快速发现异常产品。


示例 3:查询销量最高的产品

请求:

curl -X POST "http://127.0.0.1:8000/analysis/query" \
-H "Content-Type: application/json" \
-d '{
  "question": "统计每个产品的销量,并按销量从高到低排序"
}'

可能生成:

SELECT product_name, SUM(quantity) AS total_quantity
FROM orders
GROUP BY product_name
ORDER BY total_quantity DESC
LIMIT 100;

这个结果适合做产品销售排行。


十九、为什么不能直接执行 AI 生成的 SQL?

这是很多 AI 数据分析 Demo 最大的问题。

错误做法:

sql = llm.generate(question)
cursor.execute(sql)

这种写法非常危险。

因为 AI 可能生成:

DELETE FROM orders;

或者:

DROP TABLE orders;

也可能生成:

SELECT * FROM users;

如果数据库里有客户手机号、地址、邮箱、支付信息,就可能造成数据泄露。

正确做法应该是:

AI 生成 SQL
  ↓
SQL 安全校验
  ↓
权限判断
  ↓
查询限制
  ↓
执行 SQL

本文中的 SQLGuard 就是最基础的安全层。


二十、生产环境中 SQLGuard 应该如何升级?

本文为了便于理解,使用正则做了基础校验。

但生产环境中,建议继续增强。

1. 使用 SQL Parser

正则适合做简单拦截,但复杂 SQL 场景建议使用 SQL Parser。

例如:

sqlglot
sqlparse
Apache Calcite
JSqlParser

使用 SQL Parser 可以更准确地解析:

表名
字段名
where 条件
join 关系
函数调用
子查询

2. 字段级权限控制

不同角色能看的字段不同。

例如:

运营可以看销售额
客服可以看订单状态
财务可以看退款金额
普通员工不能看客户邮箱和地址

权限控制不能只做到表级别,还要做到字段级别。

3. 行级权限控制

不同用户只能看自己权限范围内的数据。

例如:

Amazon 运营只能看 Amazon 数据
Walmart 运营只能看 Walmart 数据
美国团队只能看 US 数据
日本团队只能看 JP 数据

可以在 SQL 中自动追加权限条件:

WHERE platform = 'Amazon'

或者:

WHERE country = 'US'

4. 查询超时控制

AI 可能生成复杂查询,导致数据库压力过大。

生产环境必须限制:

最大执行时间
最大返回行数
最大扫描数据量
最大并发查询数

5. 敏感字段脱敏

即使查询允许执行,返回结果也要脱敏。

例如:

手机号:138****5678
邮箱:a***@example.com
地址:只显示州和城市

二十一、企业级 AI 数据分析助手的完整架构

企业级架构图

下面是生产级AI数据分析助手的完整架构:

安全与监控层

用户问题

用户身份识别

权限系统

数据库元数据管理

Text-to-SQL模型

SQL安全校验

SQL权限改写

数据库查询

结果脱敏

AI分析总结

图表推荐

日志审计

权限控制

安全校验

审计日志

脱敏处理

限流保护

成本统计

慢查询保护

生产级关键模块

和普通Demo相比,生产系统多了以下关键模块:

企业级AI数据分析助手

安全防护

权限控制

SQL安全校验

数据脱敏

访问控制

监控审计

日志记录

操作审计

成本统计

性能监控

性能优化

查询限流

慢查询保护

缓存策略

连接池管理

可扩展性

模块化设计

插件化架构

配置化管理

多租户支持

这就是企业级AI应用开发的核心区别。


二十二、如何让前端展示得更好?

当前接口返回的是:

SQL
columns
rows
analysis

前端可以基于这些数据展示:

  1. 查询 SQL;
  2. 表格结果;
  3. AI 分析结论;
  4. 推荐图表。

例如可以增加一个字段:

{
  "chart_type": "bar"
}

让 AI 根据查询结果推荐图表类型:

平台销售额对比 -> 柱状图
销售趋势 -> 折线图
品类占比 -> 饼图
明细列表 -> 表格

后续可以扩展为:

AI 自动生成 ECharts 配置
AI 自动生成数据看板
AI 自动生成日报周报
AI 自动发现异常数据

这类功能非常适合企业 BI 场景。


二十三、这个项目可以应用在哪些业务场景?

这个 AI 数据分析助手可以扩展到很多企业场景。

1. 电商运营分析

例如:

上个月哪个平台销售额最高?
哪个 SKU 退款率最高?
最近 7 天销量下降最快的产品是什么?
广告花费和销售额是否匹配?

2. 仓库库存分析

例如:

哪些 SKU 库存低于安全库存?
哪个仓库缺货最多?
最近 30 天周转率最低的产品是什么?

3. 客服售后分析

例如:

本周投诉最多的问题类型是什么?
哪个产品售后率最高?
退款金额最高的平台是哪个?

4. 财务经营分析

例如:

本月毛利最高的产品是什么?
哪个平台退款金额最高?
本季度销售额同比增长多少?

5. 管理层经营看板

例如:

总结一下本月经营情况
找出销售异常的 SKU
分析最近 30 天销售额变化

这些场景的共同点是:

业务人员有问题,但不会写 SQL。

AI 数据分析助手正好可以解决这个问题。


二十四、企业落地时的几个关键经验

1. 不要让 AI 直接连生产库

建议中间加一层查询服务。

更安全的方式是:

AI 生成查询意图
  ↓
后端服务生成受控 SQL
  ↓
只读账号查询数据库

数据库账号必须是只读权限。

2. 不要开放所有表

一开始只开放几张高价值、低风险的表。

例如:

订单汇总表
销售统计表
库存汇总表
售后统计表

不要一开始就让 AI 查询全库。

3. 不要让 AI 查询明细隐私数据

AI 数据分析助手更适合查询汇总数据,而不是客户隐私明细。

例如优先支持:

销售额统计
退款率统计
销量排行
库存预警
平台对比

而不是:

查询某个客户的详细地址和电话

4. 数据字典非常重要

Text-to-SQL 的效果很大程度取决于数据字典。

数据字典至少要包含:

表名
字段名
字段含义
枚举值
常用指标口径
字段之间关系
权限说明

5. 必须记录日志

每次查询都应该记录:

用户是谁
问了什么问题
生成了什么 SQL
查询了哪些表
返回了多少行
耗时多久
是否命中敏感字段

这对安全审计非常重要。


十六、性能优化与生产部署建议

在企业级 AI 数据分析助手投入生产环境时,性能优化和稳定部署是确保系统可靠运行的关键。以下是可能遇到的性能瓶颈及相应的优化策略:

1. 大模型 API 延迟优化

问题:大模型 API 调用通常有较高延迟(1-3秒),在高并发场景下会成为系统瓶颈。

优化策略

  • 结果缓存:对常见查询的 SQL 和分析结果进行缓存,使用 Redis 等内存数据库存储,设置合理的 TTL(如5分钟)。
  • 请求合并:对相似查询进行去重和合并,减少重复的 API 调用。
  • 流式响应:对于长文本分析结果,采用流式返回,提升用户体验。
  • 降级策略:当大模型服务不可用时,降级到规则引擎生成简单 SQL 或返回缓存结果。
# Redis 缓存示例
import redis
import json
import hashlib

redis_client = redis.Redis(host='localhost', port=6379, db=0)

def get_cached_sql(natural_language_query: str, db_schema: str) -> Optional[str]:
    """从缓存获取 SQL"""
    cache_key = hashlib.md5(f"{natural_language_query}:{db_schema}".encode()).hexdigest()
    cached = redis_client.get(f"sql:{cache_key}")
    return cached.decode() if cached else None

def cache_sql(natural_language_query: str, db_schema: str, sql: str, ttl: int = 300):
    """缓存 SQL 结果"""
    cache_key = hashlib.md5(f"{natural_language_query}:{db_schema}".encode()).hexdigest()
    redis_client.setex(f"sql:{cache_key}", ttl, sql)

2. 数据库查询性能优化

问题:复杂查询可能导致数据库负载过高,慢查询影响整体响应时间。

优化策略

  • 查询超时控制:为所有数据库查询设置超时时间(如10秒),超时自动取消。
  • 连接池管理:使用数据库连接池,避免频繁创建连接的开销。
  • 查询优化:对生成的 SQL 进行简单优化,如避免 SELECT *、添加合适的索引提示。
  • 只读副本:将查询流量导向只读副本,减轻主库压力。
# SQL 查询超时控制示例
import asyncio
from asyncpg import connect, Connection
from contextlib import asynccontextmanager

@asynccontextmanager
async def query_with_timeout(conn: Connection, sql: str, timeout: int = 10):
    """带超时的查询执行"""
    try:
        # 创建超时任务
        query_task = asyncio.create_task(conn.fetch(sql))
        done, pending = await asyncio.wait([query_task], timeout=timeout)
        
        if query_task in done:
            yield query_task.result()
        else:
            query_task.cancel()
            raise asyncio.TimeoutError(f"查询超时 ({timeout}秒)")
    except asyncio.CancelledError:
        raise asyncio.TimeoutError("查询被取消")

3. 高并发处理

问题:大量用户同时查询时,系统资源可能成为瓶颈。

优化策略

  • API 限流:对自然语言查询接口实施限流,如令牌桶算法。
  • 异步处理:将耗时的 SQL 生成和数据分析转为异步任务,通过消息队列处理。
  • 水平扩展:无状态服务可水平扩展,数据库层考虑分库分表。
  • 负载均衡:使用负载均衡器分发请求到多个服务实例。
# 使用 Celery 进行异步任务处理示例
from celery import Celery
from app.services.analysis_service import generate_sql_and_analyze

# 初始化 Celery
celery_app = Celery('ai_data_assistant', broker='redis://localhost:6379/0')

@celery_app.task
def async_analyze_query(natural_language_query: str, user_id: str):
    """异步处理分析任务"""
    try:
        # 生成 SQL 并执行分析
        result = generate_sql_and_analyze(natural_language_query, user_id)
        
        # 将结果存储到数据库或缓存
        store_analysis_result(user_id, result)
        
        return {"status": "success", "result_id": result.id}
    except Exception as e:
        return {"status": "error", "message": str(e)}

# FastAPI 接口中调用异步任务
@app.post("/api/async-analyze")
async def async_analyze(request: AnalysisRequest):
    """异步分析接口"""
    task = async_analyze_query.delay(request.query, request.user_id)
    return {"task_id": task.id, "status": "processing"}

4. 生产部署建议

基础设施

  • 容器化部署:使用 Docker 容器化应用,便于部署和扩展。
  • 编排管理:使用 Kubernetes 或 Docker Swarm 进行容器编排。
  • 监控告警:集成 Prometheus + Grafana 监控系统指标,设置关键指标告警。
  • 日志集中:使用 ELK 或 Loki 集中管理日志,便于问题排查。

安全加固

  • 网络隔离:将服务部署在内网,通过 API 网关对外暴露。
  • 身份认证:集成企业 SSO 或 OAuth2.0 进行用户认证。
  • 审计日志:记录所有查询请求、生成的 SQL 和执行结果。
  • 定期备份:定期备份缓存数据和系统配置。

成本控制

  • 大模型用量监控:监控 API 调用次数和费用,设置用量告警。
  • 缓存命中率优化:通过分析查询模式优化缓存策略,提高命中率。
  • 自动伸缩:根据负载自动伸缩服务实例,平衡性能和成本。

5. 性能监控指标

建议监控以下关键指标:

  • API 响应时间 P95/P99
  • 大模型 API 调用延迟
  • 数据库查询耗时
  • 缓存命中率
  • 系统并发连接数
  • 错误率和异常率
  • 资源使用率(CPU、内存、磁盘)

通过以上优化策略,可以显著提升企业级 AI 数据分析助手在生产环境中的性能、稳定性和可扩展性,确保系统能够支撑大规模并发查询,同时保持较低的运营成本。

二十五、总结

本文实现了一个简化版企业级 AI 数据分析助手。

它完成了以下能力:

  • 用户自然语言提问;
  • 大模型生成 SQL;
  • SQL 安全校验;
  • 自动执行数据库查询;
  • 返回结构化表格数据;
  • AI 生成中文分析结论;
  • 防止危险 SQL 执行;
  • 支持后续扩展到企业 BI 场景。

这篇文章的重点不是展示一个简单 Demo,而是说明企业级 AI 数据分析应用的核心思想:

AI 不能直接操作数据库,必须放在安全、权限、审计、规则的框架内使用。

从架构角度看,一个真正可落地的 AI 数据分析助手,至少需要:

Text-to-SQL
SQL Guard
权限控制
只读数据库账号
查询限制
结果脱敏
日志审计
AI 分析总结

自然语言查询数据库会成为企业 AI 应用中非常重要的方向。

因为它解决的是一个长期存在的问题:

业务人员想用数据,但不会写 SQL。

AI 的价值,就是把复杂的数据查询能力,变成普通人可以直接使用的自然语言能力。

后续我会继续更新企业级 AI 应用开发实战系列,包括:

  1. AI 数据分析助手如何接入 MySQL;
  2. Text-to-SQL 如何做字段级权限控制;
  3. AI 如何自动生成 ECharts 图表;
  4. 企业级 RAG 知识库如何做权限过滤;
  5. AI Agent 如何安全调用企业内部 API;
  6. 大模型应用如何做日志、监控和成本统计。

关注我,一起系统学习企业级 AI 应用开发。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐