GeminiDB 宽表搜索与向量检索增强——CSS数据同步特性深度解析
一、背景
在分布式数据库场景中,GeminiDB Cassandra凭借其分布式架构与LSM存储引擎,在海量数据高并发写入场景下表现优异,是物联网、日志分析、用户画像等业务场景的核心数据存储组件。然而,随着业务复杂度的持续提升,越来越多的客户在以下方面提出了新的需求:
· 全文检索:客户需要对文本字段进行分词搜索、模糊匹配、相关性排序等操作,而Cassandra的CQL仅支持主键和二级索引列的等值与范围查询,缺乏全文检索能力。
· 复杂条件查询:业务场景中大量存在多列组合条件查询的需求,传统Cassandra的查询模型难以高效支持,需要在应用层进行多次查询和结果拼接,开发与维护成本高。
· 聚合分析:面对统计分析、分组聚合等分析型需求,Cassandra原生能力有限,往往需要引入额外的大数据计算组件,进一步增加了系统复杂度。
· 实时可查:客户希望写入Cassandra的数据能够在短时间内被检索到,传统离线数据同步方案延迟大、运维复杂,无法满足业务实时性要求。
综上所述,在保持GeminiDB Cassandra高并发写入优势的同时,如何便捷地获得强大的复杂查询与全文检索能力,成为客户的核心诉求。
二、传统解决方案与新方案对比
2.1 传统方案一:直接使用Elasticsearch替代Cassandra
最直接的方式是直接使用Elasticsearch替代Cassandra作为数据存储与检索引擎。Elasticsearch天生具备倒排索引和DSL查询语法,能够灵活应对全文检索和多条件组合查询。然而,Elasticsearch在数据写入吞吐、存储成本和事务一致性方面相较于Cassandra存在明显短板,难以独自承担高并发写入场景下的核心数据存储职责。
2.2 传统方案二:独立同步中间件
为兼顾写入性能与查询灵活性,业界常见的做法是在Cassandra与Elasticsearch之间搭建独立的数据同步中间件,由中间件负责将Cassandra的变更数据同步至Elasticsearch。该方案虽然功能上可行,但引入了额外的同步组件,客户需要自行开发和维护同步逻辑,处理数据一致性、故障恢复、断点续传等复杂问题,系统的运维难度和故障风险显著上升。
2.3 新方案:GeminiDB 宽表搜索与向量检索增强

针对上述痛点,GeminiDB Cassandra推出宽表搜索与向量检索增强特性,将数据写入Cassandra后自动同步至华为云CSS(Cloud Search Service,云搜索服务)集群,借助Elasticsearch强大的倒排索引和DSL查询语法满足复杂查询需求。该方案在架构层面将写入路径与同步路径彻底解耦,通过本地存储引擎的WAL与文档存储双层机制实现高可用、可恢复的异步同步链路,在写入性能、查询灵活性和运维简便性之间取得了理想平衡。
三、特性实现
本特性的核心设计理念是"写入无阻塞,同步可恢复"。整个同步链路分为构建阶段和回放阶段:构建阶段由Cassandra写入线程负责将行数据转换为CSS bulk格式的文档,暂存到本地存储引擎的文档存储列族和WAL中;回放阶段由后台线程池负责从本地存储引擎中读取WAL记录和文档数据,通过HTTP Bulk API异步发送至CSS集群,成功后清理已同步的数据。两个阶段通过本地存储引擎的事务机制保证原子性,通过WAL的序列号保证顺序性。
3.1 CSS索引同步整体流程

CSS索引同步的整体流程可概括为以下步骤:
1. 客户端连接GeminiDB Cassandra,对需要执行数据同步的用户表执行创建CSS索引的CQL语句。
2. 内核在创建索引时执行初始化流程:在本地存储引擎中创建文档存储列族,注册WAL同步源防止WAL文件被过早回收,将当前WAL末尾序列号记录为增量回放的起点,启动后台回放线程。
3. 内核触发存量数据扫描任务,遍历用户表中所有已有数据,将每行数据转换为文档格式写入OpenSearch CF和WAL。
4. 后台回放线程持续从WAL中读取记录,从OpenSearch CF中获取对应的文档数,批量发送至CSS集群的Bulk API。
5. Bulk请求成功后,更新检查点(checkpoint),清理OpenSearch CF中已同步的数据。若失败则不更新检查点,下次循环从同一位置重试。
6. 后续用户写入的增量数据通过Indexer机制自动进入构建阶段,写入OpenSearch CF和WAL后由回放线程异步同步至CSS。
整体数据流路径如下:客户端CQL写入 → Cassandra写入路径(本地存储引擎) → 构建阶段(行→文档转换,写入WAL与文档存储) → 回放阶段(读取WAL与文档,异步发送至CSS) → CSS集群。在该路径中,从写入到WAL与文档存储的步骤是同步的(与Cassandra写入在同一事务中),但从文档存储到CSS的步骤是异步的,确保Cassandra的写入延迟不受CSS集群状态影响。
3.2 CSS索引构建流程

