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

如何使用Flink CDC實(shí)現(xiàn) Oracle數(shù)據(jù)庫數(shù)據(jù)同步

 更新時(shí)間:2024年08月21日 11:55:13   作者:shandongwill  
Flink CDC是一個(gè)基于流的數(shù)據(jù)集成工具,為用戶提供一套功能全面的編程接口API, 該工具使得用戶能夠以YAML 配置文件的形式實(shí)現(xiàn)數(shù)據(jù)庫同步,同時(shí)也提供了Flink CDC Source Connector API,本文給大家介紹使用Flink CDC實(shí)現(xiàn) Oracle數(shù)據(jù)庫數(shù)據(jù)同步的方法,感興趣的朋友一起看看吧

前言

Flink CDC 是一個(gè)基于流的數(shù)據(jù)集成工具,旨在為用戶提供一套功能更加全面的編程接口(API)。 該工具使得用戶能夠以 YAML 配置文件的形式實(shí)現(xiàn)數(shù)據(jù)庫同步,同時(shí)也提供了Flink CDC Source Connector API。 Flink CDC 在任務(wù)提交過程中進(jìn)行了優(yōu)化,并且增加了一些高級特性,如表結(jié)構(gòu)變更自動(dòng)同步(Schema Evolution)、數(shù)據(jù)轉(zhuǎn)換(Data Transformation)、整庫同步(Full Database Synchronization)以及 精確一次(Exactly-once)語義。
本文通過flink-connector-oracle-cdc來實(shí)現(xiàn)Oracle數(shù)據(jù)庫的數(shù)據(jù)同步。

一、開啟歸檔日志

1)數(shù)據(jù)庫服務(wù)器終端,使用sysdba角色連接數(shù)據(jù)庫

 sqlplus / as sysdba
或
sqlplus /nolog
CONNECT sys/password AS SYSDBA;

2)檢查歸檔日志是否開啟

archive log list;

(“Database log mode: No Archive Mode”,日志歸檔未開啟)
(“Database log mode: Archive Mode”,日志歸檔已開啟)
3)啟用歸檔日志

alter system set db_recovery_file_dest_size = 10G;
alter system set db_recovery_file_dest = '/opt/oracle/oradata/recovery_area' scope=spfile;
shutdown immediate;
startup mount;
alter database archivelog;
alter database open;

注意:
啟用歸檔日志需要重啟數(shù)據(jù)庫。
歸檔日志會占用大量的磁盤空間,應(yīng)定期清除過期的日志文件
4)啟動(dòng)完成后重新執(zhí)行 archive log list; 查看歸檔打開狀態(tài)

二、創(chuàng)建flinkcdc專屬用戶

2.1 對于Oracle 非CDB數(shù)據(jù)庫,執(zhí)行如下sql

  CREATE USER flinkuser IDENTIFIED BY flinkpw DEFAULT TABLESPACE LOGMINER_TBS QUOTA UNLIMITED ON LOGMINER_TBS;
  GRANT CREATE SESSION TO flinkuser;
  GRANT SET CONTAINER TO flinkuser;
  GRANT SELECT ON V_$DATABASE to flinkuser;
  GRANT FLASHBACK ANY TABLE TO flinkuser;
  GRANT SELECT ANY TABLE TO flinkuser;
  GRANT SELECT_CATALOG_ROLE TO flinkuser;
  GRANT EXECUTE_CATALOG_ROLE TO flinkuser;
  GRANT SELECT ANY TRANSACTION TO flinkuser;
  GRANT LOGMINING TO flinkuser;
  GRANT ANALYZE ANY TO flinkuser;
  GRANT CREATE TABLE TO flinkuser;
  -- need not to execute if set scan.incremental.snapshot.enabled=true(default)
  GRANT LOCK ANY TABLE TO flinkuser;
  GRANT ALTER ANY TABLE TO flinkuser;
  GRANT CREATE SEQUENCE TO flinkuser;
  GRANT EXECUTE ON DBMS_LOGMNR TO flinkuser;
  GRANT EXECUTE ON DBMS_LOGMNR_D TO flinkuser;
  GRANT SELECT ON V_$LOG TO flinkuser;
  GRANT SELECT ON V_$LOG_HISTORY TO flinkuser;
  GRANT SELECT ON V_$LOGMNR_LOGS TO flinkuser;
  GRANT SELECT ON V_$LOGMNR_CONTENTS TO flinkuser;
  GRANT SELECT ON V_$LOGMNR_PARAMETERS TO flinkuser;
  GRANT SELECT ON V_$LOGFILE TO flinkuser;
  GRANT SELECT ON V_$ARCHIVED_LOG TO flinkuser;
  GRANT SELECT ON V_$ARCHIVE_DEST_STATUS TO flinkuser;

