构建AI Agent自主循环系统:从定时任务到智能自动化
你是不是也遇到过这样的场景:想用 AI Agent 自动化处理一些任务,比如写周报、分析数据、监控竞品,但每次都得自己手动写一个详细的 prompt 去“启动”它?更头疼的是,很多任务不是一次性的,而是需要持续、循环执行的。手动触发不仅效率低下,而且完全违背了“自动化”的初衷。
最近三个月,我们团队内部跑通了一套名为 “自主循环(Loop)系统” 的解决方案。它让 Agent 不再是“你推一下,它动一下”的提线木偶,而是变成了一个 “会自己找活干、自己安排日程”的智能员工 。这套系统稳定运行了三个月,处理了数千个周期性任务,从每日数据同步到每周竞品报告生成,完全无需人工干预。
这篇文章,我就来手把手拆解这套系统的核心设计、技术选型与落地步骤。 核心判断是:构建 Agent 的 Loop 系统,关键不在于模型本身有多强,而在于如何设计一个稳定、可观测、能自愈的任务调度与执行框架。 如果你正在为如何让 AI 真正“7x24小时”为自己工作而烦恼,那么这篇文章就是为你准备的。
1. 为什么你需要一个“会自己找活”的 Agent 系统?
在深入技术细节之前,我们先明确问题。传统的单次 Prompt 调用 AI,或者 even 一些简单的 Agent 框架,通常存在几个痛点:
- 被动触发 :每次都需要人工或外部系统事件来“唤醒”Agent。这相当于你雇了一个员工,但他从不主动看待办事项,必须你每天去拍他肩膀。
- 状态丢失 :每次执行都是全新的会话,Agent 不记得上次执行的结果、上下文和进度。对于需要持续跟踪的任务(如监控某个指标的变化),这几乎是致命的。
- 缺乏调度 :无法处理“每隔X小时/每天/每周执行一次”这类需求。你需要额外写一个 cron job 来调用 Agent,这又把调度逻辑和业务逻辑耦合在了一起。
- 异常处理薄弱 :一次执行失败就彻底停止,没有重试、没有告警、没有降级方案。
而一个理想的“自主循环”系统,应该具备以下能力:
- 任务编排与调度 :能定义周期性或条件触发的任务。
- 状态持久化 :能记住每次执行的历史、上下文和中间结果。
- 工作流引擎 :能处理复杂的、多步骤的任务流程,并支持分支判断。
- 可观测性 :能清晰地看到每个任务的状态、日志、消耗和结果。
- 自愈能力 :在遇到网络波动、API限制、内容过滤等常见错误时,能按策略重试或优雅降级。
我们的目标,就是用一个相对轻量的技术组合,搭建出具备上述核心能力的系统。
2. 核心架构与组件选型
我们不重复造轮子,而是基于成熟的开源工具进行组合。下图展示了系统的核心架构:
[用户/系统] 定义任务
|
v
+----------------------+
| 任务定义层 |
| (YAML/JSON 配置文件) |
+----------------------+
|
v
+----------------------+
| 调度中心 |
| (Apache Airflow) |
+----------------------+
|
v
+----------------------+
| 任务执行引擎 |
| (LangChain + 自定义 Agent)|
+----------------------+
|
v
+----------------------+
| 状态与记忆存储 |
| (PostgreSQL/Redis) |
+----------------------+
|
v
+----------------------+
| 结果输出与通知 |
| (文件/数据库/Webhook)|
+----------------------+
组件详解:
-
调度中心 - Apache Airflow :
- 为什么选它? Airflow 是业界公认的任务编排调度系统,其基于 DAG(有向无环图)的设计理念,完美契合复杂工作流的定义。它提供 Web UI 用于监控、强大的调度能力(cron 表达式)、任务依赖管理、失败重试、邮件告警等开箱即用的功能。
- 替代方案 :如果你觉得 Airflow 较重,可以考虑
Prefect或Dagster,它们更现代,对数据管道友好。但对于需要强调度和运维熟悉的场景,Airflow 社区资源更丰富。
-
任务执行引擎 - LangChain + 自定义 Agent :
- LangChain :它提供了构建 Agent 所需的标准化组件(Tools, Memory, Chains),以及与大模型(OpenAI, Anthropic, 国内大模型等)交互的接口。它是连接“调度”和“AI能力”的桥梁。
- 自定义 Agent :我们不会使用 LangChain 内置的最复杂的 Agent,而是基于
AgentExecutor和ReAct模式,定制一个目标明确、工具集清晰的 Agent。核心是定义好它可用的“工具”(Tools)。
-
状态与记忆存储 - PostgreSQL & Redis :
- PostgreSQL :Airflow 的元数据数据库,同时我们也用它存储任务执行的历史结果、结构化数据。它是“长期记忆”。
- Redis :用于存储 Agent 的会话缓存、临时状态以及作为 Celery(Airflow 执行器)的消息队列。它是“短期记忆”和高速缓存。
-
辅助工具 :
- Docker & Docker Compose :用于容器化部署所有组件,保证环境一致性。
- Prometheus & Grafana (可选):用于监控系统指标,如任务执行时长、成功率、Token 消耗等。
3. 环境准备与快速部署
我们使用 Docker Compose 一键部署最核心的 Airflow 和 PostgreSQL。这是最快体验的方式。
步骤 1:创建项目目录
mkdir ai-agent-loop && cd ai-agent-loop
步骤 2:创建 docker-compose.yml
version: '3.8'
services:
postgres:
image: postgres:13
environment:
POSTGRES_USER: airflow
POSTGRES_PASSWORD: airflow
POSTGRES_DB: airflow
volumes:
- postgres-db-volume:/var/lib/postgresql/data
healthcheck:
test: ["CMD", "pg_isready", "-U", "airflow"]
interval: 10s
retries: 5
start_period: 30s
redis:
image: redis:7-alpine
command: redis-server --appendonly yes
volumes:
- redis-data-volume:/data
healthcheck:
test: ["CMD", "redis-cli", "ping"]
interval: 10s
timeout: 5s
retries: 5
airflow-webserver:
image: apache/airflow:2.7.3
restart: always
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_healthy
environment:
AIRFLOW__CORE__EXECUTOR: CeleryExecutor
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow
AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql://airflow:airflow@postgres/airflow
AIRFLOW__CELERY__BROKER_URL: redis://redis:6379/0
AIRFLOW__CORE__LOAD_EXAMPLES: 'false'
AIRFLOW__WEBSERVER__SECRET_KEY: 'your-secret-key-here' # 生产环境请务必更改!
volumes:
- ./dags:/opt/airflow/dags # 挂载DAG文件目录
- ./logs:/opt/airflow/logs
- ./plugins:/opt/airflow/plugins
- ./config:/opt/airflow/config
ports:
- "8080:8080"
command: webserver
healthcheck:
test: ["CMD", "curl", "--fail", "http://localhost:8080/health"]
interval: 30s
timeout: 10s
retries: 5
airflow-scheduler:
image: apache/airflow:2.7.3
restart: always
depends_on:
- airflow-webserver
environment:
# ... 环境变量与 webserver 相同 ...
AIRFLOW__CORE__EXECUTOR: CeleryExecutor
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow
AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql://airflow:airflow@postgres/airflow
AIRFLOW__CELERY__BROKER_URL: redis://redis:6379/0
AIRFLOW__CORE__LOAD_EXAMPLES: 'false'
volumes:
- ./dags:/opt/airflow/dags
- ./logs:/opt/airflow/logs
- ./plugins:/opt/airflow/plugins
- ./config:/opt/airflow/config
command: scheduler
airflow-worker:
image: apache/airflow:2.7.3
restart: always
depends_on:
- airflow-scheduler
environment:
# ... 环境变量与 webserver 相同 ...
volumes:
- ./dags:/opt/airflow/dags
- ./logs:/opt/airflow/logs
- ./plugins:/opt/airflow/plugins
- ./config:/opt/airflow/config
command: celery worker
volumes:
postgres-db-volume:
redis-data-volume:
步骤 3:初始化 Airflow 数据库 在启动前,需要先初始化数据库。我们可以临时启动一个 Airflow 容器来完成:
# 创建一个 .env 文件存放密钥(可选但推荐)
echo "AIRFLOW_UID=$(id -u)" > .env
echo "AIRFLOW__WEBSERVER__SECRET_KEY=$(openssl rand -hex 32)" >> .env
# 初始化数据库
docker compose up airflow-init
看到 airflow-init_1 exited with code 0 表示成功。
步骤 4:启动所有服务
docker compose up -d
等待几分钟,访问 http://localhost:8080 ,默认账号密码是 airflow / airflow 。你将看到 Airflow 的 Web UI。
至此,调度中心就绪。
4. 定义你的第一个自主 Agent 任务(DAG)
DAG 是 Airflow 的核心,定义了任务的工作流。我们创建一个让 Agent 每天自动分析 GitHub Trending 的 DAG。
步骤 1:创建 DAG 文件 在项目根目录的 dags/ 文件夹下(Docker Compose 已挂载),创建文件 github_trending_agent.py 。
# dags/github_trending_agent.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
import os
import sys
# 将项目根目录加入路径,以便导入自定义模块
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from agents.trending_analyzer import analyze_github_trending
default_args = {
'owner': 'data_team',
'depends_on_past': False,
'email': ['your-email@example.com'],
'email_on_failure': True,
'email_on_retry': False,
'retries': 2,
'retry_delay': timedelta(minutes=3),
}
# 定义 DAG:每天 UTC 时间 10:30(北京时间 18:30)执行
dag = DAG(
'github_trending_daily_analysis',
default_args=default_args,
description='A daily agent to analyze GitHub trending repositories',
schedule_interval='30 10 * * *', # Cron 表达式:每天 10:30 UTC
start_date=datetime(2024, 1, 1),
catchup=False, # 不补抓历史记录
tags=['agent', 'github', 'daily'],
)
# 定义任务:调用我们的 Agent 函数
run_analysis = PythonOperator(
task_id='run_github_trending_analysis',
python_callable=analyze_github_trending,
dag=dag,
)
run_analysis
步骤 2:创建自定义 Agent 模块 在项目根目录创建 agents/ 文件夹和 trending_analyzer.py 文件。
# agents/trending_analyzer.py
import os
import requests
from datetime import datetime
from langchain.agents import AgentExecutor, create_react_agent
from langchain.tools import Tool
from langchain_community.chat_models import ChatOpenAI # 或使用其他模型
from langchain.memory import ConversationBufferMemory
from langchain.prompts import PromptTemplate
from langchain_core.pydantic_v1 import BaseModel, Field
from typing import Type
import json
# ------------------------- 第一步:定义 Agent 可用的工具 -------------------------
# 工具1:获取今日 GitHub Trending 数据
def fetch_github_trending(language: str = "", since: str = "daily") -> str:
"""
从 GitHub Trending API 获取数据。
注意:GitHub 没有官方 API,这是一个模拟函数,实际可能需要爬虫或使用第三方 API。
此处为演示,我们模拟返回数据。
"""
# 模拟数据
mock_data = [
{"name": "repo1", "stars": 450, "language": "Python", "description": "A cool AI project."},
{"name": "repo2", "stars": 320, "language": "JavaScript", "description": "Web framework."},
{"name": "repo3", "stars": 280, "language": "Go", "description": "System tool."},
]
return json.dumps(mock_data, ensure_ascii=False)
# 将函数包装成 LangChain Tool
fetch_trending_tool = Tool(
name="FetchGitHubTrending",
func=fetch_github_trending,
description="获取 GitHub Trending 仓库列表。输入参数:language(可选,如'python'), since(可选,'daily'/'weekly'/'monthly')。返回JSON字符串。"
)
# 工具2:将分析结果保存到文件(或数据库)
def save_analysis_result(analysis_text: str) -> str:
"""将分析结果保存到本地文件。"""
today = datetime.now().strftime("%Y-%m-%d")
filename = f"results/github_analysis_{today}.md"
os.makedirs("results", exist_ok=True)
with open(filename, 'w', encoding='utf-8') as f:
f.write(analysis_text)
return f"分析结果已保存至:{filename}"
save_result_tool = Tool(
name="SaveAnalysisResult",
func=save_analysis_result,
description="将文本分析结果保存到文件。输入参数:analysis_text (字符串)。"
)
# ------------------------- 第二步:配置 LLM 和 Agent -------------------------
# 初始化 LLM,请替换为你的 API Key 和 Base URL
llm = ChatOpenAI(
model="gpt-3.5-turbo", # 或 "gpt-4", "claude-3-haiku" 等
temperature=0,
openai_api_key=os.getenv("OPENAI_API_KEY"), # 从环境变量读取
# 如果使用国内模型,可能需要配置 openai_api_base
# openai_api_base="https://api.xxx.com/v1"
)
# 定义 Agent 的提示词模板
prompt_template = """你是一个专业的 GitHub 趋势分析师。你的任务是:
1. 使用工具获取今日的 GitHub Trending 仓库数据。
2. 分析这些仓库,总结出今日的技术趋势(例如:哪些语言流行?哪些领域热门?)。
3. 挑选出 1-2 个最有意思或最有潜力的仓库,简要说明理由。
4. 将最终的分析报告用中文保存下来。
请按步骤执行。你可以使用的工具有:
{tools}
开始任务前,先明确你的思考过程。
当前对话历史:
{history}
问题:请开始执行今日的 GitHub Trending 分析任务。
思考:"""
prompt = PromptTemplate.from_template(prompt_template)
# 创建 Agent
tools = [fetch_trending_tool, save_result_tool]
agent = create_react_agent(llm, tools, prompt)
# 使用 ConversationBufferMemory 让 Agent 有简单的记忆(记住上一步的结果)
memory = ConversationBufferMemory(memory_key="history", return_messages=True)
agent_executor = AgentExecutor(agent=agent, tools=tools, memory=memory, verbose=True, handle_parsing_errors=True)
# ------------------------- 第三步:定义 Airflow 将要调用的主函数 -------------------------
def analyze_github_trending():
"""Airflow PythonOperator 调用的函数"""
print(f"[{datetime.now()}] 开始执行 GitHub Trending 分析任务...")
try:
# 运行 Agent
result = agent_executor.invoke({"input": "请分析今日的 GitHub Trending。"})
print(f"Agent 执行结果: {result['output']}")
print(f"任务执行成功!")
except Exception as e:
print(f"任务执行失败,错误: {e}")
# 这里可以接入更详细的日志和告警
raise e
if __name__ == "__main__":
# 本地测试时直接运行
analyze_github_trending()
步骤 3:创建结果目录和配置环境变量
mkdir results
echo "OPENAI_API_KEY=sk-your-actual-api-key-here" > .env
# 确保 .env 文件也被加载到 Docker Compose 环境中(需要修改 docker-compose.yml 的 env_file 部分,为简化此处略过,生产环境务必配置)
5. 运行与验证
- 确保服务运行 :
docker compose ps查看所有容器状态应为Up。 - 触发 DAG :在 Airflow Web UI (
localhost:8080) 中,找到github_trending_daily_analysisDAG,点击左侧的 “播放”按钮 手动触发一次执行。 - 查看日志 :点击运行的任务实例,查看日志。你应该能看到类似以下的输出:
[2024-05-27, 10:30:00 UTC] 开始执行 GitHub Trending 分析任务... > Entering new AgentExecutor chain... 思考:我需要先获取数据,然后分析,最后保存。 行动:FetchGitHubTrending 观察:[{"name": "repo1", "stars": 450...}] 思考:现在我有了数据,我需要分析趋势... ...(分析过程)... 行动:SaveAnalysisResult 观察:分析结果已保存至:results/github_analysis_2024-05-27.md 思考:任务完成。 > Finished chain. Agent 执行结果: 已完成今日 GitHub Trending 分析,报告已保存。 - 检查结果 :查看
results/目录下生成的 Markdown 文件,里面就是 AI 生成的趋势分析报告。
至此,一个最简单的自主循环 Agent 系统就跑通了。它会在每天指定时间自动运行,无需你手动干预。
6. 系统扩展:让 Agent 真正“会自己找活”
上面的例子是定时任务。但“自主”意味着它能根据情况 主动 触发。我们可以通过以下几种模式扩展:
6.1 事件驱动模式
让 Agent 监听外部事件。例如,当监控系统发现某个 API 错误率飙升时,自动触发一个“故障诊断 Agent”。
实现思路 :
- 使用 Airflow 的
TriggerDagRunOperator或外部触发器(如 Webhook)。 - 创建一个接收 Webhook 的简单 Flask/FastAPI 服务,在收到事件后,通过 Airflow REST API 触发对应的 DAG。
- DAG 接收事件参数,并传递给 Agent。
6.2 条件循环与自省模式
Agent 完成任务后,根据结果决定下一步做什么。例如,一个“竞品监控 Agent”发现某竞品发布了新功能,它可以自动触发一个“深度分析该功能”的子任务。
实现思路 :
- 在 Agent 的工具集中,增加一个
create_sub_task(topic, priority)工具。 - 这个工具会向一个任务队列(如 Redis List)或数据库写入一条新任务记录。
- 另一个 Airflow DAG(或同一个 DAG 中的另一个传感器)定期扫描这个队列,发现有新任务时就触发执行。
6.3 记忆与上下文增强
让 Agent 记住更长的历史。上面的 ConversationBufferMemory 只存在于单次执行。要实现跨周期记忆,需要将记忆存储到外部数据库。
实现思路 :
- 使用
LangChain的PostgresChatMessageHistory或RedisChatMessageHistory。 - 为每个任务类型或任务实例创建一个唯一的
session_id(例如github_trending_20240527)。 - 每次执行时,Agent 加载该
session_id的历史消息,执行后再保存新的对话记录。
# 示例:使用 PostgreSQL 存储记忆
from langchain_community.chat_message_histories import PostgresChatMessageHistory
from langchain.memory import ConversationBufferMemory
def get_agent_with_persistent_memory(session_id: str):
message_history = PostgresChatMessageHistory(
session_id=session_id,
connection_string="postgresql://airflow:airflow@postgres/airflow", # 复用 Airflow 数据库
table_name="agent_message_store"
)
memory = ConversationBufferMemory(
memory_key="chat_history",
chat_memory=message_history,
return_messages=True
)
# ... 用这个 memory 创建 agent_executor
return agent_executor
7. 生产环境最佳实践与避坑指南
运行了三个月,我们踩过不少坑,总结出以下关键点:
| 问题领域 | 具体问题 | 解决方案与最佳实践 |
|---|---|---|
| API 稳定性与成本 | 1. 大模型 API 调用失败或超时。 2. Token 消耗不可控,成本飙升。 |
1. 重试与退避 :在 Agent 调用层实现指数退避重试。 2. 设置预算与监控 :使用模型的 usage 字段监控 token 消耗,设置每日/每月预算,超限后自动切换模型或暂停任务。 3. 缓存 :对频繁查询且结果不变的工具(如某些数据查询)增加缓存层。 |
| 任务管理与监控 | 1. 某个 Agent 任务卡死或长时间运行。 2. 难以追溯某次执行的具体输入输出。 |
1. 超时设置 :在 Airflow Operator 和 LangChain AgentExecutor 中都设置 max_execution_time 。 2. 全面日志 :将 LangChain 的 verbose 输出重定向到 Airflow 任务日志或集中式日志系统(如 ELK)。 3. 结构化存储 :将重要的中间结果和最终结果,不仅保存为文件,也结构化地存入数据库,便于查询分析。 |
| 错误处理与自愈 | 1. Agent 因内容过滤政策(Content Policy)中断。 2. 依赖的外部 API 不可用。 |
1. 结构化输出与验证 :要求 LLM 以 JSON 等格式输出,并在代码中验证,解析失败则重试或降级。 2. 备选工具/降级策略 :例如,获取数据的主 API 失败,自动切换到备用 API 或使用缓存的昨日数据。 3. 敏感词过滤 :在将用户输入或外部数据传给 LLM 前,进行一层基本的敏感词过滤,减少触发内容政策的风险。 |
| 安全与权限 | 1. Agent 工具权限过大(如能执行任意 Shell 命令)。 2. 敏感信息(API Keys)泄露。 |
1. 最小权限原则 :严格限制每个工具的能力。例如,写文件工具只能写到特定目录;执行命令工具必须限制命令白名单。 2. 秘密管理 :使用 Airflow 的 Variables 或 Connections 功能管理 API Key,或使用专业的秘密管理工具(如 HashiCorp Vault),绝不硬编码在代码中。 |
| 性能与扩展 | 1. 大量并发任务时,系统负载高。 2. Agent 思考(Reasoning)步骤多,执行慢。 |
1. 队列与资源隔离 :在 Airflow 中配置不同的队列(Queues),将重量级和轻量级任务分配到不同的 Worker。 2. 优化 Prompt 与工具 :精简 Prompt,避免开放式问题。设计更精准的工具,减少 LLM 不必要的思考循环。对于复杂任务,可以拆分成多个子 DAG 并行执行。 |
8. 总结:从“手动 Prompt”到“自主系统”的思维转变
搭建这样一个系统,最大的价值不在于技术栈的堆砌,而在于 思维模式的转变 。我们不再将 AI Agent 视为一个需要精心呵护、每次手动投喂的“展示品”,而是将其定位为一个可批量部署、有明确职责、能持续运行的“数字员工”。
这套系统的核心收益:
- 效率规模化 :从一个 Agent 到一百个 Agent,管理成本并不会线性增长。Airflow 的 UI 和日志提供了统一的管理平面。
- 过程可观测 :每一次执行都有记录、有日志、有结果。调试和优化从“黑盒”变成了“白盒”。
- 能力可复用 :定义好的工具(Tools)和 Agent 模板,可以像乐高一样快速组合成新的任务流。
- 可靠性提升 :通过调度系统的重试、告警、依赖管理,整个 AI 自动化流程的鲁棒性远超手动运行脚本。
给你的行动建议:
- 从小处着手 :不要一开始就规划庞大的系统。从你 最重复、最枯燥 的一个日常任务开始(比如每日数据汇总、信息收集),用上述框架实现它。
- 先跑通,再优化 :先用最简单的
PythonOperator+ 基础LangChain Agent实现核心逻辑。让整个循环先转起来,再去考虑事件驱动、记忆持久化等高级特性。 - 重视工具设计 :Agent 的能力边界由其工具集决定。花时间思考并封装好稳定、安全、功能单一的工具,这是整个系统长期稳定的基石。
- 建立监控看板 :至少要将 Airflow 的任务状态、LLM API 的调用次数和 Token 消耗监控起来。数据能帮你发现瓶颈和异常。
技术总是在迭代,但将不确定性高的 AI 能力封装进确定性高的自动化流程中,这个思路会长期有效。希望这套经过三个月实战检验的“自主循环”系统设计,能为你打开一扇门,让你手中的 AI Agent 真正不知疲倦地为你工作。
更多推荐



所有评论(0)