17. n8n 架构深度分析

作者:

在

n8n 是一个 Fair-code 工作流自动化平台,原生支持 AI Agent 能力。通过可视化画布结合自定义代码,连接 1500+ 集成,支持自托管或云端部署。本文从源码层面深入分析其架构设计。

1. 版本化节点系统(VersionedNodeType)

n8n 采用版本化节点架构,所有节点(包括 AI Agent)都通过 VersionedNodeType 基类管理多版本。这种设计允许节点在不破坏现有工作流的情况下持续演进。

源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/Agent.node.ts

export class Agent extends VersionedNodeType {
    constructor() {
        const baseDescription: INodeTypeBaseDescription = {
            displayName: 'AI Agent',
            name: 'agent',
            icon: 'node:ai-agent',
            iconColor: 'black',
            group: ['transform'],
            description: 'Generates an action plan and executes it. Can use external tools.',
            codex: {
                alias: ['LangChain', 'Chat', 'Conversational', 'Plan and Execute', 'ReAct', 'Tools'],
                categories: ['AI'],
                subcategories: {
                    AI: ['Agents', 'Root Nodes'],
                },
            },
            defaultVersion: 3.1,
        };

        const nodeVersions: IVersionedNodeType['nodeVersions'] = {
            1: new AgentV1(baseDescription),
            1.1: new AgentV1(baseDescription),
            // ... 1.2-1.9 均为 AgentV1
            2: new AgentV2(baseDescription),
            2.1: new AgentV2(baseDescription),
            // ... 2.2-2.3 均为 AgentV2
            3: new AgentV3(baseDescription),
            3.1: new AgentV3(baseDescription),
        };
        super(nodeVersions, baseDescription);
    }
}

架构洞察: 版本映射表将语义版本号映射到具体实现类。V1 系列(1.0-1.9)共享同一实现,V2 和 V3 各有独立实现。defaultVersion: 3.1 表示新建工作流默认使用最新版本。这种设计使得旧工作流可以继续运行,同时新功能只在新版本中提供。

2. 核心执行引擎(WorkflowExecute)

执行引擎是 n8n 的心脏,负责工作流的调度和执行。它使用 PCancelable 实现可取消的异步执行,支持部分执行(Partial Execution)优化。

源码位置: packages/core/src/execution-engine/workflow-execute.ts

export class WorkflowExecute {
    private status: ExecutionStatus = 'new';
    private readonly abortController = new AbortController();
    timedOut: boolean = false;

    constructor(
        private readonly additionalData: IWorkflowExecuteAdditionalData,
        private readonly mode: WorkflowExecuteMode,
        private runExecutionData: IRunExecutionData = createRunExecutionData(),
        private readonly storedAt: ExecutionStorageLocation = 'db',
    ) {}

    // IMPORTANT: Do not add "async" to this function, it will then convert the
    //            PCancelable to a regular Promise and does so not allow canceling
    //            active executions anymore
    run({ workflow, startNode, destinationNode, pinData, triggerToStartFrom,
          additionalRunFilterNodes }: RunWorkflowOptions): PCancelable {
        this.status = 'running';
        // Get the nodes to start workflow execution from
        startNode = startNode || workflow.getStartNode(destinationNode?.nodeName);
        // ...
        return this.processRunExecutionData(workflow);
    }
}

关键设计:

  • PCancelable: 故意不使用 async 关键字,以保持 PCancelable 的可取消特性
  • AbortController: 支持执行超时和手动取消
  • 状态机: new → running → ... 的状态转换
  • 部分执行: runPartialWorkflow2 方法通过有向图分析,只执行必要的节点子图

3. 部分执行与有向图分析

n8n 的部分执行系统是最精妙的设计之一。当用户手动执行某个节点时,引擎会构建有向图,找到触发器到目标节点的最小子图,避免执行整个工作流。

源码位置: packages/core/src/execution-engine/workflow-execute.ts

