欢迎光临
我们一直在努力

AI 工具搭建工作流全流程教程(含代码 / 操作指引)

前言:AI 工作流的价值与应用场景

在数字化办公中,重复的数据录入、文本分类、报告生成等工作占据大量人力成本。AI 工作流通过 “数据源接入→AI 模型处理→自动化执行→结果输出” 的闭环,可将此类工作效率提升 50%-80%。例如:

  • 客户服务:自动提取邮件 / 聊天中的客户诉求,分类后分配给对应部门
  • 数据分析:每日自动从数据库拉取数据,经 AI 生成分析结论后推送日报
  • 内容创作:根据产品信息自动生成宣传文案,按渠道格式调整后发布

本教程将覆盖 “开发型”(需代码)与 “低代码型”(可视化操作)两种搭建方式,适配技术 / 非技术人员需求,所有代码均经过实测可运行,操作步骤附带界面指引。

一、AI 工作流基础认知

1.1 什么是 AI 工作流?

AI 工作流是融合 “数据采集、AI 处理、自动化逻辑” 的协作体系,核心是让 AI 替代人工完成 “判断、分析、生成” 类任务,再通过自动化工具串联全流程,无需人工干预。其典型结构如下:

1.2 AI 工作流核心组件

  • 数据源:提供原始数据的载体,如 Excel/CSV、数据库(MySQL/PostgreSQL)、API 接口(电商平台 / CRM)、文档(PDF/Word)、消息工具(邮件 / Slack)
  • AI 引擎:处理核心逻辑的模块,分为两类:
      • 通用 AI:OpenAI GPT-4o、Anthropic Claude、百度文心一言(适合文本处理、生成)
      • 垂直 AI:阿里通义千问(电商场景)、科大讯飞星火(语音转文字)、本地模型(Llama 3/Mistral,适合隐私数据)
  • 自动化节点:串联流程的 “连接器”,负责触发条件设置、节点依赖管理、异常处理(如重试 / 告警)
  • 输出终端:呈现或存储结果的载体,如表格(Google Sheets/Excel)、报表工具(Tableau/Power BI)、消息工具(邮件 / Slack)、文档(Notion/Confluence)
  • 二、AI 工作流工具选型指南

    2.1 工具分类与适用场景

    工具类型

    代表工具

    技术门槛

    适用人群

    核心优势

    开发型工具

    Python+LangChain+Airflow

    中高

    程序员 / 数据分析师

    高度定制化,支持复杂逻辑

    低代码工具

    Make(原 Integromat)、Power Automate

    运营 / 产品 / 非技术人员

    可视化操作,无需代码

    AI 模型接口

    OpenAI API、百度智能云

    全人群(需 API 密钥)

    快速调用成熟 AI 能力

    数据存储工具

    MySQL、Google Sheets

    全人群

    轻量化 / 企业级数据存储

    2.2 选型核心原则

  • 需求匹配:简单流程(如 “邮件→AI 提取→表格存储”)用低代码工具;复杂逻辑(如 “多数据源关联 + 自定义 AI 模型”)用开发型工具
  • 技术门槛:非技术人员优先选 Make/Power Automate;有 Python 基础者可选 LangChain+Airflow
  • 成本控制:个人 / 小团队优先用免费额度(OpenAI 免费额度、Make 免费版);企业级需考虑 API 调用成本、工具订阅费
  • 隐私安全:处理敏感数据(如客户身份证号)需选支持本地部署的工具(如 Llama 3 本地模型 + 自建 Airflow)
  • 三、AI 工作流搭建核心步骤(附代码 / 操作)

    以 “客户反馈分析工作流” 为例(需求:从表格获取客户反馈→AI 分类→生成统计报告→推送 Slack),分步骤拆解搭建过程。

    3.1 步骤 1:需求拆解(关键前置动作)

    先将模糊需求拆分为可执行的节点,避免后续返工。以本案例为例,拆解结果如下:

  • 触发条件:每日 9 点自动读取 Google Sheets 中的新反馈数据
  • 数据处理:过滤空值反馈,提取 “客户 ID、反馈内容、提交时间” 字段
  • AI 处理:将反馈分类为 “产品质量”“服务态度”“物流问题” 三类
  • 结果存储:将分类结果写回 Google Sheets,更新 “分类状态” 列
  • 报告生成:统计各分类数量,生成 Markdown 格式日报
  • 通知推送:将日报发送至 Slack 客户服务频道
  • 3.2 步骤 2:数据接入(两种实现方式)

    方式 1:开发型(Python 代码接入)

    适用于需处理大量数据或自定义逻辑的场景,以下为接入 Excel/API/ 数据库的代码示例。

    (1)接入 Excel 文件

    # 安装依赖:pip install pandas openpyxl

    import pandas as pd

    from datetime import datetime

    def load_feedback_from_excel(file_path: str, sheet_name: str = "反馈数据"):

    """

    从Excel加载客户反馈数据,筛选当日新数据

    :param file_path: Excel文件路径

    :param sheet_name: 工作表名称

    :return: 清洗后的DataFrame

    """

    try:

    # 读取Excel文件

    df = pd.read_excel(file_path, sheet_name=sheet_name, engine="openpyxl")

    # 数据清洗:过滤空值、筛选当日数据

    df_clean = df.dropna(subset=["反馈内容", "客户ID"]) # 排除关键字段为空的行

    df_clean["提交时间"] = pd.to_datetime(df_clean["提交时间"], errors="coerce") # 转换时间格式

    today = datetime.now().date()

    df_today = df_clean[df_clean["提交时间"].dt.date == today] # 仅保留今日数据

    print(f"成功加载 {len(df_today)} 条今日反馈数据")

    return df_today

    except Exception as e:

    print(f"Excel数据加载失败:{str(e)}")

    return pd.DataFrame() # 返回空DataFrame避免后续报错

    # 调用函数(替换为你的文件路径)

    feedback_df = load_feedback_from_excel(file_path="C:/data/客户反馈.xlsx")

    (2)接入 API 接口(以电商平台反馈 API 为例)

    # 安装依赖:pip install requests

    import requests

    import json

    from datetime import datetime, timedelta

    def load_feedback_from_api(api_url: str, api_key: str):

    """

    从API接口获取近24小时客户反馈数据

    :param api_url: API接口地址(由平台提供)

    :param api_key: 接口授权密钥

    :return: 反馈数据列表

    """

    # 计算近24小时时间范围(API通常支持时间筛选)

    end_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")

    start_time = (datetime.now() – timedelta(hours=24)).strftime("%Y-%m-%d %H:%M:%S")

    # 构造请求头和参数

    headers = {

    "Authorization": f"Bearer {api_key}",

    "Content-Type": "application/json"

    }

    params = {

    "start_time": start_time,

    "end_time": end_time,

    "page_size": 100 # 单次获取最大条数

    }

    try:

    # 发送GET请求

    response = requests.get(api_url, headers=headers, params=params, timeout=10)

    response.raise_for_status() # 若HTTP状态码非200,抛出异常

    data = response.json()

    # 提取核心反馈数据(根据API返回格式调整,此处为示例)

    feedback_list = data.get("data", {}).get("feedback", [])

    print(f"从API获取 {len(feedback_list)} 条反馈数据")

    return feedback_list

    except requests.exceptions.RequestException as e:

    print(f"API请求失败:{str(e)}")

    return []

    # 调用函数(替换为你的API地址和密钥)

    api_feedback = load_feedback_from_api(

    api_url="https://api.xxx.com/v1/customer/feedback",

    api_key="sk-xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"

    )

    (3)接入 MySQL 数据库

    # 安装依赖:pip install pymysql pandas

    import pymysql

    import pandas as pd

    from datetime import datetime

    def load_feedback_from_mysql(host: str, user: str, password: str, db_name: str, table_name: str):

    """

    从MySQL数据库获取今日客户反馈数据

    :param host: 数据库地址(如localhost)

    :param user: 数据库用户名

    :param password: 数据库密码

    :param db_name: 数据库名称

    :param table_name: 表名称

    :return: 反馈数据DataFrame

    """

    # 建立数据库连接

    conn = None

    try:

    conn = pymysql.connect(

    host=host,

    user=user,

    password=password,

    database=db_name,

    charset="utf8mb4"

    )

    # 编写SQL:筛选今日数据,仅获取需要的字段

    today = datetime.now().date().strftime("%Y-%m-%d")

    sql = f"""

    SELECT customer_id, feedback_content, submit_time

    FROM {table_name}

    WHERE DATE(submit_time) = '{today}'

    AND feedback_content IS NOT NULL

    """

    # 执行SQL并转换为DataFrame

    df = pd.read_sql(sql, conn)

    print(f"从MySQL获取 {len(df)} 条今日反馈数据")

    return df

    except Exception as e:

    print(f"MySQL数据加载失败:{str(e)}")

    return pd.DataFrame()

    finally:

    if conn:

    conn.close() # 关闭连接

    # 调用函数(替换为你的数据库信息)

    mysql_feedback = load_feedback_from_mysql(

    host="localhost",

    user="root",

    password="123456",

    db_name="customer_service",

    table_name="feedback"

    )

    方式 2:低代码(Make 连接 Google Sheets)

    适用于非技术人员,无需代码,通过可视化配置接入数据:

  • 登录 Make(www.make.com),创建新场景(Scenarios)
  • 点击 “+” 添加触发器,搜索 “Google Sheets” 并选择 “Watch rows”(监控新行)
  • 授权 Google 账号:点击 “Add”,登录你的 Google 账号,允许 Make 访问 Google Sheets
  • 配置触发器参数(图 1):
      • Spreadsheet ID:从 Google Sheets URL 获取(格式:https://docs.google.com/spreadsheets/d/[SPREADSHEET_ID]/edit)
      • Sheet name:输入工作表名称(如 “反馈数据”)
      • Trigger column:选择 “提交时间”(确保新行有时间戳,触发流程)
  • 点击 “Run once” 测试:Make 会自动拉取表格中的测试行,确认数据能正常获取(图 2)
  • 图 1:Make Google Sheets 触发器配置界面

    界面左侧为触发器列表,中间为流程画布,右侧为配置面板。红框标注关键参数:“Spreadsheet ID” 输入框、“Sheet name” 下拉框、“Trigger column” 选择框,底部有 “Run once” 测试按钮。

    图 2:Make 触发器测试结果界面

    测试成功后,界面会显示拉取到的字段(如客户 ID、反馈内容、提交时间)及对应值,底部提示 “Successfully loaded 1 row”。

    3.3 步骤 3:AI 模型集成(分类 / 分析 / 生成)

    方式 1:开发型(Python 调用 AI 接口 + LangChain)
    (1)直接调用 OpenAI API 分类反馈

    # 安装依赖:pip install openai python-dotenv

    from openai import OpenAI

    from dotenv import load_dotenv # 用于加载环境变量,避免硬编码API密钥

    import os

    # 加载环境变量(建议将API密钥存放在.env文件中,格式:OPENAI_API_KEY=sk-xxx)

    load_dotenv()

    client = OpenAI(api_key=os.getenv("OPENAI_API_KEY"))

    def classify_feedback(feedback_content: str) -> str:

    """

    用GPT-4o-mini对客户反馈进行分类

    :param feedback_content: 反馈内容

    :return: 分类结果(产品质量/服务态度/物流问题)

    """

    # 提示词设计:明确任务、限制输出格式,避免AI返回多余内容

    prompt = f"""

    任务:将客户反馈按问题类型分类,仅返回以下三类中的一种,无需额外解释:

    1. 产品质量(如产品损坏、功能故障)

    2. 服务态度(如客服不耐烦、响应慢)

    3. 物流问题(如延迟发货、包裹丢失)

    客户反馈:{feedback_content}

    分类结果:

    """

    try:

    # 调用OpenAI API

    response = client.chat.completions.create(

    model="gpt-4o-mini", # 性价比高,适合分类任务

    messages=[{"role": "user", "content": prompt}],

    temperature=0.2, # 降低随机性,确保分类稳定

    max_tokens=10 # 限制输出长度,避免冗余

    )

    # 提取分类结果并清理空格

    category = response.choices[0].message.content.strip()

    # 兜底处理:若AI返回非预期结果,标记为“未分类”

    valid_categories = ["产品质量", "服务态度", "物流问题"]

    return category if category in valid_categories else "未分类"

    except Exception as e:

    print(f"AI分类失败:{str(e)}")

    return "未分类"

    # 测试函数

    test_feedback = "收到的手机屏幕有裂痕,明显是质量问题"

    print(classify_feedback(test_feedback)) # 输出:产品质量

    (2)用 LangChain 构建复杂 AI 链(分类 + 总结)

    LangChain 是连接 AI 模型与数据源的框架,适合构建多步骤 AI 逻辑(如 “分类 + 总结 + 关键词提取”):

    # 安装依赖:pip install langchain langchain-openai

    from langchain.prompts import PromptTemplate

    from langchain_openai import ChatOpenAI

    from langchain.schema.output_parser import StrOutputParser

    from langchain.chains import LLMChain

    from dotenv import load_dotenv

    import os

    load_dotenv()

    def build_feedback_chain():

    """构建反馈处理链:分类+总结"""

    # 1. 定义提示模板(支持动态传入变量)

    feedback_template = PromptTemplate(

    input_variables=["feedback_content"],

    template="""

    任务:处理客户反馈,按以下格式输出(无额外文字):

    分类:{分类结果,仅产品质量/服务态度/物流问题}

    总结:{一句话总结反馈核心内容,不超过50字}

    客户反馈:{feedback_content}

    """

    )

    # 2. 初始化AI模型

    llm = ChatOpenAI(

    model="gpt-4o-mini",

    api_key=os.getenv("OPENAI_API_KEY"),

    temperature=0.3

    )

    # 3. 构建链(模板→LLM→输出解析)

    chain = feedback_template | llm | StrOutputParser()

    return chain

    # 初始化链

    feedback_chain = build_feedback_chain()

    # 处理单条反馈

    def process_feedback(feedback_content: str) -> dict:

    """

    处理反馈:分类+总结

    :return: 包含分类和总结的字典

    """

    result = feedback_chain.invoke({"feedback_content": feedback_content})

    # 解析输出结果(按模板格式拆分)

    lines = [line.strip() for line in result.split("\\n") if line.strip()]

    category = lines[0].replace("分类:", "") if len(lines) > 0 else "未分类"

    summary = lines[1].replace("总结:", "") if len(lines) > 1 else "无总结"

    return {"分类": category, "总结": summary}

    # 测试批量处理(从Excel加载的数据)

    feedback_df = load_feedback_from_excel("C:/data/客户反馈.xlsx")

    if not feedback_df.empty:

    # 应用处理函数到每一行

    feedback_df[["分类结果", "反馈总结"]] = feedback_df["反馈内容"].apply(

    lambda x: pd.Series(process_feedback(x))

    )

    print("批量处理完成,前3条结果:")

    print(feedback_df[["客户ID", "反馈内容", "分类结果", "反馈总结"]].head(3))

    方式 2:低代码(Make 集成 OpenAI)
  • 在 Make 场景画布中,点击 “+” 添加新节点,搜索 “OpenAI” 并选择 “Create a chat completion”
  • 授权 OpenAI 账号:点击 “Add”,输入 API 密钥(从 OpenAI 控制台获取:platform.openai.com/api-keys)
  • 配置节点参数(图 3):
      • Model:选择 “gpt-4o-mini”
      • Messages:点击 “Add item”,Role 选择 “user”,Content 输入提示词:

    请将以下客户反馈分类为"产品质量"、"服务态度"或"物流问题",仅返回分类结果:{{1.Google Sheets.Watch rows.反馈内容}}

    (注:{{1.Google Sheets.Watch rows.反馈内容}}是引用前序触发器的 “反馈内容” 字段)

  • 点击 “Test” 测试:输入测试反馈内容,确认 AI 返回正确分类(图 4)
  • 图 3:Make OpenAI 节点配置界面

    右侧配置面板分 “Connection”(账号授权)、“Settings”(模型选择)、“Messages”(提示词配置)三部分,红框标注 “Content” 中的变量引用。

    图 4:Make OpenAI 测试结果界面

    测试成功后,界面显示 “Choices [0].Message.Content” 字段,值为 “产品质量”(或其他分类),底部提示 “Success”。

    3.4 步骤 4:自动化逻辑配置(触发 + 执行)

    方式 1:开发型(Airflow 调度定时任务)

    Airflow 是开源调度工具,适合定时执行 AI 工作流(如每日 9 点运行反馈分析)。

    (1)Airflow 环境搭建(Windows 为例)
  • 安装 WSL2(Windows 子系统):参考微软官方教程
  • 在 WSL2 中安装 Python3.8+:sudo apt update && sudo apt install python3 python3-pip
  • 安装 Airflow:
  • # 设置Airflow版本(建议2.8.x)

    AIRFLOW_VERSION=2.8.4

    PYTHON_VERSION="$(python3 –version | cut -d " " -f 2 | cut -d "." -f 1-2)"

    CONSTRAINT_URL="https://raw.githubusercontent.com/apache/airflow/constraints-${AIRFLOW_VERSION}/constraints-${PYTHON_VERSION}.txt"

    pip install "apache-airflow==${AIRFLOW_VERSION}" –constraint "${CONSTRAINT_URL}"

  • 初始化 Airflow:
  • airflow db init

    airflow users create \\

    –username admin \\

    –firstname Admin \\

    –lastname User \\

    –role Admin \\

    –email admin@example.com

  • 启动 Airflow:
  • airflow webserver –port 8080 # 前台运行Web界面

    # 新打开一个WSL窗口,启动调度器

    airflow scheduler

  • 访问 Airflow 界面:浏览器输入http://localhost:8080,用创建的账号登录。
  • (2)编写 Airflow DAG(定时执行工作流)

    DAG(有向无环图)定义任务依赖关系,以下为 “每日 9 点执行客户反馈分析” 的 DAG 代码:

    # 保存路径:$AIRFLOW_HOME/dags/feedback_analysis_dag.py

    from airflow import DAG

    from airflow.operators.python import PythonOperator

    from datetime import datetime, timedelta

    import pandas as pd

    from openai import OpenAI

    from dotenv import load_dotenv

    import os

    import smtplib

    from email.mime.text import MIMEText

    from email.mime.multipart import MIMEMultipart

    # 加载环境变量

    load_dotenv()

    client = OpenAI(api_key=os.getenv("OPENAI_API_KEY"))

    # ———————- 1. 定义任务函数 ———————-

    def load_data(**kwargs):

    """任务1:加载今日反馈数据(Excel)"""

    df = pd.read_excel("C:/data/客户反馈.xlsx", sheet_name="反馈数据", engine="openpyxl")

    df_clean = df.dropna(subset=["反馈内容", "客户ID"])

    df_clean["提交时间"] = pd.to_datetime(df_clean["提交时间"], errors="coerce")

    today = datetime.now().date()

    df_today = df_clean[df_clean["提交时间"].dt.date == today]

    # 用XCom传递数据(Airflow中任务间共享数据的方式)

    kwargs["ti"].xcom_push(key="feedback_data", value=df_today.to_dict("records"))

    print(f"任务1:加载 {len(df_today)} 条今日数据")

    def ai_classify(**kwargs):

    """任务2:AI分类反馈"""

    # 从XCom获取前序任务数据

    feedback_data = kwargs["ti"].xcom_pull(key="feedback_data", task_ids="load_data")

    processed_data = []

    for feedback in feedback_data:

    content = feedback["反馈内容"]

    # AI分类(复用之前的函数)

    prompt = f"分类客户反馈为产品质量/服务态度/物流问题,仅返回结果:{content}"

    try:

    response = client.chat.completions.create(

    model="gpt-4o-mini",

    messages=[{"role": "user", "content": prompt}],

    temperature=0.2

    )

    category = response.choices[0].message.content.strip()

    except Exception as e:

    category = "未分类"

    print(f"AI分类失败:{str(e)}")

    processed_data.append({

    "客户ID": feedback["客户ID"],

    "反馈内容": content,

    "分类结果": category,

    "处理时间": datetime.now().strftime("%Y-%m-%d %H:%M:%S")

    })

    # 传递处理后的数据

    kwargs["ti"].xcom_push(key="processed_data", value=processed_data)

    print(f"任务2:完成 {len(processed_data)} 条数据分类")

    def generate_report(**kwargs):

    """任务3:生成日报"""

    processed_data = kwargs["ti"].xcom_pull(key="processed_data", task_ids="ai_classify")

    df = pd.DataFrame(processed_data)

    # 统计分类数量

    category_count = df["分类结果"].value_counts().to_dict()

    total = len(df)

    # 生成Markdown报告

    report = f"""

    # 客户反馈日报({datetime.now().strftime("%Y-%m-%d")})

    ## 1. 数据概览

    – 今日总反馈数:{total} 条

    – 处理完成率:100%

    – 报告生成时间:{datetime.now().strftime("%Y-%m-%d %H:%M:%S")}

    ## 2. 分类统计

    """

    for cat, count in category_count.items():

    percentage = (count / total) * 100 if total > 0 else 0

    report += f"- {cat}:{count} 条({percentage:.1f}%)\\n"

    # 保存报告到文件

    report_path = f"C:/data/反馈日报_{datetime.now().strftime('%Y%m%d')}.md"

    with open(report_path, "w", encoding="utf-8") as f:

    f.write(report)

    # 传递报告内容和路径

    kwargs["ti"].xcom_push(key="report_content", value=report)

    kwargs["ti"].xcom_push(key="report_path", value=report_path)

    print(f"任务3:日报生成完成,路径:{report_path}")

    def send_email(**kwargs):

    """任务4:发送邮件通知"""

    report_content = kwargs["ti"].xcom_pull(key="report_content", task_ids="generate_report")

    report_path = kwargs["ti"].xcom_pull(key="report_path", task_ids="generate_report")

    # 邮件配置(替换为你的邮箱信息)

    sender = "your-email@example.com"

    receiver = "service-team@example.com"

    password = os.getenv("EMAIL_PASSWORD") # 建议从.env文件加载

    subject = f"客户反馈日报 – {datetime.now().strftime('%Y-%m-%d')}"

    # 构建邮件

    msg = MIMEMultipart()

    msg["From"] = sender

    msg["To"] = receiver

    msg["Subject"] = subject

    # 添加邮件正文(Markdown格式)

    msg.attach(MIMEText(report_content, "markdown", "utf-8"))

    # 发送邮件(以QQ邮箱为例,SMTP服务器:smtp.qq.com,端口465)

    try:

    with smtplib.SMTP_SSL("smtp.qq.com", 465) as server:

    server.login(sender, password)

    server.send_message(msg)

    print("任务4:邮件发送成功")

    except Exception as e:

    print(f"任务4:邮件发送失败:{str(e)}")

    raise e # 抛出异常,Airflow会标记任务失败并告警

    # ———————- 2. 定义DAG ———————-

    default_args = {

    "owner": "customer_service_team", # 任务所属团队

    "depends_on_past": False, # 不依赖历史执行结果

    "start_date": datetime(2025, 11, 18), # DAG开始日期

    "email_on_failure": True, # 任务失败时发送邮件

    "email": ["alert@example.com"], # 告警邮箱

    "retries": 1, # 失败后重试1次

    "retry_delay": timedelta(minutes=5), # 重试间隔5分钟

    }

    with DAG(

    dag_id="customer_feedback_daily_analysis", # DAG唯一ID

    default_args=default_args,

    description="每日客户反馈分析工作流",

    schedule_interval=timedelta(days=1), # 执行频率:每天1次

    # 定时执行时间:每日9点(用cron表达式:0 9 * * *)

    # schedule_interval="0 9 * * *",

    catchup=False, # 不补跑历史任务

    tags=["客户服务", "AI分类"] # 标签,便于筛选

    ) as dag:

    # 定义任务

    task1 = PythonOperator(

    task_id="load_data", # 任务唯一ID

    python_callable=load_data, # 关联的函数

    provide_context=True, # 允许函数接收kwargs(含ti等Airflow对象)

    )

    task2 = PythonOperator(

    task_id="ai_classify",

    python_callable=ai_classify,

    provide_context=True,

    )

    task3 = PythonOperator(

    task_id="generate_report",

    python_callable=generate_report,

    provide_context=True,

    )

    task4 = PythonOperator(

    task_id="send_email",

    python_callable=send_email,

    provide_context=True,

    )

    # 定义任务依赖:task1 → task2 → task3 → task4

    task1 >> task2 >> task3 >> task4

    (3)启动 DAG 并监控
  • 登录 Airflow Web 界面(http://localhost:8080),在 “DAGs” 页面找到 “customer_feedback_daily_analysis”
  • 点击开关启动 DAG(图 5),若需立即测试,点击 “Trigger DAG” 手动触发
  • 查看执行状态:点击 DAG 名称进入详情页,“Graph” 视图可查看各任务执行状态(绿色为成功,红色为失败)
  • 查看日志:若任务失败,点击任务节点→“Logs” 查看错误信息(如 API 密钥错误、文件路径不存在)
  • 图 5:Airflow DAG 启动界面

    DAG 列表中,“customer_feedback_daily_analysis” 行的 “On/Off” 开关切换为 “On”,右侧 “Actions” 列有 “Trigger DAG”(手动触发)、“Graph”(查看任务图)按钮。

    方式 2:低代码(Make 配置自动化逻辑)
  • 完成 “Google Sheets 触发器→OpenAI 分类” 节点后,点击 “+” 添加 “Google Sheets” 的 “Update a row” 节点(更新分类结果)
  • 配置更新节点(图 6):
      • Spreadsheet ID:与触发器一致
      • Sheet name:“反馈数据”
      • Row ID:选择 “1.Google Sheets.Watch rows.Row number”(更新触发器获取的行)
      • 找到 “分类结果” 列,输入值:{{2.OpenAI.Create a chat completion.Choices[0].Message.Content}}(引用 OpenAI 的分类结果)
  • 继续添加 “Slack” 节点(推送日报):
      • 搜索 “Slack”→选择 “Send channel message”
      • 授权 Slack 账号,选择目标频道(如 #customer-service)
      • 配置消息内容:

    【客户反馈通知】

    新收到1条客户反馈:

    • 客户ID:{{1.Google Sheets.Watch rows.客户ID}}

    • 反馈内容:{{1.Google Sheets.Watch rows.反馈内容}}

    • 分类结果:{{2.OpenAI.Create a chat completion.Choices[0].Message.Content}}

    • 提交时间:{{1.Google Sheets.Watch rows.提交时间}}

  • 设置触发方式(图 7):
      • 点击场景编辑器右上角的 “Schedule”→选择 “Immediately”(有新行时实时触发)
      • 若需定时触发(如每日 9 点),选择 “Custom”→设置 cron 表达式(如0 9 * * *)
  • 保存并激活场景:点击 “Save” 输入场景名称(如 “客户反馈自动分类”),再点击 “Activate” 启动流程
  • 图 6:Make Google Sheets 更新节点配置界面

    右侧面板 “Row ID” 选择 “Row number”,“分类结果” 列的值引用 OpenAI 节点的输出,底部有 “Test” 按钮测试更新效果。

    图 7:Make 场景触发方式设置界面

    界面顶部有 “Immediately”“Hourly”“Daily” 等快捷选项,底部 “Custom” 可设置自定义 cron 表达式,红框标注 “Activate” 激活按钮。

    3.5 步骤 5:结果输出与存储

    方式 1:开发型(Excel / 数据库 / 文件)
    (1)保存结果到 Excel

    # 承接步骤3.3.2的批量处理结果

    def save_result_to_excel(df: pd.DataFrame, output_path: str):

    """保存处理结果到Excel"""

    try:

    with pd.ExcelWriter(output_path, engine="openpyxl") as writer:

    # 保存详细结果

    df.to_excel(writer, sheet_name="详细结果", index=False)

    # 保存分类统计

    category_stats = df["分类结果"].value_counts().reset_index()

    category_stats.columns = ["分类", "数量"]

    category_stats.to_excel(writer, sheet_name="分类统计", index=False)

    print(f"结果已保存到:{output_path}")

    except Exception as e:

    print(f"保存Excel失败:{str(e)}")

    # 调用函数

    output_path = f"C:/data/客户反馈处理结果_{datetime.now().strftime('%Y%m%d')}.xlsx"

    save_result_to_excel(feedback_df, output_path)

    (2)保存结果到 MySQL 数据库

    def save_result_to_mysql(df: pd.DataFrame, host: str, user: str, password: str, db_name: str, table_name: str):

    """保存处理结果到MySQL"""

    conn = None

    try:

    conn = pymysql.connect(

    host=host,

    user=user,

    password=password,

    database=db_name,

    charset="utf8mb4"

    )

    cursor = conn.cursor()

    # 先清空今日数据(避免重复插入)

    today = datetime.now().date().strftime("%Y-%m-%d")

    delete_sql = f"DELETE FROM {table_name} WHERE DATE(处理时间) = '{today}'"

    cursor.execute(delete_sql)

    # 插入新数据

    insert_sql = f"""

    INSERT INTO {table_name} (客户ID, 反馈内容, 分类结果, 反馈总结, 处理时间)

    VALUES (%s, %s, %s, %s, %s)

    """

    # 转换DataFrame为元组列表

    data_tuples = [

    (row["客户ID"], row["反馈内容"], row["分类结果"], row["反馈总结"], row["处理时间"])

    for _, row in df.iterrows()

    ]

    cursor.executemany(insert_sql, data_tuples)

    conn.commit()

    print(f"成功插入 {len(data_tuples)} 条数据到MySQL")

    except Exception as e:

    if conn:

    conn.rollback() # 失败回滚

    print(f"保存MySQL失败:{str(e)}")

    finally:

    if conn:

    cursor.close()

    conn.close()

    # 调用函数

    save_result_to_mysql(

    df=feedback_df,

    host="localhost",

    user="root",

    password="123456",

    db_name="customer_service",

    table_name="feedback_processed"

    )

    方式 2:低代码(Google Sheets/Slack)

    已在步骤 3.4.2 中完成配置,结果会自动:

  • 写回 Google Sheets 的 “分类结果” 列
  • 推送通知到 Slack 频道,包含客户 ID、反馈内容、分类结果等信息
  • 四、实战案例精讲

    案例 1:低代码版 —— 客户反馈自动分析工作流(Make+OpenAI+Google Sheets+Slack)

    1.1 需求定义
    • 触发条件:当 Google Sheets “反馈数据” 表新增一行时
    • 核心逻辑:AI 自动分类反馈→更新表格分类结果→推送 Slack 通知
    • 输出结果:表格更新、Slack 实时通知
    1.2 工具准备
    • Make 账号(免费版支持 1000 次 / 月执行)
    • OpenAI 账号(获取 API 密钥,免费额度足够测试)
    • Google 账号(创建 Google Sheets 表格)
    • Slack 账号(创建客户服务频道)
    1.3 分步搭建(完整操作)
  • 创建 Google Sheets 表格
      • 登录 Google Sheets,创建新表格,命名为 “客户反馈管理”
      • 新建工作表 “反馈数据”,设置列名:客户 ID、反馈内容、提交时间、分类结果(图 8)
      • 手动添加 1 条测试数据:客户 ID=1001,反馈内容 =“客服回复太慢,等了半小时”,提交时间 =“2025-11-18 10:00:00”,分类结果 = 空

    图 8:Google Sheets 反馈数据表格界面

    表格有 4 列,列名分别为 “客户 ID”“反馈内容”“提交时间”“分类结果”,第一行已填入测试数据,“分类结果” 列为空。

  • 创建 Make 场景并添加触发器
      • 登录 Make→新建场景→添加 “Google Sheets” 触发器→选择 “Watch rows”
      • 授权 Google 账号→选择 “客户反馈管理” 表格→工作表 “反馈数据”→触发列 “提交时间”
      • 点击 “Run once”→选择测试行→确认获取到数据(客户 ID=1001,反馈内容 =“客服回复太慢…”)
  • 添加 OpenAI 分类节点
      • 点击 “+”→添加 “OpenAI”→“Create a chat completion”
      • 授权 API 密钥→模型选择 “gpt-4o-mini”
      • Messages 配置:Role=user,Content=“分类客户反馈为产品质量 / 服务态度 / 物流问题,仅返回结果:{{1.Google Sheets.Watch rows. 反馈内容}}”
      • 测试节点→确认返回 “服务态度”
  • 添加 Google Sheets 更新节点
      • 点击 “+”→添加 “Google Sheets”→“Update a row”
      • 选择同一表格→工作表 “反馈数据”→Row ID=“1.Google Sheets.Watch rows.Row number”
      • “分类结果” 列输入 “{{2.OpenAI.Create a chat completion.Choices [0].Message.Content}}”
      • 测试节点→查看 Google Sheets,测试行 “分类结果” 列变为 “服务态度”(图 9)

    图 9:Google Sheets 更新结果界面

    测试行的 “分类结果” 列已从空变为 “服务态度”,与 OpenAI 分类结果一致。

  • 添加 Slack 通知节点
      • 点击 “+”→添加 “Slack”→“Send channel message”
      • 授权 Slack 账号→选择频道 #customer-service
      • 消息内容:

    📢 新客户反馈通知

    客户ID:{{1.Google Sheets.Watch rows.客户ID}}

    反馈内容:{{1.Google Sheets.Watch rows.反馈内容}}

    分类结果:{{2.OpenAI.Create a chat completion.Choices[0].Message.Content}}

    提交时间:{{1.Google Sheets.Watch rows.提交时间}}

      • 测试节点→查看 Slack 频道,收到通知(图 10)

    图 10:Slack 通知界面

    频道 #customer-service 中显示新消息,包含客户 ID、反馈内容等信息,分类结果为 “服务态度”,格式清晰。

  • 激活场景
      • 点击 “Save”→命名场景 “客户反馈自动分类”
      • 开启 “Schedule”→选择 “Immediately”
      • 点击 “Activate”→场景变为 “Active” 状态(图 11)

    图 11:Make 场景激活界面

    场景名称右侧显示 “Active” 绿色标签,右上角 “Schedule” 显示 “Immediately”,底部有 “Execution history” 查看执行记录。

  • 测试全流程
      • 在 Google Sheets 添加新行:客户 ID=1002,反馈内容 =“快递延迟 3 天还没到”,提交时间 =“2025-11-18 11:00:00”
      • 等待 1-2 分钟,查看 Make “Execution history”→显示 “Success”(图 12)
      • 验证结果:Google Sheets “分类结果” 列变为 “物流问题”,Slack 收到新通知

    图 12:Make 执行历史界面

    显示最新执行记录,状态为 “Success”,耗时 2.3 秒,可点击查看各节点详细输出。

    案例 2:开发版 —— 数据分析师日报自动化工作流(Python+LangChain+Airflow+MySQL+Email)

    2.1 需求定义
    • 触发条件:每日 9 点自动执行
    • 核心逻辑:从 MySQL 拉取前一日销售数据→AI 生成分析结论→生成 Excel 日报→发送邮件给团队
    • 输出结果:Excel 日报文件、邮件通知
    2.2 环境准备

    # 保存路径:$AIRFLOW_HOME/dags/sales_daily_report_dag.py

    from airflow import DAG

    from airflow.operators.python import PythonOperator

    from datetime import datetime, timedelta

    import pandas as pd

    import pymysql

    from langchain.prompts import PromptTemplate

    from langchain_openai import ChatOpenAI

    from langchain.schema.output_parser import StrOutputParser

    from dotenv import load_dotenv

    import os

    import smtplib

    from email.mime.text import MIMEText

    from email.mime.multipart import MIMEMultipart

    from email.mime.base import MIMEBase

    from email import encoders

    # 加载环境变量

    load_dotenv()

    OPENAI_API_KEY = os.getenv("OPENAI_API_KEY")

    EMAIL_SENDER = os.getenv("EMAIL_SENDER")

    EMAIL_PASSWORD = os.getenv("EMAIL_PASSWORD")

    EMAIL_RECEIVERS = os.getenv("EMAIL_RECEIVERS").split(",") # 多个收件人用逗号分隔

    # ———————- 1. 工具函数 ———————-

    def get_prev_day() -> str:

    """获取前一天日期(格式:YYYY-MM-DD)"""

    return (datetime.now() – timedelta(days=1)).strftime("%Y-%m-%d")

    def connect_mysql() -> pymysql.connections.Connection:

    """建立MySQL连接"""

    return pymysql.connect(

    host="localhost",

    user="root",

    password="123456",

    database="sales_analysis",

    charset="utf8mb4"

    )

    # ———————- 2. 任务函数 ———————-

    def extract_sales_data(**kwargs):

    """任务1:提取前一日销售数据"""

    prev_day = get_prev_day()

    conn = None

    try:

    conn = connect_mysql()

    # 提取前一日销售数据,按产品分组

    sql = f"""

    SELECT

    product_id,

    SUM(sales_amount) AS total_sales,

    SUM(sales_volume) AS total_volume,

    AVG(sales_amount / sales_volume) AS avg_price # 客单价

    FROM sales_data

    WHERE date = '{prev_day}'

    GROUP BY product_id

    ORDER BY total_sales DESC

    """

    df = pd.read_sql(sql, conn)

    print(f"任务1:提取 {prev_day} 销售数据,共 {len(df)} 个产品")

    # 传递数据

    kwargs["ti"].xcom_push(key="sales_data", value=df.to_dict("records"))

    kwargs["ti"].xcom_push(key="prev_day", value=prev_day)

    except Exception as e:

    print(f"任务1失败:{str(e)}")

    raise e

    finally:

    if conn:

    conn.close()

    def ai_analyze_sales(**kwargs):

    """任务2:AI分析销售数据"""

    sales_data = kwargs["ti"].xcom_pull(key="sales_data", task_ids="extract_data")

    prev_day = kwargs["ti"].xcom_pull(key="prev_day", task_ids="extract_data")

    df = pd.DataFrame(sales_data)

    # 计算整体指标

    total_sales = df["total_sales"].sum()

    total_volume = df["total_volume"].sum()

    top_product = df.iloc[0]["product_id"] if len(df) > 0 else "无"

    top_sales = df.iloc[0]["total_sales"] if len(df) > 0 else 0

    # 构建AI分析提示词

    prompt_template = PromptTemplate(

    input_variables=["prev_day", "total_sales", "total_volume", "top_product", "top_sales", "sales_detail"],

    template="""

    任务:分析 {prev_day} 销售数据,生成专业数据分析师报告,包含以下部分:

    1. 核心指标:总销售额、总销量、平均客单价(总销售额/总销量)

    2. 亮点产品:销售额最高的产品及占比(该产品销售额/总销售额)

    3. 潜在问题:若有产品销售额为0或客单价异常(低于均值50%),需指出

    4. 建议:基于数据提出1-2条运营建议(如重点推广亮点产品)

    数据:

    – 总销售额:{total_sales} 元

    – 总销量:{total_volume} 件

    – 亮点产品:{top_product}(销售额 {top_sales} 元)

    – 产品详情:{sales_detail}

    要求:语言简洁专业,避免冗余,不超过500字。

    """

    )

    # 初始化AI链

    llm = ChatOpenAI(model="gpt-4o-mini", api_key=OPENAI_API_KEY, temperature=0.4)

    chain = prompt_template | llm | StrOutputParser()

    # 执行分析

    try:

    analysis = chain.invoke({

    "prev_day": prev_day,

    "total_sales": total_sales,

    "total_volume": total_volume,

    "top_product": top_product,

    "top_sales": top_sales,

    "sales_detail": df.to_string(index=False) # 产品详情表格

    })

    print(f"任务2:AI分析完成,分析内容:\\n{analysis}")

    # 传递分析结果和整体指标

    kwargs["ti"].xcom_push(key="ai_analysis", value=analysis)

    kwargs["ti"].xcom_push(key="total_sales", value=total_sales)

    kwargs["ti"].xcom_push(key="total_volume", value=total_volume)

    except Exception as e:

    print(f"任务2失败:{str(e)}")

    raise e

    def generate_excel_report(**kwargs):

    """任务3:生成Excel日报"""

    sales_data = kwargs["ti"].xcom_pull(key="sales_data", task_ids="extract_data")

    ai_analysis = kwargs["ti"].xcom_pull(key="ai_analysis", task_ids="ai_analyze")

    prev_day = kwargs["ti"].xcom_pull(key="prev_day", task_ids="extract_data")

    total_sales = kwargs["ti"].xcom_pull(key="total_sales", task_ids="ai_analyze")

    total_volume = kwargs["ti"].xcom_pull(key="total_volume", task_ids="ai_analyze")

    df = pd.DataFrame(sales_data)

    avg_price = total_sales / total_volume if total_volume > 0 else 0

    # 生成Excel

    report_path = f"/data/sales_report_{prev_day}.xlsx" # Airflow服务器路径

    try:

    with pd.ExcelWriter(report_path, engine="openpyxl") as writer:

    # 工作表1:核心指标

    metrics_df = pd.DataFrame({

    "指标名称": ["日期", "总销售额(元)", "总销量(件)", "平均客单价(元)"],

    "指标值": [prev_day, total_sales, total_volume, round(avg_price, 2)]

    })

    metrics_df.to_excel(writer, sheet_name="核心指标", index=False)

    # 工作表2:产品详情

    df.to_excel(writer, sheet_name="产品销售详情", index=False)

    # 工作表3:AI分析

    analysis_df = pd.DataFrame({"AI分析报告": [ai_analysis]})

    analysis_df.to_excel(writer, sheet_name="AI分析", index=False)

    print(f"任务3:Excel日报生成完成,路径:{report_path}")

    kwargs["ti"].xcom_push(key="report_path", value=report_path)

    except Exception as e:

    print(f"任务3失败:{str(e)}")

    raise e

    def send_email_with_attachment(**kwargs):

    """任务4:发送带Excel附件的邮件"""

    ai_analysis = kwargs["ti"].xcom_pull(key="ai_analysis", task_ids="ai_analyze")

    prev_day = kwargs["ti"].xcom_pull(key="prev_day", task_ids="extract_data")

    report_path = kwargs["ti"].xcom_pull(key="report_path", task_ids="generate</doubaocanvas>

    • Python 3.9+(安装依赖:pip install pandas pymysql langchain langchain-openai openai python-dotenv apache-airflow openpyxl)
    • MySQL 数据库(含销售数据表sales_data,字段:date、product_id、sales_amount、sales_volume)
    • Airflow(参考步骤 3.4.1 搭建)
    • 邮箱账号(如 QQ 邮箱,开启 SMTP 服务获取授权码)
    2.3 代码实现(完整 DAG)
    赞(0)
    未经允许不得转载:171主机测评 » AI 工具搭建工作流全流程教程(含代码 / 操作指引)
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址