Conductor 外部负载存储(External Payload Storage)实战指南:S3、Azure Blob 与 PostgreSQL 配置全解析
【免费下载链接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents项目地址: https://gitcode.com/GitHub_Trending/co/conductor
导读
Conductor 是事件驱动的 agentic 工作流引擎,工作流与任务的输入/输出 payload 默认持久化在其自身的数据存储中。当业务场景产生较大 JSON 负载时,将大对象直接入库会显著增大数据存储压力。本指南基于 Conductor 官方文档 externalpayloadstorage.md,系统讲解 Conductor 的软/硬两种 payload 尺寸屏障(Barrier)机制,并逐一剖析三种外部存储实现——Amazon S3、Azure Blob Storage 与 PostgreSQL——的完整配置项、默认值与底层工作原理。读完本文,你将掌握如何把超阈值负载透明地卸载到外部存储、如何设置拒绝超大 payload 的硬性上限,以及如何针对本地开发环境(如 Azurite)进行联调与排障。
⚠️适用范围提醒:外部 payload 存储目前仅由Java 客户端完整实现;其他语言客户端需要相应改造才能启用该能力(官方欢迎社区贡献)。这一点在后文"限制与注意事项"中还会展开。
为什么需要外部 Payload 存储:Context
Conductor 可以在工作流与任务两级、输入与输出两个方向上对 payload 尺寸施加限制(Barrier)。这些屏障有两重目的:
- 防止把 Conductor 当作数据持久化系统使用——工作流编排引擎的核心职责是协调执行,而不是承担数据仓库的角色;
- 降低数据存储(datastore)的压力——大对象反复写入、读取会拖累任务调度与执行性能。
在默认情况下,Conductor 对 workflow 与 task 的 input/output payload 分别设置了 8 个阈值属性(详见下文"屏障配置参数"),源码中的定义位于 ConductorProperties.java,统一使用 Spring 的DataSize类型(@DataSizeUnit(DataUnit.KILOBYTES)),因此所有阈值均以KB为单位配置。
两种屏障:Soft Barrier 与 Hard Barrier
Conductor 施加两类屏障:
软屏障(Soft Barrier)
软屏障用于缓解数据存储压力。在一些特殊的工作流用例中,payload 的体量确实有必要随工作流执行一起保存;此时 Conductor 会将这些 payload外部化存储到 S3,并在执行期间按需上传/下载。整个过程对用户与 worker 进程完全透明——调用方感知不到大对象被搬到了外部存储。
软屏障的行为:当 payload 尺寸超过软阈值(workflowInputPayloadSizeThreshold等)但未超过硬阈值时,payload 被判定为"值得入库执行",Conductor 自动将其写入外部存储,并在执行时透明地回读。从源码看,这一逻辑的载体是ExternalPayloadStorage接口及其 S3 实现:
- 接口定义在 ExternalPayloadStorage.java,包含
getLocation(Operation, PayloadType, path)、upload(path, payload, payloadSize)、download(path)三个核心方法,以及Operation(READ/WRITE)与PayloadType(WORKFLOW_INPUT / WORKFLOW_OUTPUT / TASK_INPUT / TASK_OUTPUT)两个枚举; - S3 实现在 S3PayloadStorage.java,
getLocation使用 AWS SDK 的S3Presigner生成预签名 PUT/GET URL,upload/download则直接通过S3Client读写对象。
硬屏障(Hard Barrier)
硬屏障用于保护 Conductor 后端免于持久化和处理对工作流执行而言非必要的海量数据。当 payload 超过硬阈值(maxWorkflowInputPayloadSizeThreshold等)时,Conductor 会直接拒绝该 payload,并以适当的错误信息(说明 payload 的具体尺寸)作为reasonForIncompletion终止/失败工作流执行。
源码中对两类阈值的语义注释非常清晰(ConductorProperties.java):
- 超过软阈值的 payload 将被存储在
ExternalPayloadStorage中("beyond which the payload will be stored in ExternalPayloadStorage"); - 超过最大(硬)阈值后,workflow 输入/输出会被拒绝并标记为
FAILED,task 输入/输出会被拒绝并标记为FAILED_WITH_TERMINAL_ERROR("beyond which ... will be rejected and the task will be marked as FAILED_WITH_TERMINAL_ERROR")。
此外,源码还定义了一个文档表格未覆盖的额外硬阈值:conductor.app.maxWorkflowVariablesPayloadSizeThreshold,默认256 KB,用于限制工作流变量(workflow variables)变更的负载,超过后相关 task 变更会被拒绝并标记为FAILED_WITH_TERMINAL_ERROR(ConductorProperties.java)。
屏障配置参数(Barriers Setup)
在JVM 系统属性(JVM system properties)中设置以下属性为期望值:
| Property | 描述 | 默认值(KB) | | -- | -- | -- | | conductor.app.workflowInputPayloadSizeThreshold | 工作流输入 payload 软屏障 | 5120 | | conductor.app.maxWorkflowInputPayloadSizeThreshold | 工作流输入 payload 硬屏障 | 10240 | | conductor.app.workflowOutputPayloadSizeThreshold | 工作流输出 payload 软屏障 | 5120 | | conductor.app.maxWorkflowOutputPayloadSizeThreshold | 工作流输出 payload 硬屏障 | 10240 | | conductor.app.taskInputPayloadSizeThreshold | 任务输入 payload 软屏障 | 3072 | | conductor.app.maxTaskInputPayloadSizeThreshold | 任务输入 payload 硬屏障 | 10240 | | conductor.app.taskOutputPayloadSizeThreshold | 任务输出 payload 软屏障 | 3072 | | conductor.app.maxTaskOutputPayloadSizeThreshold | 任务输出 payload 硬屏障 | 10240 |
这些默认值与 ConductorProperties.java 中DataSize字段的初始值一一对应(如DataSize.ofKilobytes(5120L)、DataSize.ofKilobytes(3072L)、DataSize.ofKilobytes(10240L))。
配置示例(-D形式,适用于启动脚本):
-Dconductor.app.workflowInputPayloadSizeThreshold=2048 \ -Dconductor.app.maxWorkflowInputPayloadSizeThreshold=8192 \ -Dconductor.app.taskOutputPayloadSizeThreshold=1024 \ -Dconductor.app.maxTaskOutputPayloadSizeThreshold=8192💡 实际部署时,阈值通常与业务中最大合理 payload 尺寸挂钩:软阈值设为"可接受入库执行的上限",硬阈值设为"无论如何都不允许入库的绝对上限"。也可以参考 docker/server/config/config.properties 中的写法,将属性放入配置文件随服务启动加载。
外部存储实现一:Amazon S3
Conductor 提供了基于Amazon S3的外部 payload 存储实现,用于外部化存储大负载。
启用方式
在 JVM 系统属性中设置:
conductor.external-payload-storage.type=S3⚠️ 该实现假定 S3 访问已在实例上配置好(即通过 AWS 凭证链、IAM 角色等方式可获得访问凭证)。源码中 S3PayloadStorage.java 的类注释明确写道:"The S3 client assumes that access to S3 is configured on the instance",其凭证机制遵循 AWS SDK for Java v2 的凭证解析顺序。
S3 配置属性
在 JVM 系统属性中设置以下属性:
| Property | 描述 | 默认值 |
|---|---|---|
| conductor.external-payload-storage.s3.bucketName | 存储 payload 的 S3 bucket | 文档表格为空;源码默认值为conductor_payloads |
| conductor.external-payload-storage.s3.signedUrlExpirationDuration | payload 预签名 URL 的过期时间(秒) | 5 |
补充说明:文档表格中 bucketName 默认值留空,但 S3Properties.java 中实际定义了
private String bucketName = "conductor_payloads",生产环境务必显式指定自己的 bucket。另外,该配置类还暴露了第三个未写入文档的属性region,默认值为us-east-1(S3Properties.java),当 bucket 位于其他区域时应同步配置。
对象存储位置与 Key 生成规则
payload 会以上述 bucket 为根,以UUID.json文件名存放在由 payload 类型决定的路径下:
- 工作流输入 →
workflow/input/<UUID>.json - 工作流输出 →
workflow/output/<UUID>.json - 任务输入 →
task/input/<UUID>.json - 任务输出 →
task/output/<UUID>.json
该规则实现在 S3PayloadStorage.java 的getObjectKey(PayloadType)私有方法中:方法按PayloadType拼接workflow/input/、workflow/output/、task/input/、task/output/前缀,再追加IDGenerator.generate()生成的 UUID 与.json后缀。若调用getLocation时显式传入了path,则该path会被直接用作对象 key(S3PayloadStorage.java)。
S3 实现的核心流程(对应ExternalPayloadStorage接口三个方法):
getLocation:按操作类型(WRITE→presignPutObject,READ→presignGetObject)生成带application/jsonContent-Type 的预签名 URL,并返回ExternalStorageLocation(含path与uri);upload:通过S3Client.putObject将 JSON 输入流连同 payload 字节数上传到指定 key;download:通过S3Client.getObject返回对象输入流。
模块相关代码与说明见 awss3-storage 模块(含 S3 文件存储实现 S3FileStorage.java,用于 Conductor 文件存储场景)。
外部存储实现二:Azure Blob Storage
Conductor 也提供了基于Azure Blob Storage的外部 payload 存储实现。
⚠️ 前置条件:需要拥有 Azure Blob Storage 账户的connection string 或 SAS Token。若希望预签名 URL 能够过期,则必须指定 Connection String。SAS Token 需具备
Service=Blob、Resource=Object上的Read与Write权限。
Azure 配置属性
在 JVM 系统属性中设置以下属性:
| Property | 描述 | 默认值 |
|---|---|---|
| workflow.external.payload.storage.azure_blob.connection_string | Azure Blob Storage 连接字符串;对 URL 签名是必需的 | (空) |
| workflow.external.payload.storage.azure_blob.endpoint | Azure Blob Storage 端点;设置了 connection_string 时可省略 | (空) |
| workflow.external.payload.storage.azure_blob.sas_token | Azure Blob Storage SAS Token;需具备Blob服务、Object资源上的Read和Write权限;设置了 connection_string 时可省略 | (空) |
| workflow.external.payload.storage.azure_blob.container_name | 存储 payload 的 Azure Blob 容器 | conductor-payloads |
| workflow.external.payload.storage.azure_blob.signedurlexpirationseconds | payload 预签名 URL 的过期时间(秒) | 5 |
| workflow.external.payload.storage.azure_blob.workflow_input_path | 工作流输入存储路径前缀(随机 UUID 文件名) | workflow/input/ |
| workflow.external.payload.storage.azure_blob.workflow_output_path | 工作流输出存储路径前缀(随机 UUID 文件名) | workflow/output/ |
| workflow.external.payload.storage.azure_blob.task_input_path | 任务输入存储路径前缀(随机 UUID 文件名) | task/input/ |
| workflow.external.payload.storage.azure_blob.task_output_path | 任务输出存储路径前缀(随机 UUID 文件名) | task/output/ |
payload 的存储路径结构与 Amazon S3 实现保持一致(workflow/input/、workflow/output/、task/input/、task/output/+ UUID 文件名)。
🔍文档与源码的属性名前缀差异(重要):上述表格中的
workflow.external.payload.storage.azure_blob.*是文档记录的属性名。而当前仓库源码中,AzureBlobProperties.java 实际使用的是@ConfigurationProperties("conductor.external-payload-storage.azureblob")前缀(对应connectionString、containerName、endpoint、sasToken、signedUrlExpirationDuration、workflowInputPath、workflowOutputPath、taskInputPath、taskOutputPath等 camelCase 字段)。部署时建议以源码中的实际前缀为准(conductor.external-payload-storage.azureblob.*),文档表格可视为历史命名记录。字段默认值(容器conductor-payloads、签名 URL 5 秒、四个路径前缀)在源码中均已确认。
使用 Azurite 进行本地测试
可以使用 Azurite 在本地模拟 Azure Storage,用于开发与测试,无需真实云资源。典型流程:本地启动 Azurite(默认监听 10000/10001 端口)→ 通过connection string或endpoint指向本地模拟器 → 配置上述属性后启动 Conductor,大 payload 即会被写入本地模拟的 Blob 容器。
Azure 排障(Troubleshooting)
Netty
setAvailableProcessors冲突:使用 Elasticsearch persistence 时,可能抛出java.lang.IllegalStateException,原因是 Netty 库调用了两次setAvailableProcessors。解决方式是在属性中设置:es.set.netty.runtime.available.processors=false切换 HTTP 客户端为 OkHttp:若希望使用
okhttp替代默认的 Netty HTTP 客户端,可添加以下依赖(${compatible version}需替换为与 Azure SDK 兼容的版本):com.azure:azure-core-http-okhttp:${compatible version}
Azure 模块的源码与测试可进一步参考 azureblob-storage 模块,其中 AzureBlobPayloadStorage.java 为负载存储实现,AzureBlobPayloadStorageTest.java 提供了对应单元测试。
外部存储实现三:PostgreSQL Storage
Frinx 提供了基于PostgreSQL的外部 payload 存储实现,适合不想引入对象存储、希望复用既有数据库基础设施的场景。
⚠️ 前置条件:需要一个具备全部所需凭证的 PostgreSQL 数据库服务器。
PostgreSQL 配置属性
与 S3 / Azure 不同,PostgreSQL 实现将属性写入application.properties:
| Property | 描述 | 默认值 |
|---|---|---|
| conductor.external-payload-storage.postgres.conductor-url | 用于从 PostgreSQL 拉取 json 配置、供 Conductor server 下载的 URL;本地开发例如{{ server_host }}(源码注释示例:http://localhost:8080) | "" |
| conductor.external-payload-storage.postgres.url | PostgreSQL 数据库连接 URL;连接数据库所必需 | (必填) |
| conductor.external-payload-storage.postgres.username | 连接 PostgreSQL 数据库的用户名;连接数据库所必需 | (必填) |
| conductor.external-payload-storage.postgres.password | 连接 PostgreSQL 数据库的密码;连接数据库所必需 | (必填) |
| conductor.external-payload-storage.postgres.table-name | payload 存储所用的 PostgreSQL schema 与表名 | external.external_payload |
| conductor.external-payload-storage.postgres.max-data-rows | PostgreSQL 数据库中数据行的最大数量;超过该限制后,最旧的数据会被删除 | Long.MAX_VALUE(9223372036854775807L) |
| conductor.external-payload-storage.postgres.max-data-days | 数据最大天数;超过限制后,最旧的数据会被删除 | 0 |
| conductor.external-payload-storage.postgres.max-data-months | 数据最大月数;超过限制后,最旧的数据会被删除 | 0 |
| conductor.external-payload-storage.postgres.max-data-years | 数据最大年数;超过限制后,最旧的数据会被删除 | 1 |
保留期的计算规则:数据库中字段的最大数据年龄为years + months + days。也就是说,默认配置(1 年 + 0 月 + 0 天)意味着超过 1 年的数据会被自动清理;行数与时间两类限制同时生效,任一被突破即触发最旧数据删除。
这些属性在源码 PostgresPayloadProperties.java 中以@ConfigurationProperties("conductor.external-payload-storage.postgres")绑定,字段默认值与文档表格完全一致(tableName = "external.external_payload"、maxDataRows = Long.MAX_VALUE、maxDataDays = 0、maxDataMonths = 0、maxDataYears = 1、conductorUrl = "")。
存储结构与 URI 生成
payload 将以 key(externalPayloadPath)为UUID.json的形式存储在 PostgreSQL 数据库中。可以通过external-postgres-payload-resourceREST 控制器为这些数据生成 URI——该控制器实现在 ExternalPostgresPayloadResource.java,配套测试见 ExternalPostgresPayloadResourceTest.java。
⚠️ 要使生成的 URI 正确工作,必须正确设置
conductor-url属性——它决定了返回给调用方的可访问地址前缀,本地开发时通常指向 Conductor server 自身(如http://localhost:8080)。
PostgreSQL 实现的完整代码与说明见 postgres-external-storage 模块,核心存储类为 PostgresPayloadStorage.java,相关单元测试位于 PostgresPayloadStorageTest.java。
工作机制小结:一次透明的大 Payload 上传/回读
结合 ExternalPayloadStorage.java 接口与各存储实现,可以梳理出大 payload 的完整生命周期:
- 写入侧:调用方提交的 workflow/task 输入超过软阈值时,Conductor 通过
getLocation(WRITE, payloadType, null)获取预签名 URL(S3/Azure)或数据库写入位置(PostgreSQL),随后用upload(path, inputStream, size)把 payload 落盘到外部存储,原 payload 位置被替换为ExternalStorageLocation引用(URI + path); - 执行侧:workflow 执行过程中需要读取该输入时,Conductor 透明地通过
download(path)回读 payload 并交给对应 task;task 产生的大输出同样走"软阈值→外部存储→引用替换"的路径; - 拒绝侧:任何方向的 payload 一旦超过硬阈值,执行立即失败——workflow 侧标记为
FAILED,task 侧标记为FAILED_WITH_TERMINAL_ERROR,错误信息中会携带 payload 尺寸详情; - 访问侧:S3/Azure 通过带过期时间的预签名 URL 提供限时访问(默认 5 秒),PostgreSQL 通过
external-postgres-payload-resourceREST 接口暴露 URI。
限制与注意事项
- 客户端语言支持:外部 payload 存储当前仅由 Java 客户端实现并启用;其他语言的客户端库需要修改后才能支持(例如在调用时主动完成上传/回读逻辑)。官方欢迎社区提交贡献。
- 配置位置差异:屏障阈值与 S3 / Azure 相关属性通过JVM system properties设置;PostgreSQL 相关属性写入application.properties。
- Azure 属性前缀:文档中的
workflow.external.payload.storage.azure_blob.*与源码中的conductor.external-payload-storage.azureblob.*存在差异,部署时应以源码实际绑定的前缀为准(参见上文分析)。 - 凭证前置:S3 依赖实例上的 AWS 访问配置;Azure 需要 connection string 或满足权限要求的 SAS Token;PostgreSQL 需要完整的数据库连接信息。
- 生产建议:显式设置 bucket/容器名与区域,避免依赖源码中的默认值;根据业务中最大合理负载设置软/硬阈值,既避免误杀合法负载,也防止把编排引擎当存储用。
延伸阅读
- 官方文档原文:externalpayloadstorage.md
- 核心配置源码:ConductorProperties.java
- 存储抽象接口:ExternalPayloadStorage.java
- S3 实现:S3PayloadStorage.java 与 S3Properties.java
- Azure 实现:AzureBlobPayloadStorage.java 与 AzureBlobProperties.java
- PostgreSQL 实现:PostgresPayloadStorage.java、PostgresPayloadProperties.java 与 ExternalPostgresPayloadResource.java
- 部署配置示例:docker/server/config/config.properties
【免费下载链接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents项目地址: https://gitcode.com/GitHub_Trending/co/conductor
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考