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

SpringBoot整合Flink CDC實(shí)現(xiàn)實(shí)時(shí)追蹤mysql數(shù)據(jù)變動(dòng)

 更新時(shí)間:2024年07月24日 09:14:13   作者:碼到三十五  
我們將整合Spring Boot和Apache Flink CDC(Change Data Capture)來(lái)實(shí)現(xiàn)實(shí)時(shí)數(shù)據(jù)追蹤,下面是一個(gè)基本的實(shí)踐流程代碼,包括搭建Spring Boot項(xiàng)目、整合Flink CDC以及實(shí)現(xiàn)數(shù)據(jù)變動(dòng)的實(shí)時(shí)追蹤,需要的朋友可以參考下

前言

Flink CDC(Flink Change Data Capture)是一種基于數(shù)據(jù)庫(kù)日志的CDC技術(shù),它實(shí)現(xiàn)了一個(gè)全增量一體化的數(shù)據(jù)集成框架。與Flink計(jì)算框架相結(jié)合,F(xiàn)link CDC能夠高效地實(shí)現(xiàn)海量數(shù)據(jù)的實(shí)時(shí)集成。其核心功能在于實(shí)時(shí)監(jiān)視數(shù)據(jù)庫(kù)或數(shù)據(jù)流中的數(shù)據(jù)變動(dòng),并將這些變動(dòng)抽取出來(lái),以便進(jìn)行進(jìn)一步的處理和分析。借助Flink CDC,用戶可以輕松地構(gòu)建實(shí)時(shí)數(shù)據(jù)管道,實(shí)時(shí)響應(yīng)和處理數(shù)據(jù)變動(dòng),為實(shí)時(shí)分析、實(shí)時(shí)報(bào)表和實(shí)時(shí)決策等場(chǎng)景提供有力支持。

Flink CDC的應(yīng)用場(chǎng)景廣泛,包括但不限于實(shí)時(shí)數(shù)據(jù)倉(cāng)庫(kù)更新、實(shí)時(shí)數(shù)據(jù)同步和遷移以及實(shí)時(shí)數(shù)據(jù)處理等。它還能確保數(shù)據(jù)一致性,并在數(shù)據(jù)發(fā)生變更時(shí)準(zhǔn)確地進(jìn)行捕獲和處理。此外,F(xiàn)link CDC支持與多種數(shù)據(jù)源進(jìn)行集成,如MySQL、PostgreSQL、Oracle等,并提供了相應(yīng)的連接器,便于數(shù)據(jù)的捕獲和處理。

接下來(lái),將詳細(xì)介紹MySQL CDC的使用。MySQL CDC連接器允許從MySQL數(shù)據(jù)庫(kù)中讀取快照數(shù)據(jù)和增量數(shù)據(jù)。

1. MySQL開啟Binlog

MySQL中開啟binlog功能,需要修改配置文件中(如Linux的/etc/my.cnf或Windows的\my.ini)的[mysqld]部分設(shè)置相關(guān)參數(shù):

[mysqld]
server-id=1
# 設(shè)置日志格式為行級(jí)格式
binlog-format=Row
# 設(shè)置binlog日志文件的前綴
log-bin=mysql-bin
# 指定需要記錄二進(jìn)制日志的數(shù)據(jù)庫(kù)
binlog_do_db=testjpa

除了開啟binlog功能外,還需要為Flink CDC配置相應(yīng)的權(quán)限,以確保其能夠正常連接到MySQL并讀取數(shù)據(jù)。這包括授予Flink CDC連接MySQL的用戶必要的權(quán)限,如SELECT、REPLICATION SLAVE、REPLICATION CLIENT、SHOW VIEW等。這些權(quán)限是Flink CDC讀取數(shù)據(jù)和元數(shù)據(jù)所必需的。

檢查是否已開啟binlog功能:

mysql> SHOW VARIABLES LIKE 'log_bin';
+---------------+-------+
| Variable_name | Value |
+---------------+-------+
| log_bin       | ON    |
+---------------+-------+

至此,MySQL的相關(guān)配置已完成。

2. 創(chuàng)建Spring Boot項(xiàng)目