runPartialWorkflow2(workflow, runData, pinData, dirtyNodeNames, destinationNode,
                     agentRequest?): PCancelable {
    // 1. 构建有向图
    let graph = DirectedGraph.fromWorkflow(workflow);

    // 2. 工具节点特殊处理:重连图结构
    if (NodeHelpers.isTool(destinationNodeType.description, destination.parameters)) {
        graph = rewireGraph(destination, graph, agentRequest);
        workflow = graph.toWorkflow({ ...workflow });
    }

    // 3. 查找触发器
    let trigger = findTriggerForPartialExecution(workflow, destinationNode.nodeName, runData);

    // 4. 查找子图
    const filteredGraph = filterDisabledNodes(graph);
    graph = findSubgraph({ graph: filteredGraph, destination, trigger });

    // 5. 查找起始节点
    let startNodes = findStartNodes({ graph, trigger, destination, runData, pinData });

    // 6. 检测并处理循环
    startNodes = handleCycles(graph, startNodes, trigger);

    // 7. 重建执行栈
    const { nodeExecutionStack, waitingExecution, waitingExecutionSource } =
        recreateNodeExecutionStack(graph, startNodes, runData, pinData ?? {});

    // 8. 执行
    return this.processRunExecutionData(workflow);
}

算法流程: DirectedGraph.fromWorkflow → filterDisabledNodes → findSubgraph → findStartNodes → handleCycles → recreateNodeExecutionStack。这是一套完整的图论算法链,确保最小化执行范围。

4. LangChain Agent 集成架构

n8n 通过 LangChain 集成实现了 AI Agent 能力。Agent 节点作为根节点,通过子节点连接(Subnodes)获取模型、工具和记忆。

源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/common.ts

export async function getChatModel(
    ctx: IExecuteFunctions | ISupplyDataFunctions | IWebhookFunctions,
    index: number = 0,
): Promise {
    const connectedModels = await ctx.getInputConnectionData(
        NodeConnectionTypes.AiLanguageModel, 0
    );

    let model;
    if (Array.isArray(connectedModels) && index !== undefined) {
        if (connectedModels.length <= index) return undefined;
        // We get the models in reversed order from the workflow
        const reversedModels = [...connectedModels].reverse();
        model = reversedModels[index] as BaseChatModel;
    } else {
        model = connectedModels as BaseChatModel;
    }

    if (!isChatInstance(model) || !model.bindTools) {
        throw new NodeOperationError(ctx.getNode(),
            'Tools Agent requires Chat Model which supports Tools calling');
    }
    return model;
}

连接模型: NodeConnectionTypes.AiLanguageModel 定义了 AI 模型连接类型。模型通过 getInputConnectionData 从子节点获取,支持主模型和备用模型(Fallback Model)。

5. Agent 序列构建与执行

Agent 的核心执行逻辑封装在 createAgentSequence 中,它将 LangChain 的 createToolCallingAgent 与输出解析器、备用模型组合成可执行的序列。

源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/V3/helpers/createAgentSequence.ts

export function createAgentSequence(
    model: BaseChatModel,
    tools: Array,
    prompt: ChatPromptTemplate,
    _options: { maxIterations?: number; returnIntermediateSteps?: boolean },
    outputParser?: N8nOutputParser,
    memory?: BaseChatMemory,
    fallbackModel?: BaseChatModel | null,
    forceToolCall = false,
) {
    const allTools = getAllTools(model, tools);
    const agent = createToolCallingAgent({
        llm: forceToolCall && model.bindTools
            ? model.bindTools(allTools, { tool_choice: 'any' })
            : model,
        tools: allTools,
        prompt,
        streamRunnable: false,
    });

    let fallbackAgent: AgentRunnableSequence | undefined;
    if (fallbackModel) {
        fallbackAgent = createToolCallingAgent({
            llm: forceToolCall && fallbackModel.bindTools
                ? fallbackModel.bindTools(fallbackTools, { tool_choice: 'any' })
                : fallbackModel,
            tools: fallbackTools,
            prompt,
            streamRunnable: false,
        });
    }

    const runnableAgent = RunnableSequence.from([
        fallbackAgent ? agent.withFallbacks([fallbackAgent]) : agent,
        getAgentStepsParser(outputParser, memory),
        fixEmptyContentMessage,
    ]) as AgentRunnableSequence;

    runnableAgent.singleAction = true;
    runnableAgent.streamRunnable = false;
    return runnableAgent;
}

设计要点:

  • forceToolCall: 强制模型调用工具(tool_choice: 'any'),用于首次迭代
  • withFallbacks: LangChain 原生的备用模型机制,主模型失败时自动切换
  • singleAction = true: 单次动作模式,每次只调用一个工具
  • 输出链: Agent → StepsParser → FixEmptyContent

6. 批处理与并发执行

V3 版本引入了批处理机制,支持配置批量大小和批次间延迟,通过 Promise.allSettled 实现并发执行。