2.2 對于Oracle CDB數(shù)據(jù)庫,執(zhí)行如下sql

  CREATE USER flinkuser IDENTIFIED BY flinkpw DEFAULT TABLESPACE logminer_tbs QUOTA UNLIMITED ON logminer_tbs CONTAINER=ALL;
  GRANT CREATE SESSION TO flinkuser CONTAINER=ALL;
  GRANT SET CONTAINER TO flinkuser CONTAINER=ALL;
  GRANT SELECT ON V_$DATABASE to flinkuser CONTAINER=ALL;
  GRANT FLASHBACK ANY TABLE TO flinkuser CONTAINER=ALL;
  GRANT SELECT ANY TABLE TO flinkuser CONTAINER=ALL;
  GRANT SELECT_CATALOG_ROLE TO flinkuser CONTAINER=ALL;
  GRANT EXECUTE_CATALOG_ROLE TO flinkuser CONTAINER=ALL;
  GRANT SELECT ANY TRANSACTION TO flinkuser CONTAINER=ALL;
  GRANT LOGMINING TO flinkuser CONTAINER=ALL;
  GRANT CREATE TABLE TO flinkuser CONTAINER=ALL;
  -- need not to execute if set scan.incremental.snapshot.enabled=true(default)
  GRANT LOCK ANY TABLE TO flinkuser CONTAINER=ALL;
  GRANT CREATE SEQUENCE TO flinkuser CONTAINER=ALL;
  GRANT EXECUTE ON DBMS_LOGMNR TO flinkuser CONTAINER=ALL;
  GRANT EXECUTE ON DBMS_LOGMNR_D TO flinkuser CONTAINER=ALL;
  GRANT SELECT ON V_$LOG TO flinkuser CONTAINER=ALL;
  GRANT SELECT ON V_$LOG_HISTORY TO flinkuser CONTAINER=ALL;
  GRANT SELECT ON V_$LOGMNR_LOGS TO flinkuser CONTAINER=ALL;
  GRANT SELECT ON V_$LOGMNR_CONTENTS TO flinkuser CONTAINER=ALL;
  GRANT SELECT ON V_$LOGMNR_PARAMETERS TO flinkuser CONTAINER=ALL;
  GRANT SELECT ON V_$LOGFILE TO flinkuser CONTAINER=ALL;
  GRANT SELECT ON V_$ARCHIVED_LOG TO flinkuser CONTAINER=ALL;
  GRANT SELECT ON V_$ARCHIVE_DEST_STATUS TO flinkuser CONTAINER=ALL;

三、指定oracle表、庫級啟用

-- 指定表啟用補(bǔ)充日志記錄:
ALTER TABLE databasename.tablename ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;
-- 為數(shù)據(jù)庫的所有表啟用
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;
-- 指定數(shù)據(jù)庫啟用補(bǔ)充日志記錄
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA;

四、使用flink-connector-oracle-cdc實(shí)現(xiàn)數(shù)據(jù)庫同步

4.1 引入pom依賴

 <dependency>
     <groupId>com.ververica</groupId>
     <artifactId>flink-connector-oracle-cdc</artifactId>
     <version>2.4.0</version>
 </dependency>

4.2 Java主代碼

package test.datastream.cdc.oracle;
import com.ververica.cdc.connectors.oracle.OracleSource;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.apache.flink.streaming.api.windowing.assigners.TumblingProcessingTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.types.Row;
import test.datastream.cdc.oracle.function.CacheDataAllWindowFunction;
import test.datastream.cdc.oracle.function.CdcString2RowMap;
import test.datastream.cdc.oracle.function.DbCdcSinkFunction;
import java.util.Properties;
public class OracleCdcExample {
    public static void main(String[] args) throws Exception {
        Properties properties = new Properties();
        //數(shù)字類型數(shù)據(jù) 轉(zhuǎn)換為字符
        properties.setProperty("decimal.handling.mode", "string");
        SourceFunction<String> sourceFunction = OracleSource.<String>builder()
//                .startupOptions(StartupOptions.latest()) // 從最晚位點(diǎn)啟動(dòng)
                .url("jdbc:oracle:thin:@localhost:1521:orcl")
                .port(1521)
                .database("ORCL") // monitor XE database
                .schemaList("c##flink_user") // monitor inventory schema
                .tableList("c##flink_user.TEST2") // monitor products table
                .username("c##flink_user")
                .password("flinkpw")
                .debeziumProperties(properties)
                .deserializer(new JsonDebeziumDeserializationSchema()) // converts SourceRecord to JSON String
                .build();
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        DataStreamSource<String> source = env.addSource(sourceFunction).setParallelism(1);// use parallelism 1 for sink to keep message ordering
        SingleOutputStreamOperator<Row> mapStream = source.flatMap(new CdcString2RowMap());
        SingleOutputStreamOperator<Row[]> winStream = mapStream.windowAll(TumblingProcessingTimeWindows.of(Time.seconds(5)))
                .process(new CacheDataAllWindowFunction());
		//批量同步
        winStream.addSink(new DbCdcSinkFunction(null));
        env.execute();
    }
}

