flink-cdc 讀取mysql數(shù)據(jù)

通過flink-cdc的Connector讀取mysql數(shù)據(jù)诈豌,并寫入到其他系統(tǒng)或者數(shù)據(jù)庫,需要先開啟mysql的binlog功能

1. 導(dǎo)入maven 依賴

 <dependencies>
    <dependency>
      <groupId>mysql</groupId>
      <artifactId>mysql-connector-java</artifactId>
      <version>8.0.19</version>
    </dependency>
 <dependency>
      <groupId>com.alibaba</groupId>
      <artifactId>fastjson</artifactId>
      <version>2.0.6</version>
    </dependency>
    <!-- 引入日志管理相關(guān)依賴-->
    <dependency>
      <groupId>org.slf4j</groupId>
      <artifactId>slf4j-api</artifactId>
      <version>${slf4j.version}</version>
    </dependency>
    <dependency>
      <groupId>org.slf4j</groupId>
      <artifactId>slf4j-log4j12</artifactId>
      <version>${slf4j.version}</version>
    </dependency>
    <dependency>
      <groupId>org.apache.logging.log4j</groupId>
      <artifactId>log4j-to-slf4j</artifactId>
      <version>2.14.0</version>
    </dependency>
    <dependency>
      <groupId>org.projectlombok</groupId>
      <artifactId>lombok</artifactId>
      <version>1.16.18</version>
    </dependency>
    <dependency>
      <groupId>com.alibaba.otter</groupId>
      <artifactId>canal.client</artifactId>
      <version>1.1.6</version>
    </dependency>
    <dependency>
      <groupId>com.ververica</groupId>
      <artifactId>flink-connector-mysql-cdc</artifactId>
      <version>2.2.1</version>
    </dependency>
    <dependency>
      <groupId>org.apache.flink</groupId>
      <artifactId>flink-streaming-java_2.12</artifactId>
      <version>1.13.6</version>
    </dependency>
    <dependency>
      <groupId>org.apache.flink</groupId>
      <artifactId>flink-clients_2.12</artifactId>
      <version>1.13.6</version>
    </dependency>
    <dependency>
      <groupId>org.apache.flink</groupId>
      <artifactId>flink-java</artifactId>
      <version>1.13.6</version>
    </dependency>
    <dependency>
      <groupId>org.apache.flink</groupId>
      <artifactId>flink-table-planner-blink_2.12</artifactId>
      <version>1.13.6</version>
      <type>test-jar</type>
    </dependency>
  </dependencies>

2. 新建Flink-cdc測(cè)試類

import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.connectors.mysql.table.StartupOptions;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.ArrayUtils;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.runtime.state.filesystem.FsStateBackend;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

@Slf4j
public class FlinkCDC {


    public static void main(String[] args) throws Exception {
        MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
            .hostname("127.0.0.1")
            .port(3306)
            .databaseList("user") // set captured database
            .tableList("user.log_info") // set captured table
            .username("root")
            .password("password")
             // 自定義反序列化方式
            .deserializer(new CustomDeserialization())
             //           StartupOptions.initial() 先全量后增量
//            .startupOptions(StartupOptions.initial())
              //   StartupOptions.latest() 從最新binlog讀取仆救,增量方式
            .startupOptions(StartupOptions.latest())
            .build();


        Configuration config = new Configuration();

//        config.setString("execution.savepoint.path", "file:///D:\\flink\\checkpoints\\cc52b93fd24977e5388f0a19a30d49d2\\chk-87");
        // 啟動(dòng)時(shí)設(shè)置
        if(ArrayUtils.isNotEmpty(args)) {
             String lasCheckpointPath =  args[0];
             // 例如 D:\flink\checkpoints\8bdf5d49bb1b4cda56aaa0a590fc2cef\chk-55
            // 重啟服務(wù)器指定最新checkpoint路徑,從該路徑指定checkpoint位置恢復(fù)讀取數(shù)據(jù)
             config.setString("execution.savepoint.path", lasCheckpointPath);
        }
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);