源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/V3/helpers/executeBatch.ts

export async function executeBatch(
    ctx: IExecuteFunctions | ISupplyDataFunctions,
    batch: INodeExecutionData[],
    startIndex: number,
    model: BaseChatModel,
    fallbackModel: BaseChatModel | null,
    memory: BaseChatMemory | undefined,
    response?: EngineResponse,
): Promise {
    // 处理 HITL(Human-in-the-Loop)工具响应
    const hitlResult = processHitlResponses(response, startIndex);
    if (hitlResult.hasApprovedHitlTools && hitlResult.pendingGatedToolRequest) {
        return { returnData: [], request: hitlResult.pendingGatedToolRequest, memoryHits };
    }

    // 检查最大迭代次数
    const maxIterations = ctx.getNodeParameter('options.maxIterations', 0, 10);

    // 并发处理批次中的所有项目
    const batchPromises = batch.map(async (_item, batchItemIndex) => {
        const itemIndex = startIndex + batchItemIndex;
        checkMaxIterations(response, maxIterations, ctx.getNode());
        const itemContext = await prepareItemContext(ctx, itemIndex, processedResponse, model);
        const executor = createAgentSequence(model, tools, prompt, options,
                                              outputParser, memory, fallbackModel, ...);
        return await runAgent(ctx, executor, itemContext, model, memory,
                              processedResponse, memoryHits);
    });

    const batchResults = await Promise.allSettled(batchPromises);
    // ... 处理结果
}

并发策略: 使用 Promise.allSettled 而非 Promise.all,确保单个项目的失败不会影响整个批次。continueOnFail() 机制允许在错误时返回错误信息而非抛出异常。

7. 流式执行与事件流处理

Agent 支持流式输出,通过 LangChain 的 streamEvents API 实现实时响应。

源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/V3/helpers/runAgent.ts

export async function runAgent(
    ctx: IExecuteFunctions | ISupplyDataFunctions,
    executor: AgentRunnableSequence,
    itemContext: ItemContext,
    model: BaseChatModel,
    memory: BaseChatMemory | undefined,
    response?: EngineResponse,
    memoryHits?: { loads: number; saves: number },
): Promise {
    const isStreamingAvailable = 'isStreaming' in ctx ? ctx.isStreaming?.() : undefined;

    if ('isStreaming' in ctx && options.enableStreaming && isStreamingAvailable
        && ctx.getNode().typeVersion >= 2.1) {
        // 流式执行路径
        const chatHistory = await loadMemory(memory, model, options.maxTokensFromMemory);
        const eventStream = executorWithTracing.streamEvents(
            { ...invokeParams, chat_history: chatHistory },
            { version: 'v2', ...executeOptions },
        );
        const result = await processEventStream(ctx, eventStream, itemIndex, finalizeOutput);

        // 如果结果包含工具调用,构建请求对象
        if (result.toolCalls && result.toolCalls.length > 0) {
            const actions = createEngineRequests(result.toolCalls, itemIndex, tools);
            return { actions, metadata: buildResponseMetadata(response, itemIndex) };
        }
        return result;
    } else {
        // 非流式执行路径
        const modelResponse = await executorWithTracing.invoke(
            { ...invokeParams, chat_history: chatHistory }, executeOptions);
        // ...
    }
}

双路径设计: 流式和非流式执行共享相同的输入准备逻辑,但在执行阶段分叉。流式路径使用 streamEvents API,非流式使用 invoke。两种路径都支持工具调用的中断和恢复。

8. 子 Agent(Sub-Agent)与工具执行

n8n 支持多 Agent 协作,通过 AgentTool 节点实现 Agent 作为工具被其他 Agent 调用。子 Agent 的工具执行通过内联解析而非引擎调度。

源码位置: packages/@n8n/nodes-langchain/nodes/agents/Agent/agents/ToolsAgent/V3/helpers/resolveSubAgentRequest.ts

export async function resolveSubAgentRequest(
    ctx: ISupplyDataFunctions,
    request: EngineRequest,
    deps: ResolveSubAgentRequestDeps,
): Promise {
    let current: INodeExecutionData[][] | EngineRequest = request;
    const node = ctx.getNode();
    const tools = (await getTools(ctx)) as ConnectedTool[];

    while (isEngineRequest(current)) {
        assertNoHitlActions(node, current.actions);
        ctx.getExecutionCancelSignal?.()?.throwIfAborted?.();

        // 内联执行所有工具调用
        const actionResponses = await Promise.all(
            current.actions.map(async (action) =>
                await executeEngineAction(node, action, tools)),
        );

        // 递归调用 Agent 执行
        current = await deps.runAgentBatch({
            actionResponses, metadata: current.metadata
        });
    }

    return current;
}

