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 或其它连接器重复修改为不同的整数值

相关新闻

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

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

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

2026/7/21 16:44:17 阅读更多 →
一位老馆长的遗憾,成了殡仪馆改造的起点

一位老馆长的遗憾,成了殡仪馆改造的起点

一位殡仪馆的老馆长说过一句话:“我这辈子最后悔的事,就是让成千上万的人在一间冷冰冰的屋子里,跟亲人说了最后一句话。” 他指的是自己馆里的告别厅——白色荧光灯管、灰白色瓷砖墙、成排的硬塑料椅,空调出风口正对着家属座位。 …

2026/7/21 16:44:12 阅读更多 →
报错记录 类型未定义警告

报错记录 类型未定义警告

type of ‘n’ defaults to ‘int’ [-Wimplicit-int]#include<stdio.h> #include<string.h>int fun(int n) {//第一层&#xff1a;5*fun(4)//第二层&#xff1a;4*fun(3)//第三层&#xff1a;3*fun(2)//第四层&#xff1a;2*fun(1)//第五层&#xff1a;return 1if…

2026/7/23 2:53:03 阅读更多 →

最新新闻

链上文脉·数启新声:文昌链持续引领文化数字资产创新实践

链上文脉·数启新声:文昌链持续引领文化数字资产创新实践

2026 年&#xff0c;文化数字化战略进入纵深推进的关键阶段&#xff0c;数字资产也正加速驶入规范发展的快车道。5 月&#xff0c;《北京日报》提出以数字藏品为载体的文化数字资产是首都文化金融改革创新核心新引擎。近日&#xff0c;浙江发布《数字文化 IP文旅消费融合发展实…

2026/7/24 10:35:32 阅读更多 →
DP83816以太网控制器:硬件过滤、WoL与电源管理实战解析

DP83816以太网控制器:硬件过滤、WoL与电源管理实战解析

1. 项目概述与核心价值 在嵌入式系统和工业网络设备的设计中&#xff0c;以太网控制器扮演着连接物理世界与数字网络的桥梁角色。它远不止是一个简单的“网卡”&#xff0c;其内部集成的智能过滤与电源管理逻辑&#xff0c;往往是决定整个系统能效、响应速度和网络健壮性的关键…

2026/7/24 10:35:32 阅读更多 →
AI-TEK转速传感器

AI-TEK转速传感器

AI-TEK转速传感器拥有覆盖无源、有源等不同技术路线的全系列产品&#xff0c;以下是主流系列及核心型号&#xff1a;一、无源速度传感器系列‌70082 系列&#xff08;侧视传感器&#xff09;‌核心定位&#xff1a;专为艾里逊自动变速箱开发的可变磁阻装置代表型号&#xff1a;…

2026/7/24 10:35:32 阅读更多 →
AI如何优化硕士开题报告:技术解析与实践指南

AI如何优化硕士开题报告:技术解析与实践指南

1. 开题报告痛点与AI解决方案第一次写硕士开题报告的研究生往往面临三重困境&#xff1a;不知道如何确定研究方向、不清楚技术路线设计是否合理、对文献综述的广度和深度把握不准。传统解决方式需要反复请教导师&#xff0c;但导师时间有限&#xff0c;学生容易陷入"改完格…

2026/7/24 10:35:32 阅读更多 →
嵌入式开发必读:TI芯片免责条款的工程化解读与合规实践

嵌入式开发必读:TI芯片免责条款的工程化解读与合规实践

1. 项目概述&#xff1a;为什么我们需要认真对待芯片厂商的“小字条款” 在嵌入式硬件开发这个行当里摸爬滚打了十几年&#xff0c;我见过太多工程师和项目经理把全部精力都扑在技术选型、原理图设计和代码调试上&#xff0c;却对采购回来的每一颗芯片、每一个模块数据手册最后…

2026/7/24 10:35:32 阅读更多 →
腾讯云大模型API调用与行业解决方案详解

腾讯云大模型API调用与行业解决方案详解

1. 腾讯云大模型服务全景概览 2026年的腾讯云大模型生态已经形成了完整的服务体系矩阵&#xff0c;覆盖从基础模型调用到行业解决方案的全链路能力。作为国内最早布局云计算AI服务的厂商之一&#xff0c;腾讯云目前提供的大模型入口主要分为四大类&#xff1a;基础API服务、行业…

2026/7/24 10:34:32 阅读更多 →

日新闻

用Highcharts 创建可拖拽三维散点立方体3D图表

用Highcharts 创建可拖拽三维散点立方体3D图表

该案例基于Highcharts scatter3d 三维散点图实现空间立方体散点可视化&#xff0c;核心特色&#xff1a;三维 X/Y/Z 三轴空间&#xff0c;所有散点分布在 0~10 立方体空间内&#xff1b;散点使用径向渐变实现立体 3D 圆球质感&#xff1b;支持鼠标 / 触屏拖拽画布&#xff0c;…

2026/7/24 0:00:29 阅读更多 →
AppCertDlls:进程创建路径上的 DLL 入口

AppCertDlls:进程创建路径上的 DLL 入口

AppCertDlls&#xff1a;进程创建路径上的 DLL 入口 AppCertDlls 位于 HKLM\System\CurrentControlSet\Control\Session Manager\AppCertDlls。本文的程序功能是只读列出这个键在 64 位和 32 位注册表视图中的全部值&#xff0c;并显示每条值的来源、名称、类型和可安全显示的数…

2026/7/24 0:00:29 阅读更多 →
我的编程之路:第一篇博客

我的编程之路:第一篇博客

大家好&#xff0c;我是一名编程初学者&#xff0c;同时这也是我编程学习之路上的第一篇博客。在这里&#xff0c;我想要向大家介绍我的一些想法和规划。a.自我介绍我是一个刚刚接触编程的新手&#xff0c;目前在学习c语言&#xff0c;我对编程世界充满了强烈的好奇。当然&…

2026/7/24 0:00:29 阅读更多 →

周新闻

Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中&#xff0c;我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源&#xff0c;还是配置文件、证书等&#xff0c;都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下&#xff0c;但这…

2026/7/24 3:59:20 阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP&#xff08;轻量级目录访问协议&#xff09;作为企业级身份认证的黄金标准&#xff0c;已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时&#xff0c;发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/24 1:23:39 阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击&#xff1a; https://intelliparadigm.com 第一章&#xff1a;AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”&#xff0c;而是以可解释、可审计、可迭代的方式&#xff0c;赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/23 17:49:47 阅读更多 →

月新闻