跳到主要内容

批量导入 ​

概述 ​

批量导入用于把一批已经转换成图模型的数据,持续、可观察地写入指定 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、消息队列或业务数据库。调用方需要先完成:

text
原始数据
  -> 解析与清洗
  -> 字段映射
  -> 实体归一和稳定 ID
  -> Schema 校验
  -> GraphNode / GraphEdge
  -> Graph SDK 导入

文档知识图谱可以先通过 extractor 按 Schema 提取实体和关系;结构化数据则可以通过自定义映射器转换。无论来源是什么,进入导入器前都应保证类型、属性名、边方向和端点身份已经确定。

3. 保证可重复的身份和顺序 ​

可靠重试需要:

  • 同一实体每次生成相同节点 ID;
  • 同一关系每次生成相同的 source、type、target 和 rank;
  • 恢复时数据源顺序与生成 checkpoint 时一致;
  • 数据源版本固定,不在恢复期间插入、删除或重新排序;
  • 节点和边的迭代过程不会静默跳过异常记录。

仅仅“再次读取同一个文件名”并不能证明内容和顺序相同。生产系统应保存内容哈希、对象版本号或数据库快照标识。

节点优先,边随后 ​

导入器始终先完整消费节点数据流,再消费边数据流:

text
节点批次 0 -> 节点批次 1 -> ... -> 边批次 N -> 边批次 N+1

这种顺序保证创建关系前,端点有机会先被写入。它也带来两个要求:

  • 不要把节点和边混在一个只有单次遍历能力的共享游标中;
  • 如果边引用的数据不在本次节点流中,应确保它们已存在于目标 Space。

节点阶段存在失败而 stopOnError=false 时,导入器仍会继续处理后续节点和边。此时依赖失败节点的边可能无法创建,所以“继续执行”只适合能够接受部分成功并有后续对账修复的任务。

同步导入 ​

同步导入在当前调用线程中完成,适合由应用自己的作业框架管理生命周期:

java
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提供节点和边数据
sourceBuilder 从中取得节点和边 Iterable
batchSize控制每次组成 Mutation 的记录数
stopOnError控制失败后立即抛出还是继续

listener、checkpoint 和 resumeFrom 由异步任务服务处理,直接调用同步 importData 时不会获得这些生命周期能力。同步使用 GraphImportSource 时,调用方还必须自行关闭数据源,推荐使用 try-with-resources:

java
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():

java
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:

java
GraphImportRequest request = GraphImportRequest.builder()
    .source(source)
    .batchSize(500)
    .checkpoint((taskId, nodesProcessed, edgesProcessed) ->
        checkpointRepository.save(
            taskId, nodesProcessed, edgesProcessed))
    .build();

任务失败或取消后,可以从最近持久化的确认位置构造恢复点:

java
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、来源证据和引用关系:

text
文档版本 -> 提取事实 -> 节点/边声明 -> 当前图状态

文件更新时先比较新旧事实,只撤销不再被任何有效来源支持的声明;文件撤回时同样按来源清理。相关流程见文档生命周期。

导入完成后的数据质量验证 ​

“任务成功”只表示所有 Writer 调用没有报告失败,不表示图谱业务上正确。至少应验证:

  • 节点和边数量是否在预期范围;
  • 是否存在引用不到端点的关系或遗漏关系;
  • 必填属性、数据类型和枚举是否符合 Schema;
  • 是否出现异常重复节点、别名未归一或错误主标签;
  • 关键边方向和 rank 是否正确;
  • 每个数据源或文档的来源记录是否完整;
  • 关键查询是否能返回预期结果;
  • 索引是否就绪,查询计划和延迟是否符合要求;
  • 报告数量、源数据数量和数据库查询数量能否合理解释。

对大型导入,建议先在隔离 Space 做抽样和全量校验,达到发布阈值后再切换上层路由。

Neo4j 与 Nebula 的差异 ​

维度Neo4jNebula Graph
单个逻辑批次普通 Writer 在数据库事务内执行多条 nGQL 顺序执行
批次失败当前批次不提交,之前批次保留当前批次内部也可能部分生效
节点类型支持多 LabelPortable 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、来源和关键查询校验;
  • 任务、错误详情和数据质量结果是否纳入监控与审计。

需要后台执行、取消和任务持久化边界时,继续阅读异步导入任务;需要多个小规模读写共同提交时,继续阅读事务。