CSS索引构建在用户执行CREATE CUSTOM INDEX语句时触发,由索引初始化任务驱动,分为初始化和存量扫描两个阶段。
3.2.1 初始化阶段
初始化阶段执行以下关键操作:
· 在本地存储引擎中创建专用于暂存待同步文档数据的列族。该列族的Key格式为token与文档ID的组合,Value格式为时间戳与JSON文档的组合。
· 注册WAL同步源,告知WAL清理组件不要回收同步所需的WAL文件,确保增量回放不会因WAL被清理而丢失数据。
· 对每个shard,初始化同步序列号:若该shard的同步序列号为0(首次创建索引),则将其设为当前WAL的最后序列号。这一设计确保增量回放从索引创建时刻开始,之前的WAL由存量构建任务覆盖。
· 为每个shard启动回放实例,并将其注册到全局回放实例表中。
· 在CSS集群中创建同名索引,并自动注入保留字段_TIMESTAMP_INTERNAL(type=long, store=false),用于范围删除时按时间戳过滤文档。
3.2.2 索引选项校验
在索引创建前,系统执行严格的校验:
· 调用CSS的/_index_template/_simulate/ API,解析settings和mappings的最终合并结果(包含组件模板),确保配置的合法性。
· 校验mappings中列名与Cassandra表列的对应关系,确保不存在无效列名。
· 校验列数据类型兼容性:基于Cassandra到Elasticsearch的类型映射表,例如bigint对应long,text对应text/keyword等。
· strict模式下mappings必须包含所有主键列,否则不同行可能会相互覆盖。
· text类型主键列必须有keyword子字段,以支持范围删除时的term/range精确匹配。
· 保留字段_TIMESTAMP_INTERNAL不可使用。
3.2.3 文档ID生成策略
文档ID生成器负责将Cassandra行映射为CSS文档ID,支持两种策略:
默认策略:将所有主键列值做Hex编码后用"@"拼接。该方式下文档ID与主键一一对应,支持从文档ID反查出分区键。例如,主键为(country="CHN", city="Beijing")的行,其文档ID为"43484e@4265696a696e67@"。
配置指定策略:使用指定text列的值作为文档ID。若该列为主键列,行为与默认策略类似;若为普通列,则删除行时需要从旧行中读取该列值以确定文档ID,且不允许更新该列的值,否则会产生孤立文档。
3.2.4 Mapping模板机制
Mapping模板支持用Cassandra普通列的值填充Elasticsearch复杂数据类型(如point、object),语法为${column_name:default_value}。模板在构建阶段逐行从Row中取值替换占位符。若列值为Null则使用默认值。例如,配置"coordinates": ["${longitude:0}", "${latitude:0}"]时,若longitude和latitude列有值则直接使用,若为Null则分别填充默认值0。
3.3 CSS索引存量数据回放流程

存量数据回放由存量扫描构建器实现,负责在索引创建时将Cassandra中已有数据同步至CSS。该过程按shard粒度并行执行,各shard之间互不依赖。
3.3.1 存量扫描流程
对每个shard执行以下流程:
1. 检查shard构建状态:从系统表读取该shard的构建状态。若已标记为SUCCESS,则跳过该shard(支持断点续传场景下的增量恢复)。
2. 打开本地存储迭代器:若存在上次构建的断点key,则seek到断点位置继续扫描;否则seek到该token范围的起始位置开始全量扫描。
3. 逐行扫描处理:对每一行数据,创建事务包装器封装行索引器,依次执行开启事务、构建文档、提交事务等步骤。构建过程中将行数据转换为文档并写入WAL和文档存储列族。
4. 定期保存断点:每扫描200,000行,将当前迭代器key作为断点保存至系统表,确保构建中断后可从断点恢复。
5. 扫描完成标记:所有行扫描完毕后,将该shard标记为SUCCESS。
3.3.2 存量数据与增量回放的衔接
存量扫描产出的文档同样写入文档存储列族和WAL,由增量回放引擎统一发送至CSS。这种设计保证了:存量和增量数据使用相同的发送通道和限流机制;存量扫描期间的新写入不会丢失(增量回放同时运行);存量扫描与增量回放之间无需额外的协调机制,初始化时将同步序列号设为当前WAL末尾即可自然衔接。
3.3.3 存量构建的行处理细节
行索引器处理存量数据时有两个特殊行为:一是跳过仅含列删除标记的行,因为存量回放看到的是已合并的行视图,列级删除在合并时已生效;二是存量数据标识跳过已有文档的更新检查逻辑,直接覆盖写入。
3.4 CSS索引增量数据回放流程
增量数据回放实现Cassandra写入数据实时同步至CSS,是本特性的核心机制。分为构建阶段(写入时文档生成)和回放阶段(异步发送至CSS)。

