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
相关产品推荐
相关产品推荐

