最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot整合Canal+RabbitMQ監(jiān)聽數(shù)據(jù)變更詳解

 更新時(shí)間:2024年12月30日 09:37:46   作者:星辰@Sea  
在現(xiàn)代分布式系統(tǒng)中,實(shí)時(shí)獲取數(shù)據(jù)庫(kù)的變更信息是一個(gè)常見的需求,本文將介紹SpringBoot如何通過整合Canal和RabbitMQ監(jiān)聽數(shù)據(jù)變更,需要的可以參考下

需求

在現(xiàn)代分布式系統(tǒng)中,實(shí)時(shí)獲取數(shù)據(jù)庫(kù)的變更信息是一個(gè)常見的需求。例如,在電商系統(tǒng)中,當(dāng)訂單表發(fā)生更新時(shí),可能需要同步這些變更到搜索服務(wù)、緩存服務(wù)或者通知其他微服務(wù)。傳統(tǒng)的解決方案包括定時(shí)輪詢數(shù)據(jù)庫(kù)或通過觸發(fā)器將變更寫入消息隊(duì)列等方法,但這些方案要么效率低下,要么實(shí)現(xiàn)復(fù)雜。而使用 Canal + RabbitMQ 可以提供一種高效且可靠的方式來(lái)捕獲 MySQL 數(shù)據(jù)庫(kù)的變更,并將其發(fā)送到 RabbitMQ 中供其他服務(wù)消費(fèi)。

Canal 是阿里巴巴開源的一個(gè)用于增量訂閱和消費(fèi) MySQL 數(shù)據(jù)庫(kù) Binlog 的工具,它模擬 MySQL 主從復(fù)制機(jī)制,無(wú)需侵入業(yè)務(wù)邏輯即可捕獲數(shù)據(jù)庫(kù)變更。RabbitMQ 是一個(gè)流行的開源消息代理,支持多種協(xié)議并提供了豐富的特性來(lái)確保消息傳遞的可靠性。結(jié)合這兩者,可以構(gòu)建一個(gè)強(qiáng)大的實(shí)時(shí)數(shù)據(jù)變更監(jiān)聽和處理系統(tǒng)。

步驟

環(huán)境搭建

整合SpringBoot與Canal實(shí)現(xiàn)客戶端

Canal整合RabbitMQ

SpringBoot整合RabbitMQ

環(huán)境搭建

1. 安裝MySQL

確保你有一個(gè)正在運(yùn)行的 MySQL 實(shí)例,并啟用了 binlog 日志記錄功能。這是 Canal 捕獲數(shù)據(jù)庫(kù)變更的基礎(chǔ)。

# 修改 MySQL 配置文件 my.cnf 或 my.ini
[mysqld]
server-id=1
log-bin=mysql-bin
binlog-format=ROW

重啟 MySQL 服務(wù)使配置生效。

2. 安裝Canal Server

下載最新版本的 Canal Server 并解壓到合適的位置。根據(jù)官方文檔進(jìn)行必要的配置,特別是 instance.properties 文件中的數(shù)據(jù)庫(kù)連接信息。

3. 安裝RabbitMQ

可以通過 Docker 快速安裝 RabbitMQ:

docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:management

訪問 http://localhost:15672 登錄管理界面,默認(rèn)用戶名/密碼為 guest/guest。

整合SpringBoot與Canal實(shí)現(xiàn)客戶端

創(chuàng)建SpringBoot項(xiàng)目

使用 Spring Initializr 創(chuàng)建一個(gè)新的 Spring Boot 項(xiàng)目,添加 Web, JPA, 和 AMQP(用于后續(xù)整合 RabbitMQ)依賴。

引入Canal依賴

在 pom.xml 中添加 Canal Client 的依賴:

<dependency>
    <groupId>com.alibaba.otter</groupId>
    <artifactId>canal.client</artifactId>
    <version>1.1.5</version>
</dependency>

編寫Canal客戶端代碼

創(chuàng)建一個(gè) Canal 客戶端類,用來(lái)監(jiān)聽 MySQL 數(shù)據(jù)庫(kù)的變化,并將變更事件轉(zhuǎn)發(fā)給 RabbitMQ。

import com.alibaba.otter.canal.client.CanalConnector;
import com.alibaba.otter.canal.protocol.CanalEntry;
import com.rabbitmq.client.Channel;

public class CanalClient {

    private final CanalConnector connector;
    private final Channel channel;

    public CanalClient(CanalConnector connector, Channel channel) {
        this.connector = connector;
        this.channel = channel;
    }

