事务日志
JanusGraph 可以自动记录事务性变更,用于额外处理或作为变更记录。要为特定事务启用日志记录,请在事务开始时指定目标日志的名称。
tx = graph.buildTransaction().logIdentifier('addedPerson').start()
u = tx.addVertex(label, 'human')
u.property('name', 'proteros')
u.property('age', 36)
tx.commit()
提交时,事务期间所做的任何更改都将记录到用户日志系统中,记录到名为 addedPerson 的日志中。**用户日志系统**是一个可配置的日志后端,具有 JanusGraph 兼容的日志接口。默认情况下,日志写入主存储后端中的单独存储,可以按如下所述进行配置。在事务开始时指定的日志标识符标识了记录更改的日志,从而允许将不同类型的更改记录到单独的日志中以进行单独处理。
tx = graph.buildTransaction().logIdentifier('battle').start()
h = tx.traversal().V().has('name', 'hercules').next()
m = tx.addVertex(label, 'monster')
m.property('name', 'phylatax')
h.addEdge('battled', m, 'time', 22)
tx.commit()
JanusGraph 提供了一个用户事务日志处理器框架来处理记录的事务性变更。事务日志处理器通过 JanusGraphFactory.openTransactionLog(JanusGraph) 对先前打开的 JanusGraph 图实例打开。然后可以为持有事务性变更的特定日志添加处理器。
import java.util.concurrent.atomic.*;
import org.janusgraph.core.log.*;
import java.util.concurrent.*;
logProcessor = JanusGraphFactory.openTransactionLog(g);
totalHumansAdded = new AtomicInteger(0);
totalGodsAdded = new AtomicInteger(0);
logProcessor.addLogProcessor("addedPerson").
setProcessorIdentifier("addedPersonCounter").
setStartTimeNow().
addProcessor(new ChangeProcessor() {
@Override
public void process(JanusGraphTransaction tx, TransactionId txId, ChangeState changeState) {
for (v in changeState.getVertices(Change.ADDED)) {
if (v.label().equals("human")) totalHumansAdded.incrementAndGet();
}
}
}).
addProcessor(new ChangeProcessor() {
@Override
public void process(JanusGraphTransaction tx, TransactionId txId, ChangeState changeState) {
for (v in changeState.getVertices(Change.ADDED)) {
if (v.label().equals("god")) totalGodsAdded.incrementAndGet();
}
}
}).
build();
在此示例中,为名为 addedPerson 的用户事务日志构建了一个**日志处理器**,以处理使用 addedPerson 日志标识符的事务中所做的更改。向此日志处理器添加了两个**更改处理器**。第一个处理器计算添加的人类数量,第二个处理器计算添加到图中的神灵数量。
当针对特定日志(如上述示例中的 addedPerson 日志)构建日志处理器时,它将在成功构建和初始化后立即开始从日志中读取事务更改记录,直到日志头部。构建器中指定的开始时间标记了日志中日志处理器将开始读取记录的时间点。或者,可以在构建器中为日志处理器指定一个标识符。日志处理器将使用该标识符定期持久化其处理状态,即它将维护上次读取的日志记录上的标记。如果日志处理器稍后以相同的标识符重新启动,它将从上次读取的记录继续读取。当日志处理器需要长时间运行并且因此可能失败时,这尤其有用。在这种失败情况下,日志处理器可以简单地以相同的标识符重新启动。必须确保 JanusGraph 集群中的日志处理器标识符是唯一的,以避免在持久化读取标记上发生冲突。
更改处理器必须实现 ChangeProcessor 接口。它的 process() 方法为从日志中读取的每个更改记录调用,带有 JanusGraphTransaction 句柄、导致更改的事务 ID 以及一个持有事务更改的 ChangeState 容器。可以查询更改状态容器以检索作为更改状态一部分的单个元素。在示例中,检索了所有添加的顶点。有关 ChangeState 上所有查询方法的描述,请参阅 API 文档。提供的事务 ID 可用于调查事务的来源,该事务由执行事务的 JanusGraph 实例的 ID (txId.getInstanceId()) 和实例特定的事务 ID (txId.getTransactionId()) 的组合唯一标识。此外,事务时间可通过 txId.getTransactionTime() 获取。
更改处理器单独执行并在多个线程中执行。如果更改处理器访问全局状态,则必须确保此类状态允许并发访问。虽然日志处理器按顺序读取日志记录,但更改在多个线程中处理,因此不能保证日志顺序在更改处理器中保留。
请注意,日志处理器为日志中的每条记录至少运行一次每个注册的更改处理器,这意味着在某些故障条件下,单个事务更改记录可能会被多次处理。不能从正在运行的日志处理器中添加或删除更改处理器。换句话说,日志处理器在构建后是不可变的。要更改日志处理,请启动一个新的日志处理器并关闭现有的日志处理器。
logProcessor.addLogProcessor("battle").
setProcessorIdentifier("battleTimer").
setStartTimeNow().
addProcessor(new ChangeProcessor() {
@Override
public void process(JanusGraphTransaction tx, TransactionId txId, ChangeState changeState) {
h = tx.V().has("name", "hercules").toList().iterator().next();
for (edge in changeState.getEdges(h, Change.ADDED, Direction.OUT, "battled")) {
if (edge.<Integer>value("time")>1000)
h.property("oldFighter", true);
}
}
}).
build();
上面的日志处理器使用单个更改处理器处理 battle 日志标识符的事务,该处理器评估添加到 Hercules 的 battled 边。此示例演示了传递到更改处理器中的事务句柄是一个正常的 JanusGraphTransaction,它可以查询 JanusGraph 图并对其进行更改。
事务日志用例
变更记录
用户事务日志可用于记录对图所做的所有更改。通过使用单独的日志标识符,可以将更改记录在不同的日志中,以区分不同的事务类型。
随时可以构建日志处理器,该处理器可以从所需的开始时间开始处理所有记录的更改。这可用于取证分析、对不同的图重放更改或计算聚合。
下游更新
JanusGraph 图集群通常是更大体系结构的一部分。用户事务日志和日志处理器框架提供了将更改广播到整个系统的其他组件所需的工具,而不会减慢导致更改的原始事务。当事务延迟需要较低和/或有许多其他系统需要收到图中更改的警报时,这尤其有用。
触发器
用户事务日志提供了实现触发器的基本基础架构,这些触发器可以扩展到大量并发事务和非常大的图。触发器注册到特定的数据更改,并触发外部系统中的事件或对图的额外更改。在大规模情况下,不建议在原始事务中实现触发器,而是通过日志处理器框架以轻微延迟处理触发器。第二个示例显示了如何评估图中的更改并触发额外的修改。
日志配置
有许多配置选项可以微调日志处理器从日志中读取的方式。有关 log 命名空间下选项的完整列表,请参阅 配置参考。要配置用户事务日志,请使用 log.user 命名空间。其中列出的选项允许配置要使用的线程数、每个批次读取的日志记录数、读取间隔以及事务更改记录是否应在可配置的时间量 (TTL) 后自动过期并从日志中删除。
用户事务日志的示例配置可能如下所示
log.user.backend=default
log.user.fixed-partition=false
log.user.max-read-time=10000 ms
log.user.max-write-time=5000 ms
log.user.read-batch-size=1024
log.user.read-interval=3000 ms
log.user.read-lag-time=2000 ms
log.user.read-threads=1
log.user.send-batch-size=256
log.user.send-delay=0 ms
log.user.ttl=120000 ms
自定义存储后端
默认情况下,JanusGraph 支持一个日志后端实现,该实现由 log.user.backend 配置选项中的保留快捷方式 default 定义,并具有以下类实现 org.janusgraph.diskstorage.log.kcvs.KCVSLogManager。
KCVSLogManager 重用图的存储后端(由 storage.backend 定义)来存储所有日志。
通常,对于默认的仅管理操作,默认存储后端足以管理日志,并且不需要自定义日志存储后端。
话虽如此,如果使用用户事务日志,在高度加载的应用程序中,底层图的存储后端可能远非最佳,对于 log.user 日志,可能首选更好的存储(即 Kafka 或其他)。
因此,JanusGraph 允许通过 log.[x].backend 选项提供自定义日志实现。
自定义日志实现需要实现 org.janusgraph.diskstorage.log.LogManager 接口,并且实现类必须具有接受 org.janusgraph.diskstorage.configuration.Configuration 作为单个参数的公共构造函数。
该 LogManager 实现的主要目的是打开扩展 org.janusgraph.diskstorage.log.Log 接口的新日志记录实现。例如,由 KCVSLogManager 打开的日志具有 org.janusgraph.diskstorage.log.kcvs.KCVSLog 实现,它们管理图的存储后端中的日志,但用户可以实现使用 Kafka 存储后端或其他的 Log。
要指定任何自定义日志实现,需要通过 log.[x].backend 配置选项提供 LogManager 实现类的完整类路径。