3.4.1 构建阶段:写入时文档生成
当Cassandra执行数据写入时,索引机制在写入事务中同步执行以下操作:
行数据处理:首先读取本地存储中的旧行,根据新旧行状态判断操作类型——若旧行不存在且新行是删除,则忽略无意义删除;若旧行存在且新行是更新,则合并新旧行。随后生成文档ID和bulk格式JSON文档,行有效时生成索引操作(包含操作行和数据行),行删除时生成删除操作(仅包含操作行)。最后将WAL日志和文档存储数据写入本地存储引擎。
分区删除和范围删除处理:构建ES DSL查询,查询条件包含分区键精确匹配条件、聚类键范围条件以及时间戳条件(确保只删除删除时间戳之前写入的文档,避免误删删除后重新写入的数据)。text/varchar类型的主键列使用keyword子字段进行精确匹配和范围查询。删除操作仅写入WAL,不写入文档存储列族。
3.4.2 WAL与OpenSearch CF的双层存储机制
增量数据存储采用WAL + OpenSearch CF双层机制,该设计是经过多轮迭代优化的结果,解决了排序查询与同ID时间戳merge两个核心诉求的冲突。
WAL日志格式保证按写入顺序查询。其结构为:1字节操作类型(包括索引、删除、范围删除、分区删除四种)+ 8字节时间戳 + 操作体(索引和删除类型为token与文档ID的组合,范围删除和分区删除类型为删除查询DSL的序列化数据)。WAL由本地存储引擎保证插入顺序,且能够根据序列号进行范围查询,实现了高性能的排序查询。
文档存储列族格式支持同ID按时间戳合并。Key为token与文档ID的组合(相同文档ID的key相同),Value为时间戳与JSON文档的组合(通过本地存储引擎的merge操作按时间戳合并,新时间戳覆盖旧值)。这种设计使得回放阶段能够按文档ID合并同一行的多次修改,只将最终状态发送至CSS。
双层机制的设计演进过程:方案一将(id, doc)直接存入文档存储列族,但Cassandra时间戳与本地存储引擎的序列号不同步,并发写入时可能读到旧数据;方案二改为(id, ts+doc)存入文档存储列族并使用merge按时间戳判断,但遍历耗时过长,已同步数据残留导致性能退化;方案三改为(ts+id, ts+doc)存入文档存储列族并使用时间戳前缀查询,但不同时间同一id的key不同,无法按id合并;最终方案引入WAL + 文档存储列族双层机制,WAL保证顺序查询,文档存储列族支持同ID按时间戳合并,两者互补,完美解决了核心冲突。
3.4.3 回放阶段:WAL到CSS的异步同步
回放阶段由回放引擎驱动,其核心是一个持续运行的后台调度循环:
全局调度:后台调度循环遍历所有已注册的回放实例(以列族ID与token为key进行管理),对每个满足继续回放条件的实例提交回放任务到回放线程池。
单实例回放逻辑:首先通过非阻塞方式获取锁,失败则返回(防止同一shard的并发回放)。然后从上次检查点+1到当前WAL末尾序列号读取WAL日志迭代器,逐条解析WAL记录:对于索引和删除类型,从文档存储列族读取对应文档值,添加到批量请求批次;对于范围删除和分区删除类型,将删除查询DSL添加到删除请求批次。解析过程中通过限流器进行限流控制,当限流触发时立即刷新当前批次。
批量发送:当批量请求批次非空时,发送bulk请求至CSS集群;成功后通过本地存储引擎的merge机制写入墓碑时间戳清理已同步数据;失败则标记请求失败。当删除请求批次非空时,逐条发送删除请求。全部成功后更新检查点:将当前序列号写入文档存储列族的检查点键,同时更新shard元数据中的同步序列号(双写保证可靠性)。若有任何失败则不更新检查点,下次循环从原检查点位置重试。
3.4.4 限流与并发控制
每个回放实例持有限流器,速率由同步速率配置控制(默认5000文档/秒),该值同时作为bulk请求的最大批次大小。回放线程池支持运行时热加载调整线程池大小和限流速率。每个回放实例持有互斥锁,回放任务使用非阻塞方式获取锁,防止同一shard的并发回放导致数据重复发送。
3.5 故障恢复的实现