    public void start() {
        // Canal 連接配置
        connector.connect();
        connector.subscribe(".*\\..*"); // 訂閱所有數(shù)據(jù)庫(kù)和表
        connector.rollback();

        while (true) {
            int batchSize = 1000;
            EntryBatch batch = connector.getWithoutAck(batchSize); // 獲取一批次數(shù)據(jù)
            long batchId = batch.getId();
            int size = batch.getEntries().size();

            if (batchId == -1 || size == 0) {
                try {
                    Thread.sleep(1000);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            } else {
                printEntry(batch.getEntries());
                connector.ack(batchId); // 提交確認(rèn)
            }

            if (Thread.currentThread().isInterrupted()) {
                break;
            }
        }

        connector.disconnect();
    }

    private void printEntry(List<Entry> entrys) {
        for (Entry entry : entrys) {
            if (entry.getEntryType() == EntryType.TRANSACTIONBEGIN || entry.getEntryType() == EntryType.TRANSACTIONEND) {
                continue;
            }

            RowChange rowChage = null;
            try {
                rowChage = RowChange.parseFrom(entry.getStoreValue());
            } catch (Exception e) {
                throw new RuntimeException("ERROR ## parser of eromanga-event has an error , data:" + entry.toString(), e);
            }

            EventType eventType = rowChage.getEventType();
            System.out.println(String.format("================&gt; binlog[%s:%s] , name[%s,%s] , eventType : %s",
                    entry.getHeader().getLogfileName(), entry.getHeader().getLogfileOffset(),
                    entry.getHeader().getSchemaName(), entry.getHeader().getTableName(),
                    eventType));

            for (RowData rowData : rowChage.getRowDatasList()) {
                if (eventType == EventType.DELETE) {
                    sendToRabbitMQ(rowData.getBeforeColumnsList());
                } else if (eventType == EventType.INSERT) {
                    sendToRabbitMQ(rowData.getAfterColumnsList());
                } else {
                    System.out.println("-------> before");
                    sendToRabbitMQ(rowData.getBeforeColumnsList());

                    System.out.println("-------> after");
                    sendToRabbitMQ(rowData.getAfterColumnsList());
                }
            }
        }
    }

    private void sendToRabbitMQ(List<Column> columns) {
        StringBuilder message = new StringBuilder();
        for (Column column : columns) {
            message.append(column.getName()).append("=").append(column.getValue()).append(",");
        }
        try {
            channel.basicPublish("", "canal_exchange", null, message.toString().getBytes());
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
}

Canal整合RabbitMQ

配置Canal Server

確保 Canal Server 已正確配置并啟動(dòng),能夠監(jiān)聽 MySQL 的 binlog 日志。修改 Canal Server 的配置文件以指向你的 MySQL 實(shí)例,并設(shè)置適當(dāng)?shù)倪^濾規(guī)則。

配置RabbitMQ Exchange

在 RabbitMQ 中創(chuàng)建一個(gè)名為 canal_exchange 的 exchange,類型可以根據(jù)需要選擇,如 fanout, direct, topic 或 headers。

rabbitmqadmin declare exchange name=canal_exchange type=fanout

SpringBoot整合RabbitMQ

添加依賴

確保在 pom.xml 中已經(jīng)包含了 RabbitMQ 的 Spring AMQP 依賴。

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

配置RabbitMQ連接信息

在 application.yml 或 application.properties 中配置 RabbitMQ 的連接參數(shù)。

spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest

創(chuàng)建消費(fèi)者

編寫一個(gè)消費(fèi)者類來(lái)接收來(lái)自 RabbitMQ 的消息。

import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
public class CanalMessageConsumer {

    @RabbitListener(queues = "canal_queue")
    public void receive(String message) {
        System.out.println("Received message: " + message);
    }
}

配置隊(duì)列和綁定

確保在應(yīng)用程序啟動(dòng)時(shí)自動(dòng)創(chuàng)建所需的隊(duì)列,并將它們綁定到之前創(chuàng)建的 exchange 上。

import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.TopicExchange;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class RabbitConfig {

    @Bean
    public Queue canalQueue() {
        return new Queue("canal_queue", false);
    }

    @Bean
    public TopicExchange canalExchange() {
        return new TopicExchange("canal_exchange");
    }

    @Bean
    public Binding binding(Queue canalQueue, TopicExchange canalExchange) {
        return BindingBuilder.bind(canalQueue).to(canalExchange).with("#");
    }
}

總結(jié)

通過以上步驟,我們成功地將 Canal 與 RabbitMQ 整合到了 Spring Boot 應(yīng)用程序中。這使得我們可以實(shí)時(shí)監(jiān)聽 MySQL 數(shù)據(jù)庫(kù)的變更,并將這些變更作為消息發(fā)布到 RabbitMQ 中供其他微服務(wù)消費(fèi)。這種方法不僅提高了系統(tǒng)的響應(yīng)速度,也簡(jiǎn)化了數(shù)據(jù)同步的過程,降低了開發(fā)和維護(hù)成本。

到此這篇關(guān)于SpringBoot整合Canal+RabbitMQ監(jiān)聽數(shù)據(jù)變更詳解的文章就介紹到這了,更多相關(guān)SpringBoot監(jiān)聽數(shù)據(jù)變更內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java數(shù)據(jù)結(jié)構(gòu)之圖的兩種搜索算法詳解

    Java數(shù)據(jù)結(jié)構(gòu)之圖的兩種搜索算法詳解

    在很多情況下,我們需要遍歷圖,得到圖的一些性質(zhì)。有關(guān)圖的搜索,最經(jīng)典的算法有深度優(yōu)先搜索和廣度優(yōu)先搜索,接下來(lái)我們分別講解這兩種搜索算法,需要的可以參考一下
    2022-11-11
  • 使用SpringBoot?+?Vue?+?Redis實(shí)現(xiàn)驗(yàn)證碼登錄功能全過程

    使用SpringBoot?+?Vue?+?Redis實(shí)現(xiàn)驗(yàn)證碼登錄功能全過程

    在現(xiàn)代web應(yīng)用中,用戶驗(yàn)證是非常重要的一部分,這篇文章主要介紹了使用SpringBoot?+?Vue?+?Redis實(shí)現(xiàn)驗(yàn)證碼登錄功能的相關(guān)資料,文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2025-10-10
  • Java如何確定兩個(gè)區(qū)間范圍是否有交集

    Java如何確定兩個(gè)區(qū)間范圍是否有交集

    這篇文章主要介紹了Java如何確定兩個(gè)區(qū)間范圍是否有交集問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2023-12-12
  • Mybatis-plus常見的坑@TableField不生效問題

    Mybatis-plus常見的坑@TableField不生效問題

    這篇文章主要介紹了Mybatis-plus常見的坑@TableField不生效問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-01-01
  • java自定義攔截器用法實(shí)例

    java自定義攔截器用法實(shí)例

    這篇文章主要介紹了java自定義攔截器用法,實(shí)例分析了java自定義攔截器的實(shí)現(xiàn)與使用技巧,需要的朋友可以參考下
    2015-06-06
  • 深度解析Java常量池中的Integer緩沖池和String常量池

    深度解析Java常量池中的Integer緩沖池和String常量池

    為了減少對(duì)象重復(fù)創(chuàng)建、提升運(yùn)行時(shí)效率,Java 內(nèi)部提供了兩種重要的優(yōu)化機(jī)制Integer 緩沖池(IntegerCache)和 String 常量池(String Pool),本文將深入剖析兩大常量池的底層實(shí)現(xiàn)、工作流程、適用范圍,并通過流程圖和代碼示例幫助你徹底掌握
    2026-05-05
  • Spring動(dòng)態(tài)管理定時(shí)任務(wù)之ThreadPoolTaskScheduler解讀

    Spring動(dòng)態(tài)管理定時(shí)任務(wù)之ThreadPoolTaskScheduler解讀

    這篇文章主要介紹了Spring動(dòng)態(tài)管理定時(shí)任務(wù)之ThreadPoolTaskScheduler解讀,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-12-12
  • Java中的XML解析技術(shù)詳析

    Java中的XML解析技術(shù)詳析

    XML文檔是一個(gè)文檔樹,從根部開始,并擴(kuò)展到樹的最底部,下面這篇文章主要給大家介紹了關(guān)于Java中XML解析技術(shù)的相關(guān)資料,文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2024-08-08
  • java 中基本算法之希爾排序的實(shí)例詳解

    java 中基本算法之希爾排序的實(shí)例詳解

    這篇文章主要介紹了java 中基本算法之希爾排序的實(shí)例詳解的相關(guān)資料,這里提供簡(jiǎn)單實(shí)現(xiàn)的實(shí)例,幫助大家學(xué)習(xí)理解此部分知識(shí),需要的朋友可以參考下
    2017-07-07
  • 一文簡(jiǎn)單了解C#?中的DataSet類

    一文簡(jiǎn)單了解C#?中的DataSet類

    這篇文章主要介紹了一文簡(jiǎn)單了解C#?中的DataSet類,文章圍繞主題展開詳細(xì)的內(nèi)容介紹,具有一定的參考價(jià)值,需要的小伙伴可以參考一下
    2022-08-08

最新評(píng)論

准格尔旗| 定州市| 柘荣县| 英德市| 虎林市| 易门县| 琼中| 盘锦市| 桃江县| 六盘水市| 西和县| 喀什市| 绿春县| 普安县| 五河县| 浙江省| 龙岩市| 原平市| 阿克| 防城港市| 南召县| 舒兰市| 工布江达县| 榆林市| 兰考县| 皮山县| 东明县| 芒康县| 宣汉县| 崇文区| 大田县| 湾仔区| 永新县| 栾城县| 肇源县| 枝江市| 楚雄市| 互助| 永清县| 中阳县| 赤峰市|