You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Aerospike批量写入记录避免数据不一致及丢失的实现方案

Aerospike 高可靠批量写入函数实现方案

前置集群配置要求

在编写函数前先确保集群配置满足可靠性基础要求:

  • Namespace 配置 replication-factor 至少为2,避免单节点故障丢失数据
  • 开启持久化存储:storage-engine 配置为ssd/pmem,禁用内存-only存储
  • 可选开启 strong-consistency 模式,满足强一致场景需求

函数核心设计思路

1. 批量分片逻辑

1000万条记录不可一次性提交,需按单批100~500条(单批总大小不超过1MB)拆分批次,避免触发服务端单事务大小限制、客户端/服务端超时问题,大幅降低丢数据概率。

2. 写入策略配置

所有写入操作必须指定以下策略,从服务端层面保证可靠性:

  • commit_level 设为 COMMIT_ALL:确保所有副本节点都写入成功才返回客户端成功,避免副本未同步时节点宕机丢数据
  • durable_write 设为True:强制写入数据落盘到持久化存储,不会因为节点掉电丢失数据
  • 主键存在性策略:新增记录用 CREATE_ONLY,仅主键不存在时才写入;更新记录用 UPDATE_ONLY,仅主键存在时才写入
  • generation 校验:更新记录时必须携带从服务端读取到的当前记录版本号,服务端仅版本号匹配时才执行更新,更新成功后版本号自动+1,避免并发覆盖

3. 错误重试与兜底机制

每批写入后必须逐行解析返回结果,仅重试写入失败的记录,不要整批重试:

  • 重试采用指数退避策略,最大重试次数建议设为3~5次,避免瞬间压垮集群
  • 超过最大重试次数仍失败的记录,抛出异常并返回失败列表,可对接本地持久化重试队列做兜底,不会直接丢失数据

4. 并发调用安全保证

  • 直接使用官方提供的线程安全客户端实例,多线程并发调用无需额外加锁
  • 新增同主键记录时,CREATE_ONLY 策略仅会让第一个请求写入成功,其余请求返回主键已存在错误,不会出现覆盖
  • 更新同主键记录时,generation 校验机制会保证只有拿到最新版本号的请求写入成功,其余请求返回版本不匹配错误,不会出现脏写。如果业务允许冲突合并,可在收到版本不匹配错误时,读取最新记录合并变更后重新携带新版本号重试即可。

Python 版本代码示例

import aerospike
from aerospike.exception import AerospikeError
import time

# 初始化全局线程安全客户端
config = {
    "hosts": [("你的Aerospike节点IP", 3000)],
    "policies": {
        "batch": {
            "commit_level": aerospike.COMMIT_ALL,
            "durable_write": True
        }
    }
}
client = aerospike.client(config).connect()

# 替换为你自己的命名空间和集合名
NAMESPACE = "business_ns"
SET = "user_data"

def reliable_batch_write(records, max_retry=3, base_backoff=1):
    """
    高可靠批量写入函数
    :param records: 待写入记录列表,每条格式为 (主键元组, 字段字典, 版本号)
                    新增记录版本号传0,更新记录传从服务端读取到的generation值
                    主键元组格式为 (NAMESPACE, SET, 业务主键ID)
    :param max_retry: 最大重试次数
    :param base_backoff: 指数退避基础时长(秒)
    :raises Exception: 超过重试次数仍有写入失败时抛出异常,附带失败记录列表
    """
    # 拆分批次,单批最大200条
    batch_size = 200
    batches = [records[i:i+batch_size] for i in range(0, len(records), batch_size)]
    
    for batch in batches:
        retry_cnt = 0
        failed = batch.copy()
        
        while retry_cnt < max_retry and failed:
            ops = []
            for key, bins, gen in failed:
                policy = {}
                if gen == 0:
                    # 新增记录,仅主键不存在时写入
                    policy["exists"] = aerospike.EXISTS_CREATE
                else:
                    # 更新记录,仅版本匹配时写入
                    policy["generation"] = gen
                    policy["exists"] = aerospike.EXISTS_UPDATE
                
                # 构造批量写入操作
                write_ops = [
                    aerospike.Operation(aerospike.OP_WRITE, k, v) 
                    for k, v in bins.items()
                ]
                ops.append(aerospike.BatchWrite(key, write_ops, policy=policy))
            
            # 执行批量写入
            try:
                batch_res = client.batch_write(ops)
            except AerospikeError:
                # 集群级错误,整批延迟后重试
                time.sleep(base_backoff ** retry_cnt)
                retry_cnt += 1
                continue
            
            # 筛选写入失败的记录
            new_failed = []
            for idx, record_res in enumerate(batch_res.batch_records):
                if record_res.result != 0:
                    new_failed.append(failed[idx])
            failed = new_failed
            retry_cnt += 1
            if failed:
                time.sleep(base_backoff ** retry_cnt)
        
        if failed:
            raise Exception(f"批量写入最终失败,共{len(failed)}条记录未写入", failed)

# 1000万条记录写入调用示例
if __name__ == "__main__":
    total_records = 10000000
    records = []
    for i in range(total_records):
        # 构造主键和字段,版本号传0表示新增
        key = (NAMESPACE, SET, f"user_{i}")
        bins = {"name": f"用户{i}", "age": 20 + i % 30, "status": 1}
        records.append((key, bins, 0))
    
    # 调用批量写入函数
    try:
        reliable_batch_write(records)
        print("1000万条记录写入完成")
        # 可选:写入完成后批量读取校验完整性
    except Exception as e:
        print(e)
        # 兜底处理失败记录,比如写入本地文件后续重试

效果验证说明

  • 1000万条记录无丢失:分片写入+逐行重试+兜底异常机制,只要返回写入成功的记录都不会丢失,写入完成后可通过批量读取接口校验全量记录的存在性和内容哈希,确保零丢失
  • 写入一致性:单条记录写入是原子性的,如需跨记录的事务一致性,可将需要保证一致的记录放在同一个Aerospike多记录事务(MRT)中提交,实现要么全部成功要么全部失败
  • 并发安全:CREATE_ONLY + generation 校验机制从服务端层面避免了并发写入的覆盖问题,多线程并发调用函数时不会出现数据丢失或覆盖

内容的提问来源于stack exchange,提问作者kkk

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 18:36:03