首先,你需要?jiǎng)?chuàng)建一個(gè)Spring Boot項(xiàng)目??梢允褂肧pring Initializr(https://start.spring.io/)來(lái)快速生成項(xiàng)目。

3. 添加依賴

pom.xml中添加Apache Flink和Flink CDC的依賴。以下是必要的依賴:

<dependencies>
    <!-- Flink dependency -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-java</artifactId>
        <version>1.14.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java_2.12</artifactId>
        <version>1.14.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-connector-mysql-cdc</artifactId>
        <version>2.0.0</version>
    </dependency>
    <!-- Spring Boot dependencies -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter</artifactId>
    </dependency>
</dependencies>

4. 配置Flink和MySQL CDC

在Spring Boot的application.ymlapplication.properties文件中配置Flink和MySQL數(shù)據(jù)庫(kù)連接:

flink:
  checkpoint:
    interval: 10000
  parallelism: 1

spring:
  datasource:
    url: jdbc:mysql://localhost:3306/your_database
    username: your_username
    password: your_password

5. 實(shí)現(xiàn)數(shù)據(jù)實(shí)時(shí)追蹤

創(chuàng)建一個(gè)服務(wù)類來(lái)實(shí)現(xiàn)數(shù)據(jù)的實(shí)時(shí)追蹤:

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.springframework.stereotype.Service;

@Service
public class FlinkCdcService {

    public void startDataStreaming() {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        final StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);

        // 使用Flink CDC連接MySQL
        String name = "inventory";
        tableEnv.executeSql("CREATE TABLE " + name + " (" +
            "  id INT," +
            "  name STRING," +
            "  description STRING," +
            "  weight DECIMAL(10, 3)" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'localhost'," +
            "  'port' = '3306'," +
            "  'username' = 'your_username'," +
            "  'password' = 'your_password'," +
            "  'database-name' = 'your_database'," +
            "  'table-name' = 'your_table'" +
            ")");

        // 查詢并打印結(jié)果
        DataStream<String> dataStream = tableEnv.sqlQuery("SELECT * FROM " + name).execute().print();

        try {
            env.execute("Flink CDC Demo");
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

6. 啟動(dòng)Spring Boot應(yīng)用

在你的Spring Boot應(yīng)用的啟動(dòng)類中調(diào)用FlinkCdcServicestartDataStreaming方法來(lái)啟動(dòng)數(shù)據(jù)追蹤:

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.CommandLineRunner;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

@SpringBootApplication
public class FlinkCdcApplication implements CommandLineRunner {

    @Autowired
    private FlinkCdcService flinkCdcService;

    public static void main(String[] args) {
        SpringApplication.run(FlinkCdcApplication.class, args);
    }

    @Override
    public void run(String... args) throws Exception {
        flinkCdcService.startDataStreaming();
    }
}

7. 運(yùn)行并測(cè)試

運(yùn)行Spring Boot應(yīng)用,并在MySQL數(shù)據(jù)庫(kù)中做出一些數(shù)據(jù)變動(dòng)。你應(yīng)該能在控制臺(tái)看到實(shí)時(shí)打印的數(shù)據(jù)變動(dòng)。

到此這篇關(guān)于SpringBoot整合Flink CDC實(shí)現(xiàn)實(shí)時(shí)追蹤mysql數(shù)據(jù)變動(dòng)的文章就介紹到這了,更多相關(guān)SpringBoot Flink CDC mysql數(shù)據(jù)變動(dòng)內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • JAVA數(shù)字千分位和小數(shù)點(diǎn)的現(xiàn)實(shí)代碼(處理金額問(wèn)題)

    JAVA數(shù)字千分位和小數(shù)點(diǎn)的現(xiàn)實(shí)代碼(處理金額問(wèn)題)

    這篇文章主要介紹了JAVA數(shù)字千分位和小數(shù)點(diǎn)的現(xiàn)實(shí)代碼(處理金額問(wèn)題),具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2020-10-10
  • SpringBoot啟動(dòng)時(shí)將數(shù)據(jù)庫(kù)數(shù)據(jù)預(yù)加載到Redis緩存的幾種實(shí)現(xiàn)方案

    SpringBoot啟動(dòng)時(shí)將數(shù)據(jù)庫(kù)數(shù)據(jù)預(yù)加載到Redis緩存的幾種實(shí)現(xiàn)方案

    在實(shí)際項(xiàng)目開發(fā)中,我們經(jīng)常需要在應(yīng)用啟動(dòng)時(shí)將一些固定的、頻繁訪問(wèn)的數(shù)據(jù)從數(shù)據(jù)庫(kù)預(yù)加載到 Redis 緩存中,以提高系統(tǒng)性能,本文將介紹幾種實(shí)現(xiàn)方案,需要的朋友可以參考下
    2025-09-09
  • JAVA?ServLet創(chuàng)建一個(gè)項(xiàng)目的基本步驟

    JAVA?ServLet創(chuàng)建一個(gè)項(xiàng)目的基本步驟

    Servlet是Server Applet的簡(jiǎn)稱,是運(yùn)行在服務(wù)器上的小程序,用于編寫Java的服務(wù)器端程序,它的主要作用是接收并響應(yīng)來(lái)自Web客戶端的請(qǐng)求,下面這篇文章主要給大家介紹了關(guān)于JAVA?ServLet創(chuàng)建一個(gè)項(xiàng)目的基本步驟,需要的朋友可以參考下
    2024-03-03
  • Java線程(Thread)四種停止方式代碼實(shí)例

    Java線程(Thread)四種停止方式代碼實(shí)例

    這篇文章主要介紹了Java線程(Thread)四種停止方式代碼實(shí)例,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-03-03
  • SpringBoot發(fā)送異步郵件流程與實(shí)現(xiàn)詳解

    SpringBoot發(fā)送異步郵件流程與實(shí)現(xiàn)詳解

    這篇文章主要介紹了SpringBoot發(fā)送異步郵件流程與實(shí)現(xiàn)詳解,Servlet階段郵件發(fā)送非常的復(fù)雜,如果現(xiàn)代化的Java開發(fā)是那個(gè)樣子該有多糟糕,現(xiàn)在SpringBoot中集成好了郵件發(fā)送的東西,而且操作十分簡(jiǎn)單容易上手,需要的朋友可以參考下
    2024-01-01
  • 一起來(lái)看看springboot集成redis的使用注解

    一起來(lái)看看springboot集成redis的使用注解

    這篇文章主要為大家詳細(xì)介紹了springboot集成redis的使用注解,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下,希望能夠給你帶來(lái)幫助
    2022-03-03
  • SpringBoot自定義轉(zhuǎn)換器應(yīng)用實(shí)例講解

    SpringBoot自定義轉(zhuǎn)換器應(yīng)用實(shí)例講解

    SpringBoot在響應(yīng)客戶端請(qǐng)求時(shí),將提交的數(shù)據(jù)封裝成對(duì)象時(shí),使用了內(nèi)置的轉(zhuǎn)換器,SpringBoot 也支持自定義轉(zhuǎn)換器,這個(gè)內(nèi)置轉(zhuǎn)換器在 debug的時(shí)候,可以看到,提供了124個(gè)內(nèi)置轉(zhuǎn)換器
    2022-08-08
  • java中構(gòu)造器內(nèi)部調(diào)用構(gòu)造器實(shí)例詳解

    java中構(gòu)造器內(nèi)部調(diào)用構(gòu)造器實(shí)例詳解

    在本篇文章里小編給大家分享的是關(guān)于java中構(gòu)造器內(nèi)部調(diào)用構(gòu)造器實(shí)例內(nèi)容,需要的朋友們可以學(xué)習(xí)下。
    2020-05-05
  • Java調(diào)用第三方http接口的常用方式總結(jié)

    Java調(diào)用第三方http接口的常用方式總結(jié)

    這篇文章主要介紹了Java調(diào)用第三方http接口的常用方式總結(jié),具有很好的參考價(jià)值,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-06-06
  • Spring?請(qǐng)求之傳遞?JSON?數(shù)據(jù)的操作方法

    Spring?請(qǐng)求之傳遞?JSON?數(shù)據(jù)的操作方法

    JSON 就是一種數(shù)據(jù)格式,有自己的格式和語(yǔ)法,使用文本表示一個(gè)對(duì)象或數(shù)組的信息,因此 JSON 本質(zhì)是字符串,主要負(fù)責(zé)在不同的語(yǔ)言中數(shù)據(jù)傳遞和交換,這篇文章主要介紹了Spring請(qǐng)求之傳遞JSON數(shù)據(jù)的相關(guān)知識(shí),需要的朋友可以參考下
    2025-04-04

最新評(píng)論

房山区| 汽车| 金塔县| 大新县| 江门市| 连平县| 卓资县| 新安县| 翁牛特旗| 临江市| 濉溪县| 南陵县| 伊通| 镇康县| 丰台区| 东台市| 锡林浩特市| 广西| 田阳县| 安化县| 郴州市| 安义县| 文安县| 丹江口市| 鹤壁市| 云林县| 利川市| 千阳县| 峨边| 永川市| 钦州市| 横山县| 大丰市| 泾川县| 绥江县| 甘泉县| 西昌市| 普定县| 梅河口市| 甘肃省| 疏勒县|