Spring Boot 整合 Debezium 实现 MySQL 增量数据监听(嵌入式版)

发布时间:2026/9/15 7:07:54

Spring Boot 整合 Debezium 实现 MySQL 增量数据监听(嵌入式版) 一、背景与选型在微服务或数据同步场景中我们经常需要实时捕获 MySQL 数据库的增删改操作并触发后续业务逻辑如刷新缓存、同步到 ES、发送消息等。常见的方案有Canal阿里开源Debezium基于 Kafka Connect但也可嵌入式运行Maxwell本文选择Debezium Embedded Engine因为它无需依赖 Kafka可直接在 Spring Boot 应用中内嵌运行支持 MySQL、PostgreSQL、MongoDB 等多种数据库提供完整的变更事件结构before/after、元数据容错性好支持断点续传offset 存储。适用场景中小型项目希望快速集成 CDCChange Data Capture功能又不希望引入额外中间件。二、环境准备2.1 软件版本JDK 17Spring Boot 3.x 要求Spring Boot 3.xMySQL 5.7 或 8.0需开启 binlogDebezium 2.7.x2.2 MySQL 开启 binlog编辑 MySQL 配置文件my.cnf或my.ini添加以下内容[mysqld] log_bin mysql-bin binlog_format ROW binlog_row_image FULL server_id 1 # 确保唯一不能与 Debezium 的 database.server.id 冲突重启 MySQL 服务后执行 SQL 验证SHOW VARIABLES LIKE log_bin; -- 应为 ON SHOW VARIABLES LIKE binlog_format; -- 应为 ROW2.3 创建测试库表和 CDC 用户CREATE USER debezium% IDENTIFIED BY dbz; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO debezium%; FLUSH PRIVILEGES; -- 创建测试库和表 CREATE DATABASE IF NOT EXISTS park; USE park; CREATE TABLE test_user ( id INT AUTO_INCREMENT PRIMARY KEY, name VARCHAR(50), age INT, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );三、创建 Spring Boot 项目使用 IDEA 或 Spring Initializr 创建一个 Spring Boot 项目引入以下依赖pom.xmlproperties java.version17/java.version debezium.version2.7.0.Final/debezium.version /properties dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency !-- Debezium 核心 API -- dependency groupIdio.debezium/groupId artifactIddebezium-api/artifactId version${debezium.version}/version /dependency !-- Debezium 嵌入式引擎 -- dependency groupIdio.debezium/groupId artifactIddebezium-embedded/artifactId version${debezium.version}/version /dependency !-- Debezium 文件存储用于 offset 和 schema history -- dependency groupIdio.debezium/groupId artifactIddebezium-storage-file/artifactId version${debezium.version}/version /dependency !-- Debezium MySQL 连接器 -- dependency groupIdio.debezium/groupId artifactIddebezium-connector-mysql/artifactId version${debezium.version}/version /dependency /dependencies四、核心代码编写-可以写在yml里面4.1 配置类DebeziumConfig集中管理连接器配置方便后续调整。package com.example.demo.config; import org.springframework.context.annotation.Configuration; import java.util.Properties; Configuration public class DebeziumConfig { public Properties getDebeziumProperties() { Properties props new Properties(); // 连接器名称 props.setProperty(name, mysql-connector); props.setProperty(connector.class, io.debezium.connector.mysql.MySqlConnector); // 偏移量存储记录消费进度 props.setProperty(offset.storage, org.apache.kafka.connect.storage.FileOffsetBackingStore); props.setProperty(offset.storage.file.filename, ./offsets.dat); props.setProperty(offset.flush.interval.ms, 60000); // 数据库连接 props.setProperty(database.hostname, localhost); props.setProperty(database.port, 3306); props.setProperty(database.user, debezium); props.setProperty(database.password, dbz); props.setProperty(database.server.id, 184054); // 必须唯一 props.setProperty(topic.prefix, dbserver); // 事件主题前缀 // 过滤只监听 park 库下的 test_user 表 props.setProperty(database.include.list, park); props.setProperty(table.include.list, park.test_user); // 时区与 SSL props.setProperty(database.connectionTimeZone, UTC); props.setProperty(database.use.ssl, false); props.setProperty(database.allowPublicKeyRetrieval, true); // Schema 历史存储用于 DDL 变更跟踪 props.setProperty(schema.history.internal, io.debezium.storage.file.history.FileSchemaHistory); props.setProperty(schema.history.internal.file.filename, ./dbhistory.dat); // 快照模式never 表示不执行初始快照只监听增量 props.setProperty(snapshot.mode, never); return props; } }参数说明snapshot.modenever启动后不进行全表快照只监听后续变更。若希望首次启动时先同步现有数据可改为initial。offset.storage.file.filename记录已消费的 binlog 位置重启后从中断处继续。schema.history.internal.file.filename记录表结构变化历史用于正确解析事件中的字段。4.2 监听器组件DebeziumListener负责启动 Debezium 引擎接收并解析变更事件。package com.example.demo.listener; import com.example.demo.config.DebeziumConfig; import com.fasterxml.jackson.databind.ObjectMapper; import io.debezium.engine.ChangeEvent; import io.debezium.engine.DebeziumEngine; import io.debezium.engine.format.Json; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import java.io.IOException; import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; Component public class DebeziumListener { private static final Logger LOG LoggerFactory.getLogger(DebeziumListener.class); Autowired private DebeziumConfig debeziumConfig; private final ExecutorService executor Executors.newSingleThreadExecutor(); private DebeziumEngineChangeEventString, String engine; private final ObjectMapper objectMapper new ObjectMapper(); PostConstruct public void start() { var props debeziumConfig.getDebeziumProperties(); this.engine DebeziumEngine.create(Json.class) .using(props) .notifying(this::handleChangeEvent) .build(); executor.execute(() - { try { engine.run(); } catch (Exception e) { LOG.error(Debezium engine runtime error: , e); } }); LOG.info(Debezium Engine started asynchronously.); } /** * 处理每个变更事件 */ private void handleChangeEvent(ChangeEventString, String event) { String value event.value(); if (value null) return; try { MapString, Object payload objectMapper.readValue(value, Map.class); MapString, Object payloadData (MapString, Object) payload.get(payload); if (payloadData null) return; // 获取源信息 MapString, Object source (MapString, Object) payloadData.get(source); String db (String) source.get(db); String table (String) source.get(table); // 操作类型cinsert, uupdate, ddelete String operation (String) payloadData.get(op); if (operation null) { LOG.warn(Operation is null, skipping event for {}.{}, db, table); return; } MapString, Object before (MapString, Object) payloadData.get(before); MapString, Object after (MapString, Object) payloadData.get(after); switch (operation) { case c - { LOG.info(插入数据 on {}.{}: {}, db, table, after); // TODO: 业务处理 } case u - { LOG.info(更新数据 on {}.{}: before{}, after{}, db, table, before, after); // TODO: 业务处理 } case d - { LOG.info(删除数据 on {}.{}: {}, db, table, before); // TODO: 业务处理 } default - LOG.debug(Unknown operation: {}, operation); } } catch (IOException e) { LOG.error(Error parsing change event: {}, e.getMessage(), e); } } PreDestroy public void stop() { if (engine ! null) { try { engine.close(); } catch (Exception e) { LOG.error(Error closing Debezium engine, e); } } executor.shutdownNow(); LOG.info(Debezium Engine stopped.); } }注意在handleChangeEvent中你可以注入 Service 层将变更数据同步到 Redis、Elasticsearch 、Mqtt。五、运行与验证5.1 启动 Spring Boot 应用启动主类观察控制台输出Debezium Engine started asynchronously.5.2 在 MySQL 中执行 DML 操作-- 插入 INSERT INTO park.test_user (name, age) VALUES (张三, 25); -- 更新 UPDATE park.test_user SET age 26 WHERE name 张三; -- 删除 DELETE FROM park.test_user WHERE name 张三;5.3 查看应用日志你会看到类似以下格式的事件日志六、进阶说明6.1 数据格式详解Debezium 输出的 JSON 结构大致如下{ payload: { before: { id: 1, name: 张三, age: 25 }, after: { id: 1, name: 张三, age: 26 }, source: { db: park, table: test_user, server_id: 184054, ts_ms: 1721212345678, gtid: null, file: mysql-bin.000001, pos: 1234 }, op: u, ts_ms: 1721212345678 } }opc插入u更新d删除r表示快照若开启。before/after分别为变更前/后的行数据删除操作只有 before。6.2 断点续传原理offsets.dat文件记录了当前消费的 binlog 位置文件名 偏移量。重启应用后Debezium 会从该位置继续读取不会丢失数据。若想重新消费全部数据只需删除offsets.dat和dbhistory.dat并将snapshot.mode改为initial。6.3 多表监听修改table.include.list为多个表用逗号分隔props.setProperty(table.include.list, park.test_user, park.another_table);也可以通过database.include.list监听的库再通过table.exclude.list排除部分表。七、常见问题及解决方法问题现象可能原因解决方案启动时连接 MySQL 失败用户权限不足 / SSL 问题检查用户授权添加database.allowPublicKeyRetrievaltrue无任何变更事件输出binlog 未开启或格式不是 ROW检查 MySQL 配置确认binlog_formatROW事件中 before/after 为 nullbinlog_row_image不是 FULL设置为FULL并重启 MySQL重启后重复消费或漏消费offset 文件损坏删除offsets.dat和dbhistory.dat设置snapshot.modeinitial重新同步解析 JSON 异常表结构变更未正确记录确保schema.history.internal.file.filename文件持久化不要删除database.server.id冲突与 MySQL 的 server_id 或其它连接器重复修改为不同的整数值
延伸阅读