4.3json轉(zhuǎn)換為row

package test.datastream.cdc.oracle.function;
import cn.com.victorysoft.common.configuration.VsConfiguration;
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.types.Row;
import org.apache.flink.types.RowKind;
import org.apache.flink.util.Collector;
import test.datastream.cdc.CdcConstants;
import java.sql.Timestamp;
import java.util.HashMap;
import java.util.Map;
import java.util.Set;
/**
 * @desc cdc json解析,并轉(zhuǎn)換為Row
 */
public class CdcString2RowMap extends RichFlatMapFunction<String, Row> {
    private Map<String,Integer> columnMap =new HashMap<>();
    @Override
    public void open(Configuration parameters) throws Exception {
        columnMap.put("ID",0);
        columnMap.put("NAME",1);
        columnMap.put("DESCRIPTION",2);
        columnMap.put("AGE",3);
        columnMap.put("CREATE_TIME",4);
        columnMap.put("SCORE",5);
        columnMap.put("C_1",6);
        columnMap.put("B_1",7);
    }
    @Override
    public void flatMap(String s, Collector<Row> collector) throws Exception {
        System.out.println("receive: "+s);
        VsConfiguration conf=VsConfiguration.from(s);
        String op = conf.getString(CdcConstants.K_OP);
        VsConfiguration before = conf.getConfiguration(CdcConstants.K_BEFORE);
        VsConfiguration after = conf.getConfiguration(CdcConstants.K_AFTER);
        Row row =null;
        if(CdcConstants.OP_C.equals(op)){
            //插入,使用after數(shù)據(jù)
            row = convertToRow(after);
            row.setKind(RowKind.INSERT);
        }else if(CdcConstants.OP_U.equals(op)){
            //更新,使用after數(shù)據(jù)
            row = convertToRow(after);
            row.setKind(RowKind.UPDATE_AFTER);
        }else if(CdcConstants.OP_D.equals(op)){
            //刪除,使用before數(shù)據(jù)
            row = convertToRow(before);
            row.setKind(RowKind.DELETE);
        }else {
            //r 操作,使用after數(shù)據(jù)
            row = convertToRow(after);
            row.setKind(RowKind.INSERT);
        }
        collector.collect(row);
    }
    private Row convertToRow(VsConfiguration data){
        Set<String> keys = data.getKeys();
        int size = keys.size();
        Row row=new Row(8);
        int i=0;
        for (String key:keys) {
            Integer index = this.columnMap.get(key);
            Object value=data.get(key);
            if(key.equals("CREATE_TIME")){
                //long日期轉(zhuǎn)timestamp
                value=long2Timestamp((Long)value);
            }
            row.setField(index,value);
        }
        return row;
    }
    private static  java.sql.Timestamp long2Timestamp(Long time){
        Timestamp timestamp = new Timestamp(time/1000);
        System.out.println(timestamp);
        return timestamp;
    }
}

