AI Engineering Academy · 课时

跨执行步骤管理状态

在多个代码执行步骤之间持久化变量、数据框和已导入的库,使智能体能够基于之前的结果继续工作,而无需重新运行早先的计算。

第 3 / 4 课13 个步骤

跨执行步骤管理状态 是 CoddyKit 上的免费 AI Engineering Academy 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 AI Engineering Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 AI Engineering Academy 课程共包含 4 节课。

无状态问题

每次 Docker 容器执行都会启动一个全新的 Python 解释器。第 1 轮中定义的变量在第 2 轮中并不存在。这种无状态性迫使代理在每个执行步骤中从头重新计算或重新加载所有内容——除非您实现明确的状态管理策略,以便在多轮执行之间持久化和恢复状态。没有这种机制,多步骤分析将无法完成。

# Iteration 1 - works fine
df = pd.read_csv('data.csv')  # df is in memory
df_cleaned = df.dropna()
print('Rows after cleaning:', len(df_cleaned))

# Iteration 2 - NEW container, df_cleaned is GONE
result = df_cleaned.groupby('category').sum()  # NameError: df_cleaned is not defined
print(result)  # This will fail!

基于文件的状态持久化

最简单且可移植的方法是将状态保存到文件中。在每个代码块结束时,代理将 DataFrames、字典或其他对象保存到工作区目录中的文件里。下一轮执行时再将它们加载回来。Parquet 非常适合保存 DataFrames,JSON 适合保存字典,而 pickle 适合保存任意 Python 对象(不过,从不受信任的代码中使用 pickle 存在安全风险)。

import pandas as pd
import json
from pathlib import Path

WORKSPACE = Path('/workspace')

# Iteration 1: process and SAVE
df = pd.read_csv(WORKSPACE / 'raw_data.csv')
df_cleaned = df.dropna().reset_index(drop=True)
df_cleaned.to_parquet(WORKSPACE / 'cleaned.parquet')  # save for next iteration

stats = {'rows': len(df_cleaned), 'columns': list(df_cleaned.columns)}
with open(WORKSPACE / 'stats.json', 'w') as f:
    json.dump(stats, f)

print('Saved cleaned data:', len(df_cleaned), 'rows')

# Iteration 2: LOAD and continue
df_cleaned = pd.read_parquet(WORKSPACE / 'cleaned.parquet')  # restore state
with open(WORKSPACE / 'stats.json') as f:
    stats = json.load(f)
result = df_cleaned.groupby('category')['value'].sum()
print(result)

教会代理使用基于文件的状态

您不能只在基础设施中实现基于文件的状态管理——LLM 也必须了解这一约定并始终遵守。请在系统提示中加入明确指令:计算出之后还会用到的结果后,务必使用描述性文件名将其保存;每轮开始时,务必先加载之前步骤生成的文件,再继续执行。

STATE_MANAGEMENT_INSTRUCTIONS = '''
State persistence rules:
- You have a persistent workspace at /workspace/
- After computing any result you will need later, SAVE it to /workspace/
  - DataFrames: use .to_parquet('/workspace/name.parquet')
  - Dicts/lists: use json.dump to /workspace/name.json
  - Text: write to /workspace/name.txt
- At the start of each code block, LOAD the files you need from previous steps
- Use clear, descriptive filenames like 'cleaned_data.parquet', not 'tmp1.parquet'
- When listing files, use: import os; print(os.listdir('/workspace/'))
'''

使用持久化 Python 进程的基于内核的状态

基于文件的状态之外,另一种方案是让一个持久化 Python 进程(例如 Jupyter 内核)在多轮执行之间持续运行,并通过 exec() 将每个代码块注入其中。在一轮中定义的变量,在下一轮中仍可访问。这种方法速度更快,也更自然,但要求每个代理会话运行一个持久化进程,而不是使用临时容器。

import jupyter_client

class PersistentKernel:
    def __init__(self):
        km, self.kc = jupyter_client.manager.start_new_kernel(kernel_name='python3')
        self.kc.wait_for_ready(timeout=30)
        print('Kernel started')

    def execute(self, code: str, timeout=60) -> tuple[str, str]:
        msg_id = self.kc.execute(code)
        outputs, errors = [], []
        
        while True:
            msg = self.kc.get_iopub_msg(timeout=timeout)
            if msg['msg_type'] == 'stream':
                if msg['content']['name'] == 'stdout':
                    outputs.append(msg['content']['text'])
                else:
                    errors.append(msg['content']['text'])
            if msg['msg_type'] == 'status' and msg['content']['execution_state'] == 'idle':
                break
        
        return ''.join(outputs), ''.join(errors)

    def shutdown(self):
        self.kc.shutdown()

