批量导入
概述
批量导入用于把一批已经转换成图模型的数据,持续、可观察地写入指定 Graph Space。它处理的对象不是原始文件或数据库表,而是应用准备好的 GraphNode 和 GraphEdge 数据流。
与单次节点与边写入相比,导入更关注:
- 数据量超过一次调用适合承载的范围;
- 节点和边需要按正确顺序分批提交;
- 某批失败后如何停止、继续或恢复;
- 如何报告进度、错误和已确认偏移;
- 如何控制内存、连接和外部数据源生命周期。
Graph SDK 提供同步导入和异步导入两条路径。同步导入适合调用方自己控制线程和任务状态;异步导入提供进程内任务、进度快照、取消和 checkpoint。两者都不是文件解析器,也不是完整的分布式 ETL 平台。
为什么不能把所有数据一次写入
小规模写入可以把节点和边放进一个 GraphMutation,但随着数据量增长,一次性处理会带来明显风险:
- Java 进程需要同时持有过多对象;
- 数据库事务过大,锁持有时间和失败成本增加;
- 网络或数据库短暂故障会让整批重新开始;
- 无法持续展示进度,也难以定位失败范围;
- 节点尚未写入时,依赖它的边可能无法创建;
- 任务取消、限流和恢复缺少明确边界。
批量导入把长任务拆成多个可管理的提交单元。代价是:整个导入通常不具备一个全局事务,调用方必须接受并管理“前面的批次已经成功,后面的批次失败”的状态。
典型导入场景
首次全量装载
新建知识库或业务图谱后,把历史实体和关系首次写入空 Space。常见流程是先应用 Schema,再导入所有节点,最后导入所有边并做数量和引用完整性校验。
持续增量导入
同一个知识库可以在几个月或几年内持续导入新文件。后续导入通常继续使用原 Space,通过稳定节点 ID 和边身份更新已有数据、增加新数据。
这里有一个重要边界:通用批量导入只执行 Upsert,不会因为“新批次中没有某条旧关系”就自动删除旧数据。若要处理文档更新、撤回、事实来源和多文档共同引用,应使用知识抽取模块中的文档状态、事实溯源和变更计划,而不是把每次文件导入当作完整图快照覆盖。
重导与故障恢复
当导入进程中断时,可以重新提交同一份确定性数据源。稳定 ID 使已成功数据再次 Upsert 时收敛到相同实体;恢复点可以减少重复扫描和写入。
恢复并不意味着数据库自动回滚到中断前,也不保证失败批次完全没有写入。恢复设计仍然依赖幂等身份、稳定顺序和写入后的数据核验。
重建新版本图谱
当 Schema 或归一规则发生不兼容变化时,直接在原 Space 中重导可能留下旧结构。更稳妥的方式通常是创建新 Space,应用新 Schema,完成全量导入和验证,再由上层切换路由。旧 Space 按保留策略清理。
导入前的准备
1. 固定目标和版本
一次导入任务至少应记录:
- 目标 Space;
- 数据集或知识库 ID;
- Schema 版本;
- 数据源版本或快照时间;
- 导入批次 ID;
- 实体归一规则版本;
- 创建人、创建时间和来源位置。
这些信息属于上层任务元数据。Graph SDK 不会替应用建立导入批次表,但可以通过任务 ID、执行上下文和日志进行关联。
2. 把原始数据转换为图数据
Graph SDK 不直接解析 CSV、Excel、PDF、消息队列或业务数据库。调用方需要先完成:
原始数据
-> 解析与清洗
-> 字段映射
-> 实体归一和稳定 ID
-> Schema 校验
-> GraphNode / GraphEdge
-> Graph SDK 导入文档知识图谱可以先通过 extractor 按 Schema 提取实体和关系;结构化数据则可以通过自定义映射器转换。无论来源是什么,进入导入器前都应保证类型、属性名、边方向和端点身份已经确定。
3. 保证可重复的身份和顺序
可靠重试需要:
- 同一实体每次生成相同节点 ID;
- 同一关系每次生成相同的 source、type、target 和 rank;
- 恢复时数据源顺序与生成 checkpoint 时一致;
- 数据源版本固定,不在恢复期间插入、删除或重新排序;
- 节点和边的迭代过程不会静默跳过异常记录。
仅仅“再次读取同一个文件名”并不能证明内容和顺序相同。生产系统应保存内容哈希、对象版本号或数据库快照标识。
节点优先,边随后
导入器始终先完整消费节点数据流,再消费边数据流:
节点批次 0 -> 节点批次 1 -> ... -> 边批次 N -> 边批次 N+1这种顺序保证创建关系前,端点有机会先被写入。它也带来两个要求:
- 不要把节点和边混在一个只有单次遍历能力的共享游标中;
- 如果边引用的数据不在本次节点流中,应确保它们已存在于目标 Space。
节点阶段存在失败而 stopOnError=false 时,导入器仍会继续处理后续节点和边。此时依赖失败节点的边可能无法创建,所以“继续执行”只适合能够接受部分成功并有后续对账修复的任务。
同步导入
同步导入在当前调用线程中完成,适合由应用自己的作业框架管理生命周期:
GraphImportRequest request = GraphImportRequest.builder()
.nodes(nodes)
.edges(edges)
.batchSize(500)
.stopOnError(true)
.build();
GraphImportReport report = graph.writer().importData(
request,
GraphOptions.ofSpace("company_knowledge"));调用返回时,所有可执行批次已经处理完成,或者在遇错即停模式下抛出异常。
同步路径的能力边界
当前同步 GraphWriter.importData(...) 使用以下请求内容:
| 配置 | 同步导入中的作用 |
|---|---|
nodes / edges | 提供节点和边数据 |
source | Builder 从中取得节点和边 Iterable |
batchSize | 控制每次组成 Mutation 的记录数 |
stopOnError | 控制失败后立即抛出还是继续 |
listener、checkpoint 和 resumeFrom 由异步任务服务处理,直接调用同步 importData 时不会获得这些生命周期能力。同步使用 GraphImportSource 时,调用方还必须自行关闭数据源,推荐使用 try-with-resources:
try (GraphImportSource source = openSource()) {
GraphImportRequest request = GraphImportRequest.builder()
.source(source)
.batchSize(500)
.stopOnError(true)
.build();
GraphImportReport report = graph.writer().importData(request, options);
}批次大小如何选择
batchSize 默认是 500,必须为正数。它控制多少节点或多少边组成一次 GraphMutation,从而影响提交边界、报告频率和失败后的重做范围。
批次并不是越大越好:
| 较小批次 | 较大批次 |
|---|---|
| 单批失败重做范围小 | 调用和提交次数较少 |
| 更快产生进度 | 可能获得更高吞吐 |
| 数据库事务和内存压力较小 | 单批内存、锁和超时风险更高 |
| 调用开销相对更高 | 失败定位和恢复粒度更粗 |
当前 Portable Writer 会在一个 Mutation 内逐个执行节点或边语句,不能假设 batchSize=500 就一定转换成后端的一条原生批量语句。它首先定义逻辑批次和提交边界;真实吞吐仍应在目标数据库、真实网络和真实 Schema 上压测。
建议从默认值开始,分别测量节点和边的吞吐、P95/P99 延迟、数据库 CPU、内存、锁等待和失败重试时间,再决定生产值。不要只根据开发机上的空数据库结果调参。
失败策略和部分成功
stopOnError=true
这是默认策略。第一批失败后,导入停止并抛出异常。它适合数据一致性要求较高、希望尽快阻止错误扩散的任务。
但“停止”不等于“回滚”:此前已经提交的批次仍然存在。调用方应保存任务上下文,通过稳定 ID 重试或执行补偿。
stopOnError=false
失败批次被记录到报告中,后续批次继续处理。适合离线清洗、允许少量坏数据进入隔离队列、并且具备完整对账能力的场景。
使用这一模式时,应明确:
- 返回报告的
isSuccess()可能为false; - 成功批次不会因其他批次失败而回滚;
- 边可能因为端点批次失败而继续失败;
- 应保存失败记录的原始位置和输入内容,不能只记录一条错误消息;
- 修复后需要重新提交失败范围,并再次做完整性检查。
后端内部也可能部分成功
Neo4j 普通 Writer 会为每个 Mutation 使用事务,因此该批次内的多项写入共同提交。Nebula SessionPool Writer 会依次发送语句,一个 Mutation 在中途失败时,前面语句可能已经生效。
所以即使 checkpoint 没有推进,失败批次也可能包含已写入数据。稳定 ID 和 Upsert 能帮助重试收敛,但删除、属性覆盖和非幂等原生操作仍需要额外核验。
导入报告
GraphImportReport 汇总当前调用或任务已经观察到的结果:
| 字段 | 含义 |
|---|---|
nodesImported | 成功批次报告的节点逻辑处理数 |
edgesImported | 成功批次报告的边逻辑处理数 |
batchesCompleted | 成功批次数 |
batchesFailed | 失败批次数 |
batchesAttempted | 成功和失败批次总数 |
errors | 失败消息列表 |
errorDetails | 包含批次序号、错误码和消息的结构化详情 |
这些数字适合进度展示和运行诊断,但不等于数据库中的最终物理数量。重复 Upsert、删除后重建、缺失端点以及后端部分成功都可能造成差异。导入完成后必须通过查询和数据质量规则对账。
同步 Writer 当前不能为结构化错误提供可靠批次序号,可能使用 -1;异步任务按执行顺序维护批次序号,更适合任务级展示。
异步任务、进度和取消
需要立即返回任务 ID、后台执行、轮询状态、取消和 checkpoint 时,使用 graph.imports():
GraphImportTask task = graph.imports().submit(
request,
GraphOptions.ofSpace("company_knowledge"));
GraphImportTask latest = graph.imports().get(task.getId());异步服务逐批调用 GraphWriter.mutate(...),并维护:
QUEUED、RUNNING、SUCCEEDED、FAILED、CANCELLED状态;- 已成功确认的节点和边偏移;
- 已处理批次序号;
- 累计导入报告;
- 监听器、取消信号和数据源关闭。
取消是尽力而为的协作行为,不是数据库回滚。取消请求到达前已经成功的批次仍然保留,正在执行的数据库语句也未必能立即停止。
默认实现使用进程内执行器和内存任务存储,不等价于分布式作业系统。完整的接口、状态和持久化边界见异步导入任务。
Checkpoint 与恢复
异步任务会在一个批次完整成功后推进节点或边偏移,并调用 GraphImportCheckpoint:
GraphImportRequest request = GraphImportRequest.builder()
.source(source)
.batchSize(500)
.checkpoint((taskId, nodesProcessed, edgesProcessed) ->
checkpointRepository.save(
taskId, nodesProcessed, edgesProcessed))
.build();任务失败或取消后,可以从最近持久化的确认位置构造恢复点:
GraphImportResumePoint resumePoint = new GraphImportResumePoint(
savedNodesProcessed,
savedEdgesProcessed,
savedBatchesProcessed);
GraphImportRequest resumed = GraphImportRequest.builder()
.source(reopenSameVersionOfSource())
.batchSize(500)
.resumeFrom(resumePoint)
.checkpoint(checkpoint)
.build();恢复点记录的是“已经确认成功的数量”,不是数据库游标,也不会替 SDK 重新打开外部文件。异步服务会从重新创建的数据流开头跳过相应数量,因此恢复必须满足:
- 数据源内容和顺序与原任务完全一致;
- 节点流和边流可以从头重新迭代;
- checkpoint 已持久化,而不是只保存在任务进程内;
- 恢复到同一个 Space 和兼容 Schema;
- 相同批次大小不是硬性要求,但变更批次策略会影响诊断和对账;
- 重试数据使用稳定节点 ID 和边身份。
单一数量偏移只能准确表达“从开头连续成功到某个位置”。因此,基于 resumeFrom 的恢复应优先配合 stopOnError=true 使用。若 stopOnError=false 时某一批失败、后续批次又成功,成功数量中会出现空洞,单一 offset 无法指出具体失败范围;这类任务应额外持久化失败批次和原始记录清单,修复后按清单重放,而不能只依赖任务快照中的恢复点。
checkpoint 回调异常会被异步服务隔离,不会让数据库批次失败。因此,生产实现需要自行监控 checkpoint 持久化失败;否则数据库已写入而恢复点未推进,恢复时会重复处理已成功批次。
GraphImportSource 的资源生命周期
GraphImportSource 可以包装文件句柄、数据库游标、对象存储流或其他外部资源,并提供节点流、边流和 close()。
在异步服务中,任务成功、失败、取消或服务关闭时,SDK 会尝试幂等关闭 source。自定义实现本身也应允许多次 close(),并避免节点流和边流在并发线程中共享不安全游标。
在同步 writer().importData(...) 中,source 不会由 Writer 自动关闭,调用方必须管理资源。不要在提交异步任务后提前关闭 source,也不要让请求结束钩子关闭仍在后台使用的数据流。
幂等、重导和撤回
重复导入
稳定节点 ID 和 GraphEdgeKey 可以让同一数据重复 Upsert,而不是无限新增实体和关系。但以下变化仍会产生新数据或留下旧数据:
- 实体归一规则变化导致节点 ID 改变;
- 边方向、类型或 rank 改变;
- 标签变化导致 Neo4j 合并身份变化;
- 新版本不再包含某条旧边,但没有显式删除;
- 属性从数据源消失,但 Upsert 只提交仍存在的字段。
因此,重导前要明确它是“追加与更新”“完整快照同步”还是“重新构建”。通用导入直接支持的是前者。
文件更新与撤回
文档知识图谱不能简单地按节点 ID 删除某个文件曾提到的所有实体,因为同一实体和关系可能同时由多个文件支持。推荐保存文档版本、事实 ID、来源证据和引用关系:
文档版本 -> 提取事实 -> 节点/边声明 -> 当前图状态文件更新时先比较新旧事实,只撤销不再被任何有效来源支持的声明;文件撤回时同样按来源清理。相关流程见文档生命周期。
导入完成后的数据质量验证
“任务成功”只表示所有 Writer 调用没有报告失败,不表示图谱业务上正确。至少应验证:
- 节点和边数量是否在预期范围;
- 是否存在引用不到端点的关系或遗漏关系;
- 必填属性、数据类型和枚举是否符合 Schema;
- 是否出现异常重复节点、别名未归一或错误主标签;
- 关键边方向和 rank 是否正确;
- 每个数据源或文档的来源记录是否完整;
- 关键查询是否能返回预期结果;
- 索引是否就绪,查询计划和延迟是否符合要求;
- 报告数量、源数据数量和数据库查询数量能否合理解释。
对大型导入,建议先在隔离 Space 做抽样和全量校验,达到发布阈值后再切换上层路由。
Neo4j 与 Nebula 的差异
| 维度 | Neo4j | Nebula Graph |
|---|---|---|
| 单个逻辑批次 | 普通 Writer 在数据库事务内执行 | 多条 nGQL 顺序执行 |
| 批次失败 | 当前批次不提交,之前批次保留 | 当前批次内部也可能部分生效 |
| 节点类型 | 支持多 Label | Portable Writer 当前只支持单 Tag |
| Schema 前置要求 | 约束和索引应提前准备 | Tag、Edge、属性必须定义并等待传播 |
| 临时拓扑错误 | 由驱动和应用策略处理 | 适配器对部分 leader 错误有限重试 |
| 显式事务 | 支持,但不建议包住大型导入 | SessionPool 适配器不支持 |
两种后端都应在真实部署拓扑上进行容量和故障测试。开发环境中的单机吞吐、传播时间和错误行为不能直接代表生产集群。
常见问题
可以直接把一个 CSV 文件交给 Graph SDK 吗?
不能。应用需要解析 CSV,完成字段映射、类型转换、稳定 ID 和 Schema 校验,再提供 GraphNode 与 GraphEdge。可以用 GraphImportSource 封装这个过程和资源生命周期。
后续导入新文件需要创建新 Space 吗?
通常不需要。同一个知识库应继续使用原 Space,并通过稳定身份和来源元数据增量更新。只有隔离边界或不兼容 Schema 发生变化时,才考虑创建新 Space。
stopOnError=true 能保证全量回滚吗?
不能。它只保证发现失败后不再开始后续批次。已经提交的批次不会自动回滚。
设置 resumeFrom 后,同步 importData 会跳过数据吗?
当前不会。恢复点由 graph.imports() 的异步任务服务消费。同步流程需要调用方自己移动数据源位置,或直接从剩余数据构造 Iterable。
为什么 checkpoint 没有推进,但数据库里已有部分数据?
checkpoint 只在整个逻辑批次成功后推进。Nebula 批次内部是多条语句,失败前的语句可能已经生效;checkpoint 持久化本身也可能失败。恢复时应依赖稳定身份重试,并对数据库状态进行核验。
是否应该把整个全量导入放进一个事务?
通常不应该。大事务会长时间占用连接和锁,扩大超时、回滚和内存风险。批量导入应以可恢复、可对账的小批次为主;只有规模很小且业务确实要求原子性时才考虑显式事务。
生产检查清单
- 是否固定了目标 Space、Schema 版本和数据源版本;
- 原始数据是否已完成清洗、映射、类型转换和 Schema 校验;
- 节点 ID、边方向、类型和 rank 是否可稳定重建;
- 节点是否先于边提供,外部已存在端点是否经过验证;
- batchSize 是否经过真实后端压测,而不是凭经验无限调大;
- 是否明确 stopOnError 策略和允许的部分成功范围;
- 是否理解 Neo4j 与 Nebula 的批次原子性差异;
- 同步和异步路径是否使用了正确的进度、恢复和资源管理方式;
- checkpoint 是否持久化并监控失败,数据源顺序是否可复现;
- 异步任务存储是否满足重启恢复和多实例要求;
- 取消、超时和未知提交结果是否有幂等重试及对账方案;
- 文件更新和撤回是否通过来源/事实状态处理,而非盲目删节点;
- 导入完成后是否执行数量、引用、Schema、来源和关键查询校验;
- 任务、错误详情和数据质量结果是否纳入监控与审计。