更多相关文章

2026/9/15 1:19:04

AI数字人直播实时美颜与背景替换技术原理分析

背景在AI数字人直播场景中,画面质量直接影响观众停留和转化。实时美颜和背景替换作为两项基础但关键的画面处理技术,其技术选型直接决定了端侧算力开销和渲染延迟。技术原理实时美颜AI直播中的美颜技术主要基于轻量级CNN(卷积神经网络&#x…

2026/9/12 22:51:03

如果关注M4Markets信息透明度,是否自然?

把如果关注信息透明度,是否自然放进真实使用情境里观察,M4Markets是否重视基础体验就会更有条理。从信息透明角度观察,平台把复杂事项拆解得更容易理解,用户自然更容易形成稳定印象。这些细节拼在一起,才构成M4Markets…

2026/9/6 5:39:10

打工人必看!一个网站=GPT+Claude+MJ,做PPT/写周报效率翻倍

别再让“工具切换”吃掉你的心流时间。 真正的效率,不是跑得更快,而是少绕弯路。先问一个扎心的问题: 你上一次专注写一份方案,中间被打断了几次? 我不是说同事消息或会议提醒,而是—— 写到一半想配张图&a…

2026/9/15 7:06:38

基于51单片机的低频迷你示波器设计:采样率、ADC与PCB全解析

简介:面向51单片机学习者的迷你示波器设计资料,完整覆盖从电路原理、PCB布局到固件实现的全流程。资源包为RAR压缩格式,共31个文件、约3.88MB,内含原理图与PCB设计文件(schdoc、pcbdoc、schlib、pcblib)及预…

2026/9/15 7:06:38

数字仿真中的时间片与Delta Cycle详解

1. 理解数字仿真中的时间片与Delta Cycle在数字电路仿真领域,时间片(Time Slot)和Delta Cycle是理解仿真器内部工作机制的两个核心概念。作为一名使用Modelsim/QuestaSim多年的工程师,我发现很多初学者虽然能完成基本仿真&#xf…

2026/9/15 7:06:38

面向毕业设计的物业管理系统开发:C++与Qt+MySQL实战解析

简介:这是一个基于 C 与 Qt 框架结合 MySQL 数据库开发的物业管理系统毕业设计源码包,适合计算机相关专业学生在毕业设计、课程设计或 Qt 项目实训中参考,解决物业管理中用户登录、权限区分、房间与投诉管理等常见业务需求。项目代码均已在本…

2026/9/15 7:06:38

擦窗机器人实测:10款主流产品优缺点全解析

1. 擦窗机器人市场现状与用户痛点我花了三个月时间测试了市面上主流的10款擦窗机器人,总花费超过3万元。作为一个家住28层、每周都要请保洁擦窗的上班族,最初我对这类产品充满期待——毕竟高空擦窗不仅危险,每次200元的保洁费用长期下来也是笔…

2026/9/15 7:06:38

Python实现番茄钟:GUI开发与时间管理实践

1. 为什么选择Python实现番茄钟?番茄钟作为一种经典的时间管理工具,其核心逻辑是通过25分钟工作5分钟休息的循环来提升专注力。选择Python来实现主要基于三点考虑:首先,Python的标准库time和tkinter已经包含了我们需要的所有基础功…

2026/9/15 7:01:38

打家劫舍全解:一维动态规划从递归到滚动数组优化

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/15 4:54:30

拯救者Y7000黑屏故障排查与维修实战指南

1. 项目概述:一台黑屏的拯救者Y7000,到底卡在哪一步? 联想拯救者Y7000系列笔记本,从2018年第一代搭载i5-8300H开始,到后来的i7-9750H、i7-10750H、i5-11400H,再到2023年款的R7-7840HS,它始终是学…

2026/9/15 0:01:16

AI英语单词APP开发:自适应学习算法与移动端优化实践

1. 项目概述 作为一名在移动应用开发领域摸爬滚打多年的老手,我最近完成了一个AI英语单词APP的开发项目。这个项目将传统单词记忆方法与现代AI技术相结合,打造了一款能够智能适应不同用户学习习惯的英语学习工具。 市面上大多数单词APP都存在一个通病&a…

2026/9/15 0:01:16

Flutter与OpenHarmony结合开发手语学习APP实战

1. 项目背景与核心价值作为一名同时接触过Flutter和OpenHarmony的开发者,最近我完成了一个基于Flutter for OpenHarmony的手语学习APP实战项目。这个项目最大的特点在于实现了跨平台框架与国产操作系统深度结合的创新实践——用Flutter开发的应用能完美运行在OpenHa…

2026/9/15 0:01:16

六个月成为机器人工程师:从ROS2到SLAM的实战路径

1. 六个月的紧迫感从哪来:先搞清楚你要成为哪种机器人工程师说实话,六个月的期限并不是一个宽松的时间线。市面上任何一本正经的机器人学教材都超过五百页,ROS2的官方文档可以翻到你怀疑人生,再加上ABB、KUKA这些工业机器人厂家动…

2026/9/14 11:59:31

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/14 13:53:59

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/14 11:22:57

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码