到此這篇關(guān)于使用Flink CDC實(shí)現(xiàn) Oracle數(shù)據(jù)庫數(shù)據(jù)同步的文章就介紹到這了,更多相關(guān)Flink CDC Oracle數(shù)據(jù)同步內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 優(yōu)化Oracle停機(jī)時(shí)間及數(shù)據(jù)庫恢復(fù)

    優(yōu)化Oracle停機(jī)時(shí)間及數(shù)據(jù)庫恢復(fù)

    優(yōu)化Oracle停機(jī)時(shí)間及數(shù)據(jù)庫恢復(fù)...
    2007-03-03
  • Oracle存儲過程與函數(shù)的詳細(xì)使用教程

    Oracle存儲過程與函數(shù)的詳細(xì)使用教程

    存儲過程和函數(shù)在Oracle中被稱為子程序,是指被命名的PL/SQL塊,這種塊可以帶有參數(shù),可以被多次調(diào)用,下面這篇文章主要給大家介紹了關(guān)于Oracle存儲過程與函數(shù)的詳細(xì)使用,文中通過實(shí)例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2022-07-07
  • Informatica bulk與normal模式的深入詳解

    Informatica bulk與normal模式的深入詳解

    本篇文章是對Informatica bulk與normal模式進(jìn)行了詳細(xì)的分析介紹,需要的朋友參考下
    2013-05-05
  • oracle給新項(xiàng)目建表實(shí)操

    oracle給新項(xiàng)目建表實(shí)操

    這篇文章主要介紹了oracle給新項(xiàng)目建表,文章圍繞oracle建表的相關(guān)資料展開內(nèi)容,需要的小伙伴可以參考一下,希望對你有所幫助
    2021-10-10
  • oracle 層次化查詢(行政區(qū)劃三級級聯(lián))

    oracle 層次化查詢(行政區(qū)劃三級級聯(lián))

    現(xiàn)在將上面的行政區(qū)劃按代碼分為三個(gè)級別:省(后四位為0)/市(后兩位為0)/縣,同時(shí)分別標(biāo)出他們的級別,這樣的話,便于后期根據(jù)不同的級別查詢。
    2009-07-07
  • Oracle如何查看impdp正在執(zhí)行的內(nèi)容

    Oracle如何查看impdp正在執(zhí)行的內(nèi)容

    這篇文章主要給大家介紹了關(guān)于Oracle如何查看impdp正在執(zhí)行的內(nèi)容的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對大家學(xué)習(xí)或者使用Oracle具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-06-06
  • Oracle數(shù)據(jù)庫升級或數(shù)據(jù)遷移方法研究

    Oracle數(shù)據(jù)庫升級或數(shù)據(jù)遷移方法研究

    本文詳細(xì)論述了oracle數(shù)據(jù)庫升級的升級前的準(zhǔn)備、升級過程和升級后的測試與調(diào)整工作,并對各種升級方法在多種操作系統(tǒng)平臺上作了測試。
    2016-07-07
  • 如何使用GDAL庫的ogr2ogr將GeoJSON數(shù)據(jù)導(dǎo)入到PostgreSql中

    如何使用GDAL庫的ogr2ogr將GeoJSON數(shù)據(jù)導(dǎo)入到PostgreSql中

    本文主要介紹了PyTorch中的masked_fill函數(shù)的基本知識和使用方法,masked_fill函數(shù)接受一個(gè)輸入張量和一個(gè)布爾掩碼作為主要參數(shù),掩碼的形狀必須與輸入張量相同,掩碼操作根據(jù)掩碼中的布爾值在輸出張量中填充指定的值或保留輸入張量中的值
    2024-10-10
  • Oracle Instr函數(shù)實(shí)例講解

    Oracle Instr函數(shù)實(shí)例講解

    instr函數(shù)為字符查找函數(shù),其功能是查找一個(gè)字符串在另一個(gè)字符串中首次出現(xiàn)的位置,instr函數(shù)在Oracle/PLSQL中是返回要截取的字符串在源字符串中的位置,這篇文章主要介紹了Oracle Instr函數(shù)實(shí)例講解,需要的朋友可以參考下
    2022-11-11
  • Oracle 12.2處理sysaux空間占滿問題

    Oracle 12.2處理sysaux空間占滿問題

    今天處理別的問題查看告警日志偶然發(fā)現(xiàn)大量的報(bào)錯(cuò),無法擴(kuò)展SYSAUX表空間,于是登錄系統(tǒng),查看系統(tǒng)表空間使用情況,發(fā)現(xiàn)SYSAUX表空間用滿了,所以本文給大家介紹了Oracle 12.2處理sysaux空間占滿問題,需要的朋友可以參考下
    2024-02-02

最新評論

漾濞| 新蔡县| 海丰县| 新河县| 上饶市| 嘉义市| 武清区| 乌拉特后旗| 宜丰县| 浏阳市| 都昌县| 龙州县| 凌海市| 正宁县| 清流县| 沙河市| 屏东县| 周口市| 辽宁省| 即墨市| 邯郸县| 宁河县| 康平县| 马边| 青田县| 安丘市| 称多县| 山丹县| 周口市| 绩溪县| 革吉县| 名山县| 景宁| 万年县| 准格尔旗| 伊金霍洛旗| 海宁市| 齐河县| 隆德县| 绥滨县| 壶关县|