# Variables persist across execute() calls!
kernel = PersistentKernel()
kernel.execute('x = 42')  # x is now defined in kernel
output, _ = kernel.execute('print(x)')  # prints 42 - state persists!
kernel.shutdown()

使用状态对象管理状态

对于代理编排层(位于沙箱之外),请维护一个明确的状态对象,用于跟踪代理工作的当前状态:已经创建了哪些文件、已经完成了哪些计算、代理当前处于哪一步,以及最终目标是什么。在每轮执行中,将该状态对象的摘要传递给 LLM,使其始终了解自己在整体任务中的位置。

from dataclasses import dataclass, field
from typing import Any

@dataclass
class ExecutionState:
    task: str
    iteration: int = 0
    workspace_files: list[str] = field(default_factory=list)
    computed_values: dict[str, Any] = field(default_factory=dict)
    completed_steps: list[str] = field(default_factory=list)
    last_output: str = ''

    def to_context_summary(self) -> str:
        return f'''Current task: {self.task}
Iteration: {self.iteration}
Completed steps: {', '.join(self.completed_steps) or 'None yet'}
Workspace files: {', '.join(self.workspace_files) or 'None yet'}
Key values: {self.computed_values}
Last output: {self.last_output[:500]}'''

    def after_execution(self, output: str, step_name: str):
        import os
        self.workspace_files = os.listdir('/workspace')
        self.completed_steps.append(step_name)
        self.last_output = output
        self.iteration += 1

将状态上下文注入提示

每次调用 LLM 获取下一个代码块时,都要将当前状态上下文注入用户消息。这样,LLM 就能准确了解哪些文件可用、哪些内容已经计算完成,以及还有哪些工作待完成。没有这些上下文,LLM 可能会尝试重新创建已经存在的文件,或者跳过已经完成的步骤。

def build_iteration_prompt(task: str, state: ExecutionState, last_observation: str) -> str:
    return f'''ORIGINAL TASK: {task}

CURRENT STATE:
{state.to_context_summary()}

LAST EXECUTION OUTPUT:
{last_observation}

What is the next step? Write Python code to continue. 
Remember:
- Load data from workspace files if needed
- Save any results you will need in future steps
- If the task is complete, say TASK COMPLETE and summarize the result'''

# Use in the main loop
for step_num in range(max_iterations):
    context = build_iteration_prompt(task, state, last_observation)
    response = llm.complete(messages + [{'role': 'user', 'content': context}])
    code = extract_code_block(response)
    if not code:
        break  # done
    output, err = execute_in_sandbox(code)
    state.after_execution(output, step_name=f'step_{step_num}')
    last_observation = format_observation(output, err)

处理大型中间数据

数据分析代理通常会生成大型中间数据集,而在每轮执行中重新加载这些数据的成本很高。请采用延迟加载:只加载当前步骤所需的数据。使用 Parquet 等支持高效按列读取的列式格式。对于真正大型的数据集(100MB 以上),请保留一个持久化内核,让 DataFrame 在多轮执行之间驻留在内存中,而不是每次都进行序列化和反序列化。

# Good: load only needed columns
df = pd.read_parquet('/workspace/full_data.parquet', columns=['date', 'revenue', 'region'])

# Bad: load everything even if you only need 2 columns
# df = pd.read_parquet('/workspace/full_data.parquet')  # loads 50 columns you don't need

# Good: filter early before loading full dataset
df = pd.read_parquet('/workspace/full_data.parquet', filters=[('region', '=', 'EMEA')])

# For very large files, tell the LLM about data shape upfront
df_info = {'shape': (1_000_000, 50), 'size_mb': 850, 'columns': [...]}
# Include df_info in state so LLM plans accordingly

为长时间运行的任务创建检查点

对于跨越多轮执行的任务,请实现检查点机制:定期将完整的代理状态(已完成的步骤、工作区文件、当前进度)保存到持久存储中,以便任务在中断后恢复。这对于可能需要 30 分钟以上、并且可能因网络问题、API 错误或服务器重启而中断的长时间数据分析任务尤其重要。

import json
import time