        // enable checkpoint
        env.enableCheckpointing(3000);
//      env.getCheckpointConfig().setCheckpointStorage("file:///D:\\flink\\checkpoints");
// 設(shè)置checkpoint保存位置矫渔,這里設(shè)置為本地文件存儲(chǔ)
        env.setStateBackend(new FsStateBackend("file:///D:\\flink\\checkpoints"));
        env
            .fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
            // set 4 parallel source tasks
            .setParallelism(4)
            .addSink(new CustomSink()).setParallelism(1);
//            .print().setParallelism(1); // use parallelism 1 for sink to keep message ordering

        env.execute("flinkcdc");
    }
}

3.自定義反序列化類

import com.alibaba.fastjson2.JSONObject;
import com.ververica.cdc.debezium.DebeziumDeserializationSchema;
import io.debezium.data.Envelope;
import java.util.List;
import org.apache.flink.api.common.typeinfo.BasicTypeInfo;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.util.Collector;
import org.apache.kafka.connect.data.Field;
import org.apache.kafka.connect.data.Schema;
import org.apache.kafka.connect.data.Struct;
import org.apache.kafka.connect.source.SourceRecord;

/**
 * 自定義序列化器
 */
public class CustomDeserialization implements DebeziumDeserializationSchema<String> {



    @Override
    public void deserialize(SourceRecord sourceRecord, Collector<String> collector)
        throws Exception {

        JSONObject res = new JSONObject();

        // 獲取數(shù)據(jù)庫和表名稱
        String topic = sourceRecord.topic();
        String[] fields = topic.split("\\.");
        String database = fields[1];
        String tableName = fields[2];

        Struct value = (Struct) sourceRecord.value();
        // 獲取before數(shù)據(jù)
        Struct before = value.getStruct("before");
        JSONObject beforeJson = new JSONObject();
        if (before != null) {
            Schema beforeSchema = before.schema();
            List<Field> beforeFields = beforeSchema.fields();
            for (Field field : beforeFields) {
                Object beforeValue = before.get(field);
                beforeJson.put(field.name(), beforeValue);
            }
        }

        // 獲取after數(shù)據(jù)
        Struct after = value.getStruct("after");
        JSONObject afterJson = new JSONObject();
        if (after != null) {
            Schema afterSchema = after.schema();
            List<Field> afterFields = afterSchema.fields();
            for (Field field : afterFields) {
                Object afterValue = after.get(field);
                afterJson.put(field.name(), afterValue);
            }
        }

        //獲取操作類型 READ DELETE UPDATE CREATE
        Envelope.Operation operation = Envelope.operationFor(sourceRecord);
        String type = operation.toString().toLowerCase();
        if ("create".equals(type)) {
            type = "insert";
        }

        // 將字段寫到j(luò)son對(duì)象中
        res.put("database", database);
        res.put("tableName", tableName);
        res.put("before", beforeJson);
        res.put("after", afterJson);
        res.put("type", type);

        //輸出數(shù)據(jù)
        collector.collect(res.toString());
    }

    @Override
    public TypeInformation<String> getProducedType() {
        return BasicTypeInfo.STRING_TYPE_INFO;
    }
}

4.自定義Sink輸出

import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;

public class CustomSink extends RichSinkFunction {


    @Override
    public void invoke(Object value,Context context) throws Exception {
        String v = value.toString();

         TableData<LogInfo>  tableData = JSON.parseObject(v,  new 
            TypeReference<TableData<LogInfo>>() {});
           System.out.println(t.toString());
            // TODO 保存到其他系統(tǒng)/中間件(mq等)/其他數(shù)據(jù)庫,同學(xué)可以自己根據(jù)情況實(shí)現(xiàn)
            //  發(fā)送到消息隊(duì)列rabbitmq或者kafaka中處理
           // rabbitmqtemplate.send(tableData );
          //  保存數(shù)據(jù)庫
          //  testMapper.insert(tableData.getAfter());
    }
}


@Data
@NoArgsConstructor
@AllArgsConstructor
public class TableData<T> {

    private String database;

    private String tableName;

    private String update;

    private T before;

    private T after;
}