故障恢复是本特性最核心的设计考量之一。系统在多个层面实现了容错和恢复机制,确保在各种异常场景下数据不丢失、同步可恢复。
3.5.1 WAL与检查点双保险机制
增量回放的故障恢复核心依赖于WAL和检查点的配合。文档在构建阶段同时写入文档存储列族和WAL,且与Cassandra的写入在同一事务中,保证了原子性。回放引擎仅在CSS确认请求成功后才推进检查点,因此若进程崩溃,检查点将停留在最后一次成功确认的位置,未确认的文档仍保留在文档存储列族中,重启后回放引擎从检查点位置重新读取WAL,确保数据不丢失。检查点持久化在两个位置:文档存储列族中的检查点键(主要,重启后优先使用)和shard元数据(次要,当文档存储列族中的检查点缺失时回退使用),双写保证了检查点的可靠性。
3.5.2 存量构建的断点续传
存量构建过程中的断点续传通过将断点保存至系统表实现。构建过程中每扫描200,000行保存一次断点(当前迭代器key),构建状态从STARTED到SUCCESS逐步推进。若构建中断,重启后从最近保存的断点位置继续扫描,已标记SUCCESS的shard直接跳过,无需全量重跑。
3.5.3 初始化重试机制
索引初始化过程中,初始化任务捕获任何异常后,以固定间隔(默认3秒)重新调度初始化任务。存量构建失败时同样触发重试(文档ID相关的配置异常除外),确保初始化过程的鲁棒性。打开新shard时也实现了失败重试逻辑。
3.5.4 HTTP请求重试
HTTP连接器实现了重试机制:对于网络异常等非HTTP错误响应,自动重试一次;对于CSS返回HTTP错误码,直接抛出异常不重试,避免对不可恢复错误的无效重试。批量请求失败时不推进检查点,下次回放循环从原检查点位置重试整个批次。
3.5.5 WAL文件保护
初始化阶段注册WAL同步源,告知WAL清理组件不要回收同步所需的WAL文件。这一机制确保增量回放不会因WAL被垃圾回收而丢失数据,是故障恢复的基础保障。
3.5.6 强制恢复手段
对于异常场景,系统提供了两种恢复手段:一是全量重放,将检查点重置为起始位置,强制从文档存储列族头开始全量重放;二是强制残留回放,直接遍历文档存储列族中所有残留数据并发送至CSS(跳过WAL,用于WAL损坏等极端场景的恢复)。两种方法均为管理端操作,不在正常流程中自动触发。
四、使用流程与示例
4.1 创建用户表
在GeminiDB Cassandra中创建一张用户表:
CREATE TABLE test_ks.users(
country text,
city text,
name text,
age int,
PRIMARY KEY(country, city)
);
4.2 创建CSS索引
创建CSS索引,指定Elasticsearch的settings和mappings配置,数据将自动同步至CSS集群的同名索引:
CREATE CUSTOM INDEX user_css_idx ON test_ks.users ()
USING 'OpenSearchIndex'
WITH OPTIONS = {
'settings': '{"number_of_shards": 3}',
'mappings': '{
"properties": {
"country": {"type": "keyword"},
"city": {"type": "keyword"},
"name": {"type": "text"},
"age": {"type": "integer"}
}
}'
};
索引创建后,已有数据自动回放至CSS。新写入数据将实时同步。
4.3 写入数据与查询
向表中写入数据:
INSERT INTO test_ks.users(country, city, name, age)
VALUES('China', 'Beijing', 'ZhangSan', 28);
INSERT INTO test_ks.users(country, city, name, age)
VALUES('China', 'Shanghai', 'LiSi', 32);
INSERT INTO test_ks.users(country, city, name, age)
VALUES('Japan', 'Tokyo', 'WangWu', 25);
数据同步至CSS后,在CSS集群中查询user_css_idx索引,即可看到与Cassandra写入数据对应的文档:
// 文档1
{"country": "China", "city": "Beijing", "name": "ZhangSan", "age": 28}
// 文档2
{"country": "China", "city": "Shanghai", "name": "LiSi", "age": 32}
// 文档3
{"country": "Japan", "city": "Tokyo", "name": "WangWu", "age": 25}
4.4 更新与删除操作
当Cassandra端执行UPDATE操作时,CSS端对应文档同步更新:
UPDATE test_ks.users SET age = 29 WHERE country = 'China' AND city = 'Beijing';
// CSS端文档1更新为:
{"country": "China", "city": "Beijing", "name": "ZhangSan", "age": 29}
当Cassandra端执行DELETE操作时,CSS端对应文档同步删除:
DELETE FROM test_ks.users WHERE country = 'Japan' AND city = 'Tokyo';
// 删除后,CSS端将查询不到country为Japan、city为Tokyo的文档
4.5 使用Mapping模板
创建包含经纬度列的用户表,并通过cassandra_mapping_templates选项使用Elasticsearch的point类型:
CREATE TABLE test_ks.users_geo(
country text,
city text,
longitude int,
latitude int,
PRIMARY KEY(country, city)
);
CREATE CUSTOM INDEX users_geo_idx ON test_ks.users_geo ()
USING 'OpenSearchIndex'
WITH OPTIONS = {
'settings': '{"number_of_shards": 3}',
'mappings': '{
"properties": {
"country": {"type": "keyword"},
"city": {"type": "keyword"},
"longitude": {"type": "long"},
"latitude": {"type": "long"}
}
}',
'cassandra_mapping_templates': '{
"location": {
"type": "point",
"coordinates": ["${longitude}", "${latitude}"]
}
}'
};
写入数据后,CSS端文档中将自动包含由经纬度列值生成的location字段:
INSERT INTO test_ks.users_geo(country, city, longitude, latitude)
VALUES('China', 'Beijing', 116, 39);
// CSS端查询结果:
{
"country": "China",
"city": "Beijing",
"longitude": 116,
"latitude": 39,
"location": {
"type": "point",
"coordinates": [116, 39]
}
}
当普通列值为Null时,可使用默认值避免Null填充,语法为${column_name:default_value}。例如:
'cassandra_mapping_templates': '{
"location": {
"type": "point",
"coordinates": ["${longitude:90}", "${latitude:180}"]
}
}'
4.6 修改与删除CSS索引
如需修改索引的mappings配置,可通过ALTER INDEX命令在线修改:
ALTER INDEX user_css_idx WITH options = {
'mappings': '{
"dynamic": "strict",
"properties": {
"country": {"type": "keyword"},
"city": {"type": "keyword"},
"name": {"type": "text"},
"age": {"type": "integer"}
}
}'
};
删除CSS索引与普通索引操作一致:
DROP INDEX user_css_idx;
五、总结
GeminiDB 宽表搜索与向量检索增强特性通过"写入无阻塞,同步可恢复"的架构设计,在保持Cassandra高并发写入优势的同时,赋予了系统强大的复杂查询与全文检索能力。该特性的核心价值体现在以下方面:
架构层面的彻底解耦:构建阶段与回放阶段分离,Cassandra写入与CSS同步完全异步,写入延迟不受CSS集群状态影响。即使CSS集群不可用或网络异常,也不影响Cassandra的正常写入,系统会自动重试直至同步成功。
WAL + OpenSearch CF双层存储机制:经过多轮迭代优化,解决了排序查询与同ID时间戳merge两个核心诉求的冲突。WAL保证按写入顺序高效查询,OpenSearch CF支持同一文档ID按时间戳合并,两者互补实现了高性能、高一致性的数据暂存。
多层级故障恢复:从WAL检查点双写到存量构建断点续传,从初始化重试到HTTP请求重试,从WAL文件保护到强制恢复手段,系统在多个层面实现了容错和恢复机制,确保在各种异常场景下数据不丢失、同步可恢复。
全操作覆盖:支持INSERT、UPDATE、DELETE全类型数据同步,包括行删除、分区删除和范围删除。分区删除和范围删除通过deleteByQuery DSL实现,并利用_TIMESTAMP_INTERNAL保留字段确保只删除删除时间戳之前写入的文档,避免误删删除后重新写入的数据。
灵活的配置能力:支持自定义文档ID生成策略、Mapping模板机制(使用Cassandra列值填充Elasticsearch复杂数据类型)、运行时限流与线程池热加载,以及丰富的监控指标(同步速率、时延、积压量、失败次数等),满足不同业务场景的精细化管控需求。
当前版本支持CSS Elasticsearch 7.10.2。未来将进一步提升同步性能、扩展支持的Elasticsearch版本,并持续优化监控与运维体验。
- 点赞
- 收藏
- 关注作者
评论(0)