def checkpoint_state(state: ExecutionState, checkpoint_id: str, db):
    data = {
        'task': state.task,
        'iteration': state.iteration,
        'completed_steps': state.completed_steps,
        'workspace_files': state.workspace_files,
        'timestamp': time.time()
    }
    db.execute(
        'INSERT INTO agent_checkpoints (id, state) VALUES (?, ?) ON CONFLICT(id) DO UPDATE SET state=excluded.state',
        (checkpoint_id, json.dumps(data))
    )
    print(f'Checkpoint saved at iteration {state.iteration}')

def resume_from_checkpoint(checkpoint_id: str, db) -> ExecutionState | None:
    row = db.execute('SELECT state FROM agent_checkpoints WHERE id = ?', (checkpoint_id,)).fetchone()
    if not row:
        return None
    data = json.loads(row[0])
    state = ExecutionState(task=data['task'])
    state.iteration = data['iteration']
    state.completed_steps = data['completed_steps']
    print(f'Resumed from iteration {state.iteration}')
    return state

完成后清理状态

代理工作区会不断累积文件,并且可能随时间变得很大。请始终实现一个清理阶段,在任务完成或失败时运行:删除中间文件(cleaned.parquet、tmp_output.csv),只保留用户真正需要的最终输出文件。对于云存储,请设置生命周期策略,在保留期限结束后自动删除工作区文件。

import shutil
from pathlib import Path

INTERMEDIATE_PATTERNS = ['*.parquet', 'tmp_*.csv', 'step_*.json', 'debug_*.txt']
FINAL_OUTPUT_PATTERNS = ['report.pdf', 'final_*.csv', 'summary.json']

def cleanup_workspace(workspace_dir: str, keep_final=True):
    workspace = Path(workspace_dir)
    final_outputs = []
    
    if keep_final:
        for pattern in FINAL_OUTPUT_PATTERNS:
            final_outputs.extend(workspace.glob(pattern))
        # Move final outputs to output directory
        output_dir = workspace.parent / 'outputs'
        output_dir.mkdir(exist_ok=True)
        for f in final_outputs:
            shutil.move(str(f), str(output_dir / f.name))
    
    # Delete the workspace
    shutil.rmtree(workspace_dir)
    print(f'Workspace cleaned up. Kept {len(final_outputs)} output files.')
    return [str(f) for f in final_outputs]

用于调试的状态版本控制

当代码代理产生错误结果时,您需要沿着执行历史回溯,以找出问题所在。请通过在每轮执行后保存工作区快照来实现状态版本控制。这样,您就可以重放代理的执行过程、检查中间状态,并确定错误发生的具体轮次。为节省空间,只保存差异内容(发生变化的文件)。

上下文窗口与外部状态

在代码代理设计中,将状态放入LLM 上下文窗口(即时但有限)还是放入外部存储(容量不限但需要明确管理)之间存在根本性的权衡。最佳策略是:将当前步骤的数据保留在上下文中,将大型中间结果保存在文件中,将任务目标和高层计划保留在上下文中,并始终将原始数据保存在磁盘上。LLM 管理任务,文件系统管理数据。

快速检查

请检验您对本课跨代码执行步骤状态管理的理解。

课程回顾

本课您学到了:使用 Parquet 和 JSON 实现的基于文件的状态持久化,是 ​​在临时容器执行之间共享状态的最具可移植性的方法;持久化内核通过让 Python 进程在多轮执行之间保持运行,消除了重新加载的开销;状态上下文注入使 LLM 在每轮执行时都能准确了解可用文件和已完成步骤。接下来我们将构建一个完整的数据分析代理。

免费开始

用 AI 导师学习 Python — 免费

在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。

课程
30
课程
120

常见问题解答

「跨执行步骤管理状态」课时是免费的吗?

是的 — 「跨执行步骤管理状态」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 AI Engineering Academy 课程的其余内容,请升级到 CoddyKit PRO。 AI Engineering Academy 课程共包含 4 节课。

「跨执行步骤管理状态」这节课中我会学到什么?

在多个代码执行步骤之间持久化变量、数据框和已导入的库,使智能体能够基于之前的结果继续工作,而无需重新运行早先的计算。 你通过在浏览器中直接运行的动手代码来练习 AI Engineering Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 AI Engineering Academy 需要有经验吗?

无需任何先前经验。CoddyKit 上的 AI Engineering Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。

「跨执行步骤管理状态」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 AI Engineering Academy 课中编写并运行代码吗?

能。每节 AI Engineering Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 代码执行循环
  2. 使用 Docker 和 RestrictedPython 实现沙箱隔离
  3. 跨执行步骤管理状态
  4. 构建数据分析智能体
← 返回 AI Engineering Academy