ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Bangumi Server Canal 消费者源码解析:Debezium 实时同步搜索索引

Bangumi Server Canal 消费者源码解析:Debezium 实时同步搜索索引 Bangumi Server Canal 消费者源码解析Debezium 实时同步搜索索引【免费下载链接】serverAPI server for bgm.tv项目地址: https://gitcode.com/gh_mirrors/server17/serverBangumibgm.tv的 API server 开源项目在canal/目录下实现了一套基于 Debezium 和 Kafka 的 binlog 订阅机制用来实时同步搜索索引。简单来说Debezium 监听 MySQL 的 binlog把数据变更转成事件写入 Kafka项目中的 Canal 消费者再消费这些事件把变更同步到 Meilisearch 搜索索引同时顺手处理会话撤销、缓存清理等脏活。这篇文章将以源码为主线带你搞懂 Canal 消费者实时同步搜索索引的完整链路。为什么要做 binlog 实时同步传统做法是写代码时手动更新搜索索引但这样容易漏更、错更比如直接在数据库改了数据忘了触发索引更新多个写入入口每个入口都要维护一段同步逻辑数据被重定向、封禁后索引里还残留旧内容。通过Debezium 订阅 MySQL binlog所有写入无论来自哪个入口都会变成统一的事件流Canal 消费者只需关心表变了 → 更新索引逻辑收敛到一处可靠且不容易遗漏。Canal 消费者整体架构整个模块的入口是 cmd/canal/main.go它调用 canal/canal.go 里的Main()完成组装和启动。核心结构如下用go.uber.org/fx做依赖注入把数据库、Redis、搜索客户端、S3 等组件装配好通过 canal/stream_kafka.go 创建 Kafka 消费者消费组 ID 为go-canal启动一个 HTTP 服务暴露/metrics方便用 Prometheus 监控消费情况消息处理统一交给eventHandler按表名分发。用一张图表示整个数据流MySQL binlog ── Debezium ── Kafka Topic ── Canal 消费者 ── Meilisearch 搜索索引 ├── 撤销用户会话改密码 ├── Redis 发布通知 └── 清理 S3 头像缓存从 Kafka 拉取消息一个可靠的消费者循环canal/stream_kafka.go 使用segmentio/kafka-go库实现消费逻辑kafkaStream.Read()就是一个无限循环FetchMessage拉取一条消息交给onMessage处理处理成功后CommitMessages提交 offset遇到网络错误或处理失败记录日志后继续拉取不阻塞整个消费组。这种处理成功才提交的模式属于at-least-once至少一次语义保证消息不丢失个别重复处理对索引更新来说是幂等的更新同一份文档代价可接受。消息分发一张表对应一个处理器Canal 消费者的核心分发逻辑在 canal/event.go 的onMessage中。它先解析 Debezium 的 Payload再根据source.table把事件路由到对应的处理方法MySQL 表处理器主要动作chii_subjects/chii_subject_fieldsOnSubject/OnSubjectField同步条目搜索索引chii_charactersOnCharacter同步角色搜索索引chii_personsOnPerson同步人物搜索索引chii_membersOnUserChange撤销会话、发布通知、清理缓存每种实体条目、角色、人物的变更处理器结构几乎一致比如 canal/on_subject.go 里根据 Debezium 的op字段决定动作ccreate新增→ 触发索引新增uupdate更新→ 触发索引更新ddelete删除→ 触发索引删除rsnapshot快照→ 生产环境 Debezium 已禁用快照代码里忽略处理。Debezium 消息格式看懂 Payload 就够了Debezium 写入 Kafka 的消息体是一个 JSONCanal 消费者只关心其中几个字段见 canal/event.go 的Payload结构before变更前的数据删除/更新时才有after变更后的数据新增/更新时才有source.table来自哪张表op操作类型c / u / d / r。还有一个细节Debezium 在删除记录时会发一条tombstone墓碑消息value 为空消费者遇到空 value 直接忽略即可这是 Kafka 压缩清理日志的机制不是业务事件。搜索索引同步核心逻辑在这里事件解析完成后会调用 internal/search/search.go 中Client接口的三个方法EventAdded、EventUpdate、EventDelete。它内部按目标维护了三套搜索器条目搜索器internal/search/subject/角色搜索器internal/search/character/人物搜索器internal/search/person/以条目为例internal/search/subject/event.go 的OnUpdate逻辑非常清晰从数据库 repo 重新读取该条目的最新数据如果条目已被重定向Redirect ! 0或封禁Ban ! 0则退化为删除索引文档避免把坏数据留在搜索结果里否则把数据抽取成搜索文档extract调用 Meilisearch 的UpdateDocumentsWithContext更新索引删除操作则直接调用DeleteDocumentWithContext移除文档。这样每次数据变更搜索索引都能在秒级内拿到数据库的最新状态用户搜索时就能看到刚修改的内容。顺带处理的脏活会话撤销与缓存清理除了同步搜索索引Canal 消费者还承担了一些只有感知到数据变更才能做好的事情见 canal/on_user.go 的OnUserChange密码被修改立即调用session.RevokeUser撤销该用户的所有登录会话强制重新登录防止旧密码会话继续有效新通知数变化向 Redis 的event-user-notify-{uid}频道发布消息实时通知 WebSocket 等在线端头像更换异步分页列出 S3 上该头像的所有缩略图缓存并删除让新头像尽快生效代码里还兼容了hd1的高清路径前缀。这些附加处理体现了 binlog 订阅的价值一个事件流多处受益。目录结构速览canal/canal.go消费者主流程与依赖装配canal/event.go消息解析与表分发canal/stream_kafka.goKafka 消费循环canal/on_subject.go、on_character.go、on_person.go三类实体的索引同步canal/on_user.go用户变更的附加处理internal/search/Meilisearch 搜索索引的读写封装总结Bangumi Server 的 Canal 消费者是一个教科书式的CDCChange Data Capture落地案例用 Debezium 捕获 MySQL binlog用 Kafka 做事件管道再用一个专注的 Go 消费者把变更实时同步到 Meilisearch 搜索索引同时兼顾会话安全与缓存一致性。对于想给自己的项目加数据库变更自动同步搜索能力的同学来说这份源码思路清晰、代码量小非常值得参考。【免费下载链接】serverAPI server for bgm.tv项目地址: https://gitcode.com/gh_mirrors/server17/server创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表