客户端、Bulk 与重试:批量写入怎样处理局部失败
一次 Bulk 请求可以同时返回三种结果:第一条新建成功,第二条价格格式错误,第三条要求新建的文档已经存在。服务端用一个 HTTP 响应承载这些结果,应用仍要分别处理每条输入。
写入程序要保存原操作与响应项的对应关系,分别记录已完成、待修正和执行结果未知的操作。
客户端怎样组织连接与批次
连接应随应用复用
Java 应用通常创建一个长生命周期的 OpenSearch 客户端,复用其 transport 和连接池。为每次请求新建客户端,会重复建立 TCP/TLS 连接,也容易遗留线程和连接。客户端可由多个业务线程共享;请求构造器和可变批次列表不要跨线程无保护地修改。
OpenSearch Java Client 将类型化请求转换为 HTTP,再由 transport 执行网络通信。这里使用 opensearch-java 3.9.0 和其依赖的 Apache HttpClient 5 transport,对接 OpenSearch 3.8.0。客户端版本号与服务端版本号独立,依赖入口见 Java client。
完整工程的 Maven 依赖解析结果为 HttpClient 5.6、HttpCore 5.4.2、Jackson 2.21.2。升级时检查实际依赖树和客户端兼容性,不要为了匹配一段旧配置而强行降级某个传递依赖。坐标可对照 3.9.0 发布 POM。
连接配置至少包含:
| 配置对象 | 需要确定的内容 |
|---|---|
| 地址与发现方式 | 直连哪些节点,还是访问负载均衡、托管服务入口 |
| TLS | CA 信任链、目标主机名校验、证书轮换 |
| 身份 | 按索引授权的服务账号或签名认证,凭据从运行配置注入 |
| 连接池 | 总连接数、每路由连接数、空闲连接维护 |
| 等待预算 | 连接建立、池中租借、响应等待、整次业务操作截止时间 |
| 生命周期 | 停止接收新任务,处理未完成批次,关闭 transport |
开启 sniffing 前,应确认返回的节点地址对客户端实际可达。在容器私网、反向代理或托管服务中,发现出的内部地址可能绕开原来的入口和认证路径。
生产 TLS 应验证证书和主机名。使用自签名证书时导入受信 CA,不能把“信任所有证书”或关闭主机名校验作为常规修复方式。本机实验关闭了安全插件,因此示例中的 HTTP 只适用于专用回环服务。
一条批次受多种上限约束
批量器通常在三个条件之一触发时发送:达到条数上限、达到编码后字节上限、最早一条已经等待足够久。并发中的批次数还需要单独限制。
只按条数分批无法控制报文大小。一百份小商品文档与一百份包含长描述、规格和附件元数据的文档,序列化体积可能差几个数量级。字节数应计算最终 UTF-8 NDJSON,而不是 Java 字符串长度。
待写入操作
└─ 校验单项大小
└─ 有界队列
└─ 按条数 / 字节 / 等待时间成批
└─ 有界并发发送
└─ 逐项确认或进入修复、重试队列队列满时,需要暂停上游拉取、返回明确拒绝,或写入可靠的持久缓冲。把队列改成无限容量,只会把服务端的吞吐问题转移到应用堆内存。
同一商品的多个版本可以进入不同批次。若允许并行发送,必须靠按键串行或业务版本防止旧状态覆盖新状态;批次的提交顺序不能替代每条记录的更新顺序。
Bulk 协议与逐项结果
NDJSON 的四种操作
Bulk 使用逐行 JSON。操作行说明要做什么,后续内容行提供文档或更新指令:
| action | 后面是否有内容行 | 行为 |
|---|---|---|
| index | 有,完整文档 | 创建文档,或替换相同 ID 的原文档 |
| create | 有,完整文档 | 仅在 ID 尚不存在时创建;已存在通常返回 409 |
| update | 有,doc、script、upsert 等 | 读取并更新文档;缺失时按指定 upsert 策略处理 |
| delete | 无 | 按 ID 删除;不存在与执行错误要结合 result、status 判断 |
{"index":{"_id":"p1"}}
{"price":100}
{"create":{"_id":"p2"}}
{"price":200}
{"update":{"_id":"p1"}}
{"doc":{"price":120}}
{"delete":{"_id":"p2"}}末尾必须有换行。action 和 source 各占一行,不要对整个报文做美化缩进;delete 后也不要额外补一行空对象。通过 curl 发送文件用 --data-binary,Content-Type 采用 application/x-ndjson。具体格式、ID 长度和 action 参数见 Bulk API。
索引名既可以放在路径 POST /products/_bulk 中,也可以为操作指定 _index。应用应使用允许写入的索引或写 alias,禁止用户随意指定元数据目标。
update 的 doc 是部分更新;script 则可能进行加减、数组操作或条件更新。retry_on_conflict 用于 update 内部遭遇版本冲突后的再尝试,不能自动赋予业务脚本幂等性。更新细节见 Update Document API。
HTTP 层与操作层分别判断
服务端解析 Bulk 后按目标主分片组织写入,各项独立返回结果。批次没有跨文档原子事务:某项失败时,其他项已经成功的修改会保留。
三种情况对应不同处理入口:
| 观察结果 | 已知信息 | 下一步 |
|---|---|---|
| HTTP 400,整个 NDJSON 无法解析 | 请求格式有问题,未获得可靠逐项结果 | 修复报文;若错误发生在不明确的执行阶段,核对实际数据 |
HTTP 200,errors: true | 已返回每项 status、error 或 result | 按 items 与输入序号一一对应处理 |
| 网络中断,没有完整响应 | 客户端不知道哪些项已完成 | 依据幂等契约重放,或读取实际状态核对 |
即使顶层 errors 为 false,也应保留操作结果及所需的版本、目标索引信息。删除不存在的文档等业务可接受情况,要根据具体 action 的语义确认,而非只规定“所有项必须是 201”。
响应 items 按请求操作顺序排列。关联时使用原始序号或持久事件编号;一个批次中可能多次操作同一个 ID,仅按 ID 建 Map 会覆盖关联记录。
用真实混合批次检查局部失败
下载并解压OpenSearch 实验工程,进入 opensearch-lab。环境为 Linux Bash、Docker/Compose、curl、jq;使用有 Docker 权限的普通宿主账号。OpenSearch 单节点限制 2 GiB,安全插件关闭,19222 端口仅绑定回环;不接入生产数据。镜像部署参数见 Docker 安装说明。
set -euo pipefail
command -v docker curl jq
docker compose -p search22 up -d --wait --wait-timeout 180
SEARCH_URL=http://127.0.0.1:19222
INDEX="lab22-bulk-$$"
curl -q --noproxy '*' --fail-with-body --silent --show-error \
"$SEARCH_URL/" | jq -e '.version.number == "3.8.0"'
curl -q --noproxy '*' --fail-with-body --silent --show-error \
-X PUT "$SEARCH_URL/$INDEX" -H 'Content-Type: application/json' \
--data-binary @src/test/resources/bulk/index.json | jq -e '.acknowledged'
curl -q --noproxy '*' --fail-with-body --silent --show-error \
-X PUT "$SEARCH_URL/$INDEX/_doc/existing" \
-H 'Content-Type: application/json' --data '{"price":10}' \
| jq -e '.result == "created"'索引创建失败即停止,不要覆盖已有同名资源。工程内 mixed.ndjson 包含三项:valid 的价格是整数 100,bad 的价格是字符串 oops,最后要求 create 已有文档 existing。
RESULT=$(curl -q --noproxy '*' --fail-with-body --silent --show-error \
-X POST "$SEARCH_URL/$INDEX/_bulk" \
-H 'Content-Type: application/x-ndjson' \
--data-binary @src/test/resources/bulk/mixed.ndjson)
printf '%s\n' "$RESULT" | jq -e \
'.errors == true and ([.items[] | to_entries[0].value.status] == [201,400,409])'
printf '%s\n' "$RESULT" | jq '.items[] | to_entries[0] | {action:.key,result:.value}'这个请求传输成功,HTTP 为 200;失败来自两条操作。bad 对应 mapper_parsing_exception,existing 对应 version_conflict_engine_exception。先修正 bad 的价格,再只发送这一项:
curl -q --noproxy '*' --fail-with-body --silent --show-error \
-X POST "$SEARCH_URL/$INDEX/_bulk" \
-H 'Content-Type: application/x-ndjson' \
--data-binary @src/test/resources/bulk/repaired.ndjson \
| jq -e '.errors == false and (.items | length) == 1'
curl -q --noproxy '*' --fail-with-body --silent --show-error \
"$SEARCH_URL/$INDEX/_doc/existing" | jq -e '._source.price == 10'existing 仍是 10。将 create 改成 index 会改变业务语义并覆盖它;只有调用方确实要求替换时才应这样做。
Java 客户端仍须检查 items
工程中的类型化请求使用以下结构,完整初始化、异常处理和断言在 BulkTest.java:
var response = client.bulk(b -> b.index(name)
.operations(o -> o.index(i -> i.id("valid")
.document(Map.of("price", 100))))
.operations(o -> o.index(i -> i.id("bad")
.document(Map.of("price", "oops"))))
.operations(o -> o.create(i -> i.id("existing")
.document(Map.of("price", 999)))));
for (int position = 0; position < response.items().size(); position++) {
var item = response.items().get(position);
// position 对应原始操作;按 status、error、action 和业务版本决定处理。
}Java 方法正常返回,仅表示得到了可解析的 Bulk 响应。逐项失败不会统一变成一个 Java 异常。真实执行得到 [201, 400, 409],只修复一项后,索引内共有 valid、bad、existing 三份文档。
重试怎样保持写入含义
先分类,再确定是否重发
| 失败类别 | 常见情况 | 处理方式 |
|---|---|---|
| 输入或映射错误 | 400、字段类型错误、未知字段 | 隔离输入,修数据或 schema,修复后重放 |
| 身份或授权错误 | 401、403 | 修凭据、证书或索引权限,停止无效重试 |
| 版本冲突 | 409 | 区分重复 create、旧业务版本、OCC 冲突 |
| 资源拒绝 | 429,伴随具体 error.type | 降低并发、限制积压,在语义安全前提下退避 |
| 服务或连接故障 | 部分 5xx、连接重置、响应丢失 | 按剩余预算与幂等条件重试;结果未知时核对状态 |
HTTP 状态不足以独立决定重试。例如 403 可能来自权限,也可能伴随具体的索引限制错误;要保留 error.type 和脱敏后的 reason,找出真正需要修改的配置。429 同样可能与队列、压力控制等不同机制有关。
退避一般采用增长等待加随机抖动,同时限制尝试次数和总时长。客户端、业务服务和消息消费者如果各自重试三次,实际尝试会相乘;指定一个主要重试层,其他层保留必要的连接恢复即可。超时与重试设计进一步讨论操作预算与重复副作用。
服务恢复时不要一次清空全部重试积压。逐步提高发送并发,观察拒绝、写入延迟和消费滞后;否则排队的旧请求会再次压垮服务。
固定 ID、完整状态与单调版本
确定 ID 能让同一业务对象始终写到同一文档;完整状态覆盖则比“库存加一”更容易重放。还需要版本阻止较晚到达的旧事件:
p1 / sourceVersion=5 / price=100
p1 / sourceVersion=3 / price=80 ← 迟到旧事件使用 version_type=external 时,只有版本严格大于当前版本的写入才会通过。相等版本也返回 409。若它来自一次响应丢失后的重放,应比较持久事件标识,或读取版本和规范化后的内容,确认是同一变更。
external_gte 接受大于或等于的版本。同版本不同内容会被覆盖,因此只有在“同键、同版本对应唯一载荷”被生产者和存储共同保证时才可使用。排序规则见 Index Document API。
业务版本应在该对象的所有写入路径中可比较。两个来源各自从 1 开始计数,或直接把不同机器的当前时间当版本,会引入冲突、回退或误覆盖。需要统一排序来源,或由投影层明确合并策略。
OpenSearch 的 if_seq_no 和 if_primary_term 则用于基于已读版本的乐观并发控制:只有文档仍是读到的状态才更新。它适合读改写,不应直接当作跨索引重建时沿用的业务版本。
响应丢失后,写入可能已经完成
工程的故障代理将请求转发给 OpenSearch,在收到服务端响应后关闭客户端 TCP 连接。服务端已经返回写入结果,客户端却无法读到它。
带 external 版本 5 的 Bulk 实验结果为:
TCP response lost after actual item 201
replay item=409
version/payload verified; documents=1客户端第一次得到传输异常,但直接 GET 已能读到价格 100、版本 5。恢复连接后重放相同操作,版本冲突阻止重复写;比较版本和完整 source 后确认这条事件已经应用。
对增量脚本做相同实验:
increment response lost after count=1
blind replay changed count to 2脚本 ctx._source.count += params.delta 第一次确实执行了。第二次重发又执行一次,固定文档 ID 也不能阻止这个副作用。可以改为带版本的完整计数状态,或设计持久事件去重;只给请求加一个未被服务端使用的 requestId 没有作用。
删除也有版本寿命
使用 external 版本删除后,服务端会暂时保留删除版本,以拒绝随后到达的旧写入。这份信息受 index.gc_deletes 控制,默认保留 60 秒,之后不能再依赖它永久挡住历史事件。删除规则见 Delete Document API。
长期可重放的数据管道可保留带最新版本的软删除文档,例如:
{"version": 8, "deleted": true}查询排除 deleted=true,旧版本 7 写入仍会被版本机制拒绝。软删除文档必须保留到所有可能迟到或重放的旧事件都不再可达,才能考虑清理;清理策略要与上游日志保留、重建和备份恢复一起制定。
实验验证了近期硬删除 v6 拒绝 v5,以及保留的软删除 v8 拒绝 v7。它没有等待默认删除版本保留期结束,长期寿命规则依据官方 API 说明。
失败如何归档、重放与结束
消费位置只能越过已处理的连续前缀
假设同一消息分区的三条事件依次为 10、11、12。10 和 12 写入成功,11 格式错误。此时有两种可选策略:
- 停在 11,修复后继续;12 可以重放,但写入应幂等。
- 将 11 的原事件、错误原因、业务键和版本可靠归档到修复队列,归档成功后将其视为“已转交处理”,再推进消费位置。
只保存在进程内存中的失败列表无法支撑第二种策略。进程崩溃后,消费位点已经前进,坏数据却没有可恢复入口。
修复队列中的敏感字段需要按权限保存或脱敏。重放还应记录原事件身份、修复原因和修复后的版本,避免在不知情时覆盖已有新状态。归档与消费者提交的具体事务方式依赖所用消息系统;Kafka 手动提交的语义可查 KafkaConsumer。
正常停机要等待未完成的批次
收到停止信号后,先停止从上游接收新事件,再发送剩余批次,等待已有请求得到逐项结果。未完成事件保留在可重放的上游或持久队列,最后关闭客户端 transport。
整个停机也要有截止时间。如果来不及确认某批结果,就按“结果未知”保留重放能力,不能为了让进程正常退出而把这批标成成功。refresh 控制的是搜索可见性;写入确认后是否立即可搜,另见 Refresh API。
运行四项真实实验
仍在 opensearch-lab 目录,使用宿主 UID/GID 构建,缓存目录由该账号创建:
mkdir -p .m2
docker run --rm --memory 2g --entrypoint mvn \
--user "$(id -u):$(id -g)" --network search22_default \
-e MAVEN_CONFIG=/tmp/maven -e SEARCH_URL=http://opensearch:9200 \
-v "$PWD":/work -v "$PWD/.m2":/m2 -w /work \
maven:3.9.12-eclipse-temurin-17 \
-B -Duser.home=/tmp -Dmaven.repo.local=/m2 -Dtest=BulkTest clean verify预期四项通过:官方客户端混合响应与单项修复、真实写入后断连与版本核对、增量脚本重复副作用、删除及软删除排序。测试通过前会检查服务版本,所有临时索引由测试独立创建并删除。Java 25 可使用 maven:3.9.12-eclipse-temurin-25 复跑相同命令。
准备阶段无法访问服务会直接令测试失败。响应丢失实验则先确认服务端已写入,再检查客户端传输异常和重放后的 source;这几个结果应一起出现在 target/surefire-reports 中。
结束手动实验时删除自己创建的索引,停止专用 Compose,保留数据卷:
curl -q --noproxy '*' --fail-with-body --silent --show-error \
-X DELETE "$SEARCH_URL/$INDEX" | jq -e '.acknowledged'
docker compose -p search22 down权威资料与规范地址
服务部署与 Java 客户端
- Java 客户端:https://docs.opensearch.org/latest/clients/java/
- Java Client 3.9.0 依赖:https://repo.maven.apache.org/maven2/org/opensearch/client/opensearch-java/3.9.0/opensearch-java-3.9.0.pom
- Docker 安装:https://docs.opensearch.org/latest/install-and-configure/install-opensearch/docker/
批量写入、版本与消费进度
- Bulk 协议:https://docs.opensearch.org/latest/api-reference/document-apis/bulk/
- 部分更新与脚本:https://docs.opensearch.org/latest/api-reference/document-apis/update-document/
- 版本与并发写入:https://docs.opensearch.org/latest/api-reference/document-apis/index-document/
- 删除与版本保留:https://docs.opensearch.org/latest/api-reference/document-apis/delete-document/
- Kafka 消费与提交:https://kafka.apache.org/41/javadoc/org/apache/kafka/clients/consumer/KafkaConsumer.html
- 刷新与可见性:https://docs.opensearch.org/latest/api-reference/index-apis/refresh/