最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
  • 序言:七十年代末彤蔽,一起剝皮案震驚了整個(gè)濱河市,隨后出現(xiàn)的幾起案子庙洼,更是在濱河造成了極大的恐慌顿痪,老刑警劉巖,帶你破解...
    沈念sama閱讀 218,755評(píng)論 6 507
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件油够,死亡現(xiàn)場(chǎng)離奇詭異蚁袭,居然都是意外死亡,警方通過查閱死者的電腦和手機(jī)石咬,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 93,305評(píng)論 3 395
  • 文/潘曉璐 我一進(jìn)店門揩悄,熙熙樓的掌柜王于貴愁眉苦臉地迎上來,“玉大人碌补,你說我怎么就攤上這事虏束∶奕模” “怎么了厦章?”我有些...
    開封第一講書人閱讀 165,138評(píng)論 0 355
  • 文/不壞的土叔 我叫張陵,是天一觀的道長照藻。 經(jīng)常有香客問我袜啃,道長,這世上最難降的妖魔是什么幸缕? 我笑而不...
    開封第一講書人閱讀 58,791評(píng)論 1 295
  • 正文 為了忘掉前任群发,我火速辦了婚禮,結(jié)果婚禮上发乔,老公的妹妹穿的比我還像新娘熟妓。我一直安慰自己,他們只是感情好栏尚,可當(dāng)我...
    茶點(diǎn)故事閱讀 67,794評(píng)論 6 392
  • 文/花漫 我一把揭開白布起愈。 她就那樣靜靜地躺著,像睡著了一般。 火紅的嫁衣襯著肌膚如雪抬虽。 梳的紋絲不亂的頭發(fā)上官觅,一...
    開封第一講書人閱讀 51,631評(píng)論 1 305
  • 那天,我揣著相機(jī)與錄音阐污,去河邊找鬼休涤。 笑死,一個(gè)胖子當(dāng)著我的面吹牛笛辟,可吹牛的內(nèi)容都是我干的功氨。 我是一名探鬼主播,決...
    沈念sama閱讀 40,362評(píng)論 3 418
  • 文/蒼蘭香墨 我猛地睜開眼手幢,長吁一口氣:“原來是場(chǎng)噩夢(mèng)啊……” “哼疑故!你這毒婦竟也來了?” 一聲冷哼從身側(cè)響起弯菊,我...
    開封第一講書人閱讀 39,264評(píng)論 0 276
  • 序言:老撾萬榮一對(duì)情侶失蹤纵势,失蹤者是張志新(化名)和其女友劉穎,沒想到半個(gè)月后管钳,有當(dāng)?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體钦铁,經(jīng)...
    沈念sama閱讀 45,724評(píng)論 1 315
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡,尸身上長有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 37,900評(píng)論 3 336
  • 正文 我和宋清朗相戀三年才漆,在試婚紗的時(shí)候發(fā)現(xiàn)自己被綠了牛曹。 大學(xué)時(shí)的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片。...
    茶點(diǎn)故事閱讀 40,040評(píng)論 1 350
  • 序言:一個(gè)原本活蹦亂跳的男人離奇死亡醇滥,死狀恐怖黎比,靈堂內(nèi)的尸體忽然破棺而出,到底是詐尸還是另有隱情鸳玩,我是刑警寧澤阅虫,帶...
    沈念sama閱讀 35,742評(píng)論 5 346
  • 正文 年R本政府宣布,位于F島的核電站不跟,受9級(jí)特大地震影響颓帝,放射性物質(zhì)發(fā)生泄漏。R本人自食惡果不足惜窝革,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 41,364評(píng)論 3 330
  • 文/蒙蒙 一购城、第九天 我趴在偏房一處隱蔽的房頂上張望。 院中可真熱鬧虐译,春花似錦瘪板、人聲如沸。這莊子的主人今日做“春日...
    開封第一講書人閱讀 31,944評(píng)論 0 22
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽史侣。三九已至,卻和暖如春魏身,著一層夾襖步出監(jiān)牢的瞬間惊橱,已是汗流浹背。 一陣腳步聲響...
    開封第一講書人閱讀 33,060評(píng)論 1 270
  • 我被黑心中介騙來泰國打工箭昵, 沒想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留税朴,地道東北人。 一個(gè)月前我還...
    沈念sama閱讀 48,247評(píng)論 3 371
  • 正文 我出身青樓家制,卻偏偏與公主長得像正林,于是被迫代替她去往敵國和親。 傳聞我的和親對(duì)象是個(gè)殘疾皇子颤殴,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 44,979評(píng)論 2 355

推薦閱讀更多精彩內(nèi)容