首页
/ Conductor OpenSearch 2.x 索引后端接入指南:os-persistence-v2 模块配置、迁移与源码原理

Conductor OpenSearch 2.x 索引后端接入指南:os-persistence-v2 模块配置、迁移与源码原理

2026-09-09 20:28:19作者:舒璇辛Bertina

导读

本文聚焦 Conductor 仓库中的 os-persistence-v2 模块(os-persistence-v2/README.md),完整讲解如何将 OpenSearch 2.x(2.0 ~ 2.18+)配置为 Conductor 的工作流与任务索引后端,覆盖全部配置项、单节点开发环境、外部索引管理、从旧版 opensearch 类型的平滑迁移,以及通过 docker-compose-redis-os2.yaml 一键启动的完整栈。与此同时,文章将深入 OpenSearchConfigurationOpenSearchPropertiesOpenSearchRestDAO 等源码实现,解释每个配置参数背后的运行机制、异步批量索引流程与依赖隔离方案,帮助读者既拿到可直接落地的配置,也理解其底层原理。

模块定位与设计背景

Conductor 使用独立的索引后端来支持工作流与任务的全文检索(例如在 UI 中按 workflowIdtaskId、状态或时间范围搜索)。os-persistence-v2 正是面向 OpenSearch 2.x 集群 的索引持久化模块,其核心职责包括:

  • 为 workflow、task、task_log、event、message 五类文档建立索引;
  • 提供 IndexDAO 实现,接收 Conductor 运行时产生的索引写入请求;
  • 提供工作流/任务的搜索(search)能力,支撑 UI 与 API 查询。

从源码结构看,模块按职责划分为三个包(os-persistence-v2/src/main/java/org/conductoross/conductor/os2):

  • config:Spring Boot 自动配置与属性绑定,包括 OpenSearchConfigurationOpenSearchPropertiesOpenSearchConditions
  • dao/index:数据访问层,核心是 OpenSearchRestDAO(实现 IndexDAO)、OpenSearchBaseDAO,以及 BulkRequestWrapper / BulkRequestBuilderWrapper
  • dao/query/parser:一套用于将 Conductor 查询表达式解析为 OpenSearch 查询的解析器(ExpressionGroupedExpressionBooleanOpComparisonOp 等)。

该模块使用 opensearch-java 2.18.0 客户端,并通过对客户端类进行依赖遮蔽(dependency shading)实现命名空间隔离,确保与 os-persistence-v3 模块(面向 OpenSearch 3.x)能在同一 classpath 共存而不冲突。

启用 OpenSearch 2.x 索引

最小配置

在 Conductor 服务器的配置文件(如 docker/server/config/config-redis-os2.properties)中加入以下内容即可启用:

conductor.indexing.enabled=true
conductor.indexing.type=opensearch2

# URL of the OpenSearch cluster (comma-separated for multiple nodes)
conductor.opensearch.url=http://localhost:9200

# Index prefix (default: conductor)
conductor.opensearch.indexPrefix=conductor

其中 conductor.indexing.type=opensearch2 是模块激活的开关。从源码看,条件装配逻辑位于 OpenSearchConditions.javaOpenSearchV2Enabled 使用 AllNestedConditions 同时要求 conductor.indexing.enabled=true(未配置时默认视为 true)且 conductor.indexing.type=opensearch2,两个条件同时满足后 OpenSearchConfiguration 才会被加载。

conductor.indexing.type 的默认值在服务启动入口 Conductor.java 中读取,默认 memory(不启用索引),因此显式设置为 opensearch2 是启用本模块的关键一步。

完整配置项参考

以下配置统一使用 conductor.opensearch.* 命名空间(默认值以 OpenSearchProperties.java 源码为准):

