Drasi Core
Drasi-core 是 Drasi 用于实现连续查询的库。
连续查询(Continuous Queries),顾名思义,是持续运行的查询。为了理解其独特之处,将其与开发人员通常针对数据库执行的即时查询进行对比是很有用的。
当您执行即时查询时,您是在某个时间点针对数据库运行该查询。数据库计算查询结果并返回。在处理这些结果时,您处理的是数据的静态快照,并且不知道在运行查询后数据可能发生的任何更改。如果您定期运行相同的即时查询,由于其他进程对数据所做的更改,每次查询结果可能不同。但要了解发生了什么变化,您需要将最新的结果与之前的结果进行比较。
连续查询一旦启动,就会持续运行直到被停止。在运行期间,连续查询将处理来自一个或多个数据源的变更,计算查询结果受到的影响,并输出差异。
连续查询是以 Cypher 查询语言编写的图查询。使用声明式图查询语言意味着您可以在单个查询中表达丰富的查询逻辑,同时考虑所查询数据的属性以及数据之间的关系。
Drasi-core 是 Drasi 用于实现连续查询的内部库。Drasi 本身是一个更广泛的解决方案,包含许多更多组件。 Drasi-core 可以独立于 Drasi 用于嵌入式场景,其中连续查询可以在应用程序内进程运行。
示例
在此场景中,我们有一组 Vehicles 和一组车辆可以所在的 Zones。 Drasi 中的概念数据模型是带标签的属性图,因此我们将车辆和区域作为图中的节点添加,并使用 LOCATED_IN 关系连接它们。
我们将创建一个连续查询来观察 Parking Lot 区域,以便在任何车辆进入或离开该区域时收到通知。
MATCH
(v:Vehicle)-[:LOCATED_IN]->(:Zone {type:'Parking Lot'})
RETURN
v.color AS color,
v.plate AS plate
当添加或删除 LOCATED_IN 关系时,Continuous Query 会发出一个 diff,说明 Vehicle 已被添加到查询结果中或已从查询结果中移除。并且,更改 Vehicle 的某个属性,例如 color,将导致查询发出一个 diff,说明该 Vehicle 已被更新。
让我们看看如何使用 QueryBuilder 来配置 Continuous Query。
let query_str = "
MATCH
(v:Vehicle)-[:LOCATED_IN]->(:Zone {type:'Parking Lot'})
RETURN
v.color AS color,
v.plate AS plate";
let function_registry = Arc::new(FunctionRegistry::new()).with_cypher_function_set();
let parser = Arc::new(CypherParser::new(function_registry.clone()));
let query_builder = QueryBuilder::new(query_str, parser)
.with_function_registry(function_registry);
let query = query_builder.build().await;
让我们将一个 Vehicle(v1)和一个 Zone(z1)作为节点加载到查询中。
我们可以通过向查询中处理一个 SourceChange::Insert 来实现这一点,这反过来需要一个 Element,它可以是 Element::Node 或 Element::Relation,它们代表图模型中可查询的节点和关系。 在构建一个 Element 时,您还需要提供包含其唯一标识符(ElementReference)、在带标签属性图中要应用于它的任何标签以及有效起始时间的 ElementMetadata。
query.process_source_change(SourceChange::Insert {
element: Element::Node {
metadata: ElementMetadata {
reference: ElementReference::new("", "v1"),
labels: Arc::new([Arc::from("Vehicle")]),
effective_from: 0,
},
properties: ElementPropertyMap::from(json!({
"plate": "AAA-1234",
"color": "Blue"
}))
},
}).await;
query.process_source_change(SourceChange::Insert {
element: Element::Node {
metadata: ElementMetadata {
reference: ElementReference::new("", "z1"),
labels: Arc::new([Arc::from("Zone")]),
effective_from: 0,
},
properties: ElementPropertyMap::from(json!({
"type": "Parking Lot"
})),
},
}).await;
我们可以使用 process_source_change 函数对连续查询计算数据变更对查询结果的影响。
query.process_source_change(SourceChange::Insert {
element: Element::Relation {
metadata: ElementMetadata {
reference: ElementReference::new("", "v1-location"),
labels: Arc::new([Arc::from("LOCATED_IN")]),
effective_from: 0,
},
properties: ElementPropertyMap::new(),
out_node: ElementReference::new("", "z1"),
in_node: ElementReference::new("", "v1"),
},
}).await;
Result: [Adding {
after: {"color": String("Blue"), "plate": String("AAA-1234")}
}]
query.process_source_change(SourceChange::Update {
element: Element::Node {
metadata: ElementMetadata {
reference: ElementReference::new("", "v1"),
labels: Arc::new([Arc::from("Vehicle")]),
effective_from: 0,
},
properties: ElementPropertyMap::from(json!({
"plate": "AAA-1234",
"color": "Green"
}))
},
}).await;
Result: [Updating {
before: {"color": String("Blue"), "plate": String("AAA-1234")},
after: {"color": String("Green"), "plate": String("AAA-1234")}
}]
query.process_source_change(SourceChange::Delete {
metadata: ElementMetadata {
reference: ElementReference::new("", "v1-location"),
labels: Arc::new([Arc::from("LOCATED_AT")]),
effective_from: 0,
},
}).await;
Result: [Removing {
before: {"color": String("Green"), "plate": String("AAA-1234")}
}]
更多示例
更多示例可在 examples 文件夹中找到。
动态插件
Drasi Core 包含一个 xtask 构建工具,用于构建、列出和发布动态插件——这些是运行时由 Drasi Server 加载的共享库(.so/.dylib/.dll)。
什么使一个 Crate 成为插件?
如果一个 crate 同时满足以下两个条件,它将被自动发现为动态插件:
- 在其
Cargo.toml中定义了dynamic-plugin特性 - 遵循命名约定
drasi-{type}-{kind},其中{type}是source、reaction或bootstrap之一
例如,drasi-source-postgres、drasi-reaction-log、drasi-bootstrap-mssql。
先决条件
- Rust 工具链(参见
rust-toolchain.toml) - 系统依赖项:
jq、libjq-dev、protobuf-compiler(Linux)或jq、protobuf(macOS) - 用于交叉编译:
cross(带有 Docker 的 Linux 主机)
xtask 命令
list-plugins — 发现并列出所有动态插件
扫描工作区中符合插件条件的 crate,并打印每个插件的类型、种类、版本和清单路径。同时显示工作区的 SDK、Core 和 Lib 版本。
cargo run -p xtask -- list-plugins
# or
make list-plugins
build-plugins — 构建插件共享库
将所有发现的插件 crate 构建为 cdylib 共享库。每个插件二进制文件放置在 target/<profile>/plugins/ 下(交叉构建时为 target/<triple>/<profile>/plugins/),并附带一个包含插件元数据的 metadata.json 伴随文件。
Flags:
| Flag | Description |
|---|---|
--release | 以 release 模式构建(默认:debug) |
--jobs N / -j N | 并行构建任务数 |
--target TRIPLE | 针对目标三元组进行交叉编译(例如 aarch64-unknown-linux-gnu) |
Examples:
# Build all plugins (debug)
make build-plugins
# Build all plugins (release)
make build-plugins-release
# Cross-compile for ARM Linux
cargo run -p xtask -- build-plugins --release --target aarch64-unknown-linux-gnu
交叉编译行为:
- 在 Linux 主机上,针对 Linux 和 Windows 目标使用
cross(基于 Docker)。 - 在 macOS 主机上,直接使用
cargo。仅支持 macOS 目标;Linux/Windows 目标将退出并显示清晰的错误消息。 - 同一操作系统上的跨架构构建(例如 macOS x86 → macOS ARM)使用
cargo配合--target。
生成的元数据(metadata.json):
每个插件二进制文件旁边都会写入一个 JSON 文件,包含以下字段:
{
"name": "drasi-source-postgres",
"kind": "postgres",
"type": "source",
"version": "0.1.8",
"sdk_version": "0.1.0",
"core_version": "0.1.0",
"lib_version": "0.1.0",
"target_triple": "x86_64-unknown-linux-gnu",
"description": "...",
"license": "Apache-2.0"
}
publish-plugins — 将插件发布为 OCI 制品
将已构建的插件发布到 OCI 容器注册表。每个插件都作为包含两个层的 OCI 制品推送:
| 层 | 媒体类型 |
|---|---|
插件二进制文件 (.so/.dylib/.dll) | application/vnd.drasi.plugin.v1+binary |
| 元数据 JSON | application/vnd.drasi.plugin.v1+metadata |
发布所有插件后,该命令还会更新插件目录——一个特殊的 OCI 包 (drasi-plugin-directory),其中每个标签代表一个已知插件(例如 source.postgres、reaction.storedproc-mssql)。这使得无需预先知道插件名称即可发现插件。
标志:
| 标志 | 描述 |
|---|---|
--registry <URL> | OCI 注册表(默认:ghcr.io/drasi-project) |
--plugins-dir <DIR> | 覆盖插件目录 |
--release | 在发布构建目录中查找 |
--target <TRIPLE> | 指定目标三元组以定位交叉编译的插件 |
--tag <TAG> | 覆盖所有插件的版本标签 |
--pre-release <LABEL> | 追加预发布标签(例如 dev.1) |
--arch-suffix <SUFFIX> | 在标签后追加架构后缀(例如 linux-amd64) |
--dry-run | 显示将要发布的内容而不实际推送 |
示例:
# Dry run
make publish-plugins-dry-run ARCH_SUFFIX=linux-amd64
# Publish release build for a single architecture
make publish-plugins-release ARCH_SUFFIX=linux-amd64
# Publish with a pre-release label
make publish-plugins-release ARCH_SUFFIX=linux-amd64 PRE_RELEASE=dev.1
# Publish to a custom registry
make publish-plugins-release ARCH_SUFFIX=linux-amd64 REGISTRY=ghcr.io/my-org
身份验证
设置以下环境变量:
export OCI_REGISTRY_USERNAME=<your-github-username>
export OCI_REGISTRY_PASSWORD=<your-pat-with-write-packages-scope>
PAT 需要 write:packages 权限,并且你的 GitHub 账户需要对目标组织的包具有写入访问权限。
标签格式
插件使用带平台后缀的标签 —— 架构始终作为标签后缀附加。不存在多架构清单索引;客户端在拉取时会自动附加正确的后缀。
| 格式 | 示例 |
|---|---|
| 正式发布 | ghcr.io/drasi-project/source/postgres:0.1.8-linux-amd64 |
| 预发布 | ghcr.io/drasi-project/source/postgres:0.1.8-dev.1-linux-amd64 |
| Musl | ghcr.io/drasi-project/source/postgres:0.1.8-linux-musl-amd64 |
支持的架构后缀:
| 后缀 | 目标三元组 |
|---|---|
linux-amd64 | x86_64-unknown-linux-gnu |
linux-arm64 | aarch64-unknown-linux-gnu |
linux-musl-amd64 | x86_64-unknown-linux-musl |
linux-musl-arm64 | aarch64-unknown-linux-musl |
windows-amd64 | x86_64-pc-windows-gnu |
darwin-amd64 | x86_64-apple-darwin |
darwin-arm64 | aarch64-apple-darwin |
插件目录
注册表中维护了一个名为 drasi-plugin-directory 的特殊 OCI 包。每个标签代表一个已知插件,使用 {type}.{kind} 格式(例如 source.postgres、reaction.storedproc-mssql)。使用 . 分隔符是因为插件类型从不包含点,从而避免了与插件类型名称中的连字符产生歧义。
这使得 Drasi Server 中的 plugin search 命令能够通过列出目录标签来发现可用插件,然后从每个匹配的插件包中获取版本信息。
发布所有架构
publish-all 目标按顺序为所有 7 种支持的架构构建并发布插件:
# Publish all architectures
make publish-all
# Dry run
make publish-all-dry-run
# With pre-release label
make publish-all PRE_RELEASE=dev.1
CI 工作流
.github/workflows/publish-plugins.yml 工作流自动化了跨所有 7 种架构的发布。它使用构建矩阵:
- Linux 目标(
x86_64-unknown-linux-gnu、aarch64-unknown-linux-gnu、x86_64-unknown-linux-musl、aarch64-unknown-linux-musl)通过cross在ubuntu-latest上构建 - macOS 目标(
x86_64-apple-darwin、aarch64-apple-darwin)通过cargo在macos-latest上构建 - Windows 目标(
x86_64-pc-windows-gnu)通过cross在ubuntu-latest上构建
所有构建成功后,一个可见性步骤会将所有软件包(包括 drasi-plugin-directory)在 GHCR 上设置为公开。通过 GitHub Actions UI 中的 workflow_dispatch 触发该工作流。
Makefile 参考
| 目标 | 描述 |
|---|---|
make list-plugins | 列出所有发现的插件 crate |
make build-plugins | 构建所有插件(调试) |
make build-plugins-release | 构建所有插件(发布) |
make test-host-sdk | 构建测试插件并运行 host-sdk 集成测试 |
make publish-plugins | 发布插件(调试构建) |
make publish-plugins-release | 发布插件(发布构建) |
make publish-plugins-dry-run | 预览将要发布的内容 |
make publish-all | 为所有 7 种架构构建并发布 |
make publish-all-dry-run | publish-all 的试运行 |
所有发布目标接受可选变量:REGISTRY、PRE_RELEASE、ARCH_SUFFIX。
运行 Host-SDK 集成测试
# Build test plugins and run integration tests
make test-host-sdk
存储实现
Drasi 维护内部索引,用于计算数据变更对查询结果的影响。默认情况下,这些索引位于内存中,但连续查询可以配置为使用持久化存储。 目前,存在针对 Redis、Garnet 和 RocksDB 的存储实现。
发布状态
drasi-core 库是 Drasi 早期发布的一个组件,旨在让社区了解并试验该平台。请告诉我们您的想法,并在发现 bug 或希望请求新功能时提交 Issue。Drasi 尚不适用于生产工作负载。
贡献
请参阅 贡献指南