function assertNoHitlActions(node: INode, actions: Action[]): void {
    for (const action of actions) {
        if (action.metadata?.hitl) {
            throw new UserError(
                `Human-in-the-Loop nodes cannot be used inside a sub-agent. Move "${action.nodeName}" to the top-level workflow.`);
        }
    }
}

关键限制: 子 Agent 不支持 HITL(Human-in-the-Loop)节点,因为审批流程需要顶层引擎的参与。这是一个合理的架构约束。

9. 路由节点与声明式 API 集成

n8n 的 RoutingNode 类实现了声明式 REST API 集成,支持分页、批处理和凭证管理。

源码位置: packages/core/src/execution-engine/routing-node.ts

export class RoutingNode {
    constructor(
        private readonly context: ExecuteContext,
        private readonly nodeType: INodeType,
        private readonly credentialsDecrypted?: ICredentialsDecrypted,
    ) {}

    async runNode(): Promise {
        const items = (inputData[NodeConnectionTypes.Main] ??
            inputData[NodeConnectionTypes.AiTool])[0] as INodeExecutionData[];
        const returnData: INodeExecutionData[] = [];

        const { credentials, credentialDescription } = await this.prepareCredentials();

        const { batching } = context.getNodeParameter('requestOptions', 0, {}) as {
            batching: { batch: { batchSize: number; batchInterval: number } };
        };

        const batchSize = batching?.batch?.batchSize > 0 ? batching?.batch?.batchSize : 1;
        const batchInterval = batching?.batch.batchInterval;

        for (let itemIndex = 0; itemIndex  0 && batchSize >= 0 && batchInterval > 0) {
                if (itemIndex % batchSize === 0) {
                    await sleep(batchInterval);
                }
            }
            // ...
        }
    }
}

设计模式: 路由节点支持 Main 和 AiTool 两种输入连接类型,使得同一个节点既可以作为工作流节点使用,也可以作为 AI Agent 的工具使用。

10. 节点加载器与插件系统

n8n 的节点加载器支持从目录动态加载节点和凭证类型,支持自定义节点和第三方包。

源码位置: packages/core/src/nodes-loader/directory-loader.ts

export abstract class DirectoryLoader implements NodeLoader {
    isLazyLoaded = false;
    loadedNodes: INodeTypeNameVersion[] = [];
    nodeTypes: INodeTypeData = {};
    credentialTypes: ICredentialTypeData = {};
    known: KnownNodesAndCredentials = { nodes: {}, credentials: {} };
    types: Types = { nodes: [], credentials: [] };
    nodesByCredential: Record = {};

    constructor(
        readonly directory: string,
        protected excludeNodes: string[] = [],
        protected includeNodes: string[] = [],
    ) {
        // If `directory` is a symlink, we try to resolve it to its real path
        try {
            this.directory = realpathSync(directory);
        } catch (error) {
            if (error.code !== 'ENOENT') throw error;
        }
    }
}

插件架构: DirectoryLoader 是抽象基类,具体实现包括 PackageDirectoryLoader(npm 包)、LazyPackageDirectoryLoader(懒加载)和 CustomDirectoryLoader(自定义目录)。known 字段存储节点和凭证的位置信息,支持懒加载场景。


架构总结

n8n 的架构设计体现了以下核心理念:

  1. 版本化演进: 节点通过版本号管理,旧工作流永不被破坏
  2. 图论驱动执行: 有向图分析实现最小化部分执行
  3. LangChain 深度集成: 不是简单封装,而是深度利用 LangChain 的 Agent、Tool、Memory 原语
  4. 可取消异步: PCancelable + AbortController 实现优雅的执行控制
  5. 子 Agent 协作: 通过内联解析实现多 Agent 嵌套调用
  6. 声明式集成: RoutingNode 支持声明式 REST API 定义
  7. 插件化节点: DirectoryLoader 支持动态加载和懒加载

这种架构使得 n8n 既能作为通用工作流引擎运行,又能作为 AI Agent 平台提供智能自动化能力。

评论

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注