Property Default Description
conductor.opensearch.url localhost:9201 逗号分隔的 OpenSearch 节点 URL 列表,支持 http://https:// 协议
conductor.opensearch.indexPrefix conductor 创建索引时使用的前缀
conductor.opensearch.clusterHealthColor green 启动前等待的集群健康色(greenyellow
conductor.opensearch.indexBatchSize 1 启用异步索引时每批写入的文档数
conductor.opensearch.asyncWorkerQueueSize 100 异步索引任务队列大小
conductor.opensearch.asyncMaxPoolSize 12 异步索引线程池最大线程数
conductor.opensearch.asyncBufferFlushTimeout 10s 异步缓冲在无活动情况下被刷写前的保留时间
conductor.opensearch.indexShardCount 5 每个索引的分片数
conductor.opensearch.indexReplicasCount 0 每个索引的副本数
conductor.opensearch.taskLogResultLimit 10 单次查询返回的最大任务日志条数
conductor.opensearch.restClientConnectionRequestTimeout -1 从连接管理器请求连接的超时时间(毫秒,-1 表示不限制)
conductor.opensearch.autoIndexManagementEnabled true 是否由 Conductor 自动创建和管理索引
conductor.opensearch.username (none) Basic 认证用户名
conductor.opensearch.password (none) Basic 认证密码

此外,源码中还暴露了 conductor.opensearch.version(默认 2,支持值 12)与 conductor.opensearch.documentTypeOverride(默认空字符串)两个参数,前者用于版本相关 API 处理并会在非法版本时抛出 IllegalArgumentException(见 OpenSearchProperties#validateVersion),后者用于在外部管理索引的场景下覆盖文档类型。

各配置项的源码级解读

  • URL 解析与多节点支持OpenSearchProperties#toURLs() 会将 conductor.opensearch.url 按逗号切分,并对未携带协议的主机自动补全 http:// 前缀。OpenSearchConfiguration#osRestClientBuilder 随后通过 convertToHttpHosts 将 URL 列表转换为 HttpHost 数组交给 RestClient.builder(...) 构建,因此多节点写法为:
conductor.opensearch.url=http://os-node1:9200,http://os-node2:9200,http://os-node3:9200
  • Basic 认证osRestClientBuilder 中,当用户名与密码均非空时会构造 BasicCredentialsProvider 并挂到 HTTP 客户端回调上,日志会打印 "Configure OpenSearch with BASIC authentication";否则打印 "Configure OpenSearch with no authentication."(OpenSearchConfiguration.java)。

  • 连接请求超时:仅当 restClientConnectionRequestTimeout > 0 时才会设置 setConnectionRequestTimeout,默认 -1 表示不施加限制。

  • 启动健康检查OpenSearchRestDAO#setup()@PostConstruct 阶段首先调用 waitForHealthyCluster(),通过 GET /_cluster/health 并携带 timeout=30swait_for_status=<clusterHealthColor> 等待集群达到期望健康状态,之后才进行索引初始化(OpenSearchRestDAO.java)。

  • 重试策略OpenSearchConfiguration#osRetryTemplate 提供 RetryTemplate 实例,采用 FixedBackOffPolicy 固定 1000ms 退避,供索引写入(如 indexWithRetry)失败时重试使用。

基础认证配置

连接带安全认证的 OpenSearch 集群:

conductor.opensearch.username=myuser
conductor.opensearch.password=mypassword

对应的底层实现即上文提到的 BasicCredentialsProvider 注入逻辑。需要注意的是,密码应避免硬编码在提交到版本库的配置文件中,生产环境建议通过环境变量或密钥管理服务注入。

单节点 / 开发集群配置

单节点 OpenSearch 集群无法达到 green 健康状态——副本分片没有第二个节点可分配,因此需要将健康色放宽为 yellow,同时将副本数设为 0:

conductor.opensearch.clusterHealthColor=yellow
conductor.opensearch.indexReplicasCount=0

这也是 docker/server/config/config-redis-os2.propertiesindexReplicasCount=0clusterHealthColor=green 组合背后的考量:官方 compose 中 OpenSearch 只有单节点(docker-compose-redis-os2.yaml),但该配置组合仍可正常运行。

外部索引管理

如果索引的创建与管理由外部完成(例如通过 ILM 生命周期策略或 Terraform 管理索引模板、别名与滚动策略),可以关闭自动索引管理:

conductor.opensearch.autoIndexManagementEnabled=false

源码层面,OpenSearchRestDAO#setup() 会先执行 waitForHealthyCluster(),然后判断 properties.isAutoIndexManagementEnabled():仅在开启时才调用 createIndexesTemplates()createWorkflowIndex()createTaskIndex()OpenSearchRestDAO.java)。也就是说,即便关闭自动管理,健康检查依然生效,只是不再创建索引与模板。

自动管理模式下,Conductor 会:

  • 通过 initIndexTemplatetask_logeventmessage 三类文档创建 template_task_logtemplate_eventtemplate_message 索引模板(PUT /_index_template/<template>),模板资源位于 os-persistence-v2/src/main/resources
  • 通过 addIndex(index, mapping) 创建 conductor_workflow(映射 mappings_docType_workflow.json)与 conductor_task(映射 mappings_docType_task.json)两个固定索引,分片数与副本数取自 indexShardCount / indexReplicasCount
  • task_logeventmessage 使用 yyyyMMww(ISO 周)粒度滚动索引,即 conductor_task_log_<yyyymmww> 等命名,并每小时(scheduleAtFixedRate(..., 0, 1, TimeUnit.HOURS))更新一次当前索引名。

template_task_log.json 为例,模板声明 index_patterns: ["*task*log*"]refresh_interval: 1s,并定义 createdTime(long)、log(text)、taskId(keyword)字段映射,OpenSearchBaseDAO#applyIndexPrefixToTemplate 会在加载模板资源时自动将 *task*log* 模式改写为带前缀的形式(如 conductor_task_log_*),这也是 indexPrefix 能同时影响模板与索引名的原因。

从旧版 opensearch 索引类型迁移

若此前使用的是 conductor.indexing.type=opensearch(旧命名空间 conductor.elasticsearch.*),迁移只需两步:

# Before
conductor.indexing.type=opensearch
conductor.elasticsearch.url=http://localhost:9200

# After
conductor.indexing.type=opensearch2
conductor.opensearch.url=http://localhost:9200

conductor.elasticsearch.* 命名空间出于向后兼容仍然被接受,但已被标记为废弃,检测到旧属性时启动阶段会打印弃用警告日志。其兼容逻辑在 OpenSearchProperties#init() 中实现:遍历 urlversionindexName(映射到 indexPrefix)、clusterHealthColorindexBatchSizeasyncWorkerQueueSizeasyncMaxPoolSizeindexShardCountindexReplicasCounttaskLogResultLimitrestClientConnectionRequestTimeoutautoIndexManagementEnableddocumentTypeOverrideusernamepassword 等旧属性,只要对应新属性未配置就回退读取旧值;当新旧命名空间同时存在时,conductor.opensearch.* 优先(hasNewProperty 判定)。同时需注意特殊约定:conductor.elasticsearch.version=0 是旧 ES7 自动配置的禁用标记,不会被采纳为 OpenSearch 版本号。

这些兼容行为都有单元测试覆盖,例如 OpenSearchPropertiesTest.java 中的 testLegacyUrlPropertyFallbacktestLegacyVersionPropertyFallbacktestVersionZeroIsIgnoredAsES7DisableFlagtestAllLegacyPropertiesFallback 等用例。

Docker Compose 一键启动

仓库提供了完整的 OpenSearch 2.x 演示栈:

docker compose -f docker/docker-compose-redis-os2.yaml up

该编排会启动三个服务(docker-compose-redis-os2.yaml):

  • conductor-server:以 config-redis-os2.properties 为配置(其中已启用 conductor.indexing.type=opensearch2,索引指向 http://os:9200),暴露 8000→8080(API)与 8127→5000(调试端口);
  • conductor-redisredis:6.2.3-alpine,承担 DB 与队列(配置中 conductor.db.type=redis_standaloneconductor.queue.type=redis_standalone);
  • conductor-opensearchopensearchproject/opensearch:2.18.0,单节点(plugins.security.disabled=trueOPENSEARCH_INITIAL_ADMIN_PASSWORD 已设置),持久化到 osdata2-conductor 卷,并配置了集群健康检查。

启动后 Conductor 会自动创建 conductor_workflowconductor_task 索引及日志/事件/消息模板,工作流运行时产生的数据将异步写入索引,可通过 Conductor UI 或 REST API 检索。

异步批量索引机制

os-persistence-v2 的写入采用"异步 + 批量"设计,源码集中于 OpenSearchRestDAO

  • 双线程池:主线程池 executorService 核心 6 线程、最大 asyncMaxPoolSize(默认 12)、队列容量 asyncWorkerQueueSize(默认 100),负责 workflow/task 文档索引;logExecutorService 核心 1、最大 2 线程,负责 task_log / event / message 的索引(OpenSearchRestDAO.java)。队列满时采用拒绝策略并记录 recordDiscardedIndexingCount 指标。

  • 批量聚合indexDocument 将文档以 IndexRequest 形式追加到按 docType 分桶的 BulkRequest 缓冲中,一旦 numberOfActions >= indexBatchSize 就触发 indexBulkRequest 批量提交(OpenSearchRestDAO.java)。这也是 conductor.opensearch.indexBatchSize 直接控制吞吐与延迟的原因。

  • 定时兜底刷写:构造器中注册了 scheduleAtFixedRate(this::flushBulkRequests, 60, 30, TimeUnit.SECONDS) 的定时任务,flushBulkRequests 会检查每个 docType 缓冲距离上次刷写是否超过 asyncBufferFlushTimeout(默认 10s),超时且缓冲非空则强制批量提交,防止实例终止时缓冲中的文档丢失(OpenSearchRestDAO.java)。

  • 失败重试与监控indexWithRetry 通过 retryTemplate(固定 1s 退避)执行 openSearchClient.bulk(...),并记录 recordESIndexTime 索引耗时指标与两个队列深度指标。

  • 优雅关闭@PreDestroy shutdown() 对两个线程池先 shutdown() 再等待最多 30 秒,超时则 shutdownNow(),保证存量缓冲尽可能刷完。

对应的测试用例包括 TestOpenSearchRestDAOBatch.java(以 @TestPropertySource(properties = "conductor.elasticsearch.indexBatchSize=2") 验证批量行为)与 TestBulkRequestBuilderWrapper.java,可用于理解批量索引的边界行为。

依赖隔离:v2 与 v3 共存方案

OpenSearch 2.x 与 3.x 的 Java 客户端使用完全相同的包名(org.opensearch.client.*),如果两个模块同时进入服务器 classpath,会产生类冲突。为此,os-persistence-v2 通过 Shadow 插件(Shadow plugin)将所有 OpenSearch 客户端类搬迁到隔离命名空间:

org.opensearch.client → org.conductoross.conductor.os2.shaded.opensearch.client

这样 os-persistence-v2os-persistence-v3 可以共存于同一进程,由 conductor.indexing.type 决定启用哪一个(opensearch2opensearch3)。这种隔离策略允许 Conductor 在迁移 OpenSearch 大版本期间平滑过渡,无需停机切换索引后端。

与 OpenSearch 3.x 模块及配置文档的关系

  • 若集群为 OpenSearch 3.x,应使用 os-persistence-v3 模块(conductor.indexing.type=opensearch3),两模块共享 conductor.opensearch.* 命名空间,仅 conductor.indexing.type 取值不同;
  • 更完整的 OpenSearch 集成说明可参考 docs/documentation/advanced/opensearch.md,其中给出了 v2/v3 的版本支持矩阵(v2 模块对应 2.0 – 2.18+,v3 模块对应 3.0+)与更多配置细节。

常见问题与排查建议

  • 启动时集群健康等待超时:确认 conductor.opensearch.url 可达且集群状态至少为配置的 clusterHealthColor;单节点环境请配置为 yellow 并设置 indexReplicasCount=0
  • 索引未自动创建:检查 conductor.indexing.enabled=trueconductor.indexing.type=opensearch2 是否同时满足;若关闭了 autoIndexManagementEnabled,请确认外部已按 conductor_workflowconductor_taskconductor_task_log_* / conductor_event_* / conductor_message_* 约定创建索引与模板。
  • 写入延迟或队列丢弃:观察 asyncWorkerQueueSizeasyncMaxPoolSizeindexBatchSize 三个参数是否与写入压力匹配;被丢弃的请求会通过 recordDiscardedIndexingCount("indexQueue") / recordDiscardedIndexingCount("logQueue") 指标暴露。
  • 同时部署 v2 与 v3 出现类冲突:确认两个模块均使用带遮蔽(shaded)的客户端构建,且 conductor.indexing.type 明确指向唯一后端。
  • 迁移旧配置后行为异常:优先检查启动日志中的 DEPRECATION WARNING 是否列出旧属性,并核对 conductor.elasticsearch.version 特殊值 0 未被误用作 OpenSearch 版本。

小结

os-persistence-v2 是 Conductor 面向 OpenSearch 2.x 的正式索引后端模块:一条 conductor.indexing.type=opensearch2 即可接入,十余个 conductor.opensearch.* 参数覆盖连接、认证、健康检查、分片副本、异步批量与索引生命周期管理;conductor.elasticsearch.* 兼容层保证了旧配置平滑迁移;Shadow 依赖隔离则为 v2/v3 双版本共存铺平道路。理解 OpenSearchRestDAO 中"双线程池 + 批量缓冲 + 定时兜底刷写"的写入流水线,有助于在生产环境中针对吞吐与延迟正确调优 indexBatchSizeasyncWorkerQueueSizeasyncBufferFlushTimeout 等关键参数。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.76 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
860
1.35 K
docsdocs
暂无描述
Markdown
899
5.83 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
925
1.85 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.84 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
533
601
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.03 K
525
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.37 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
395