mariadb 10.4使用者身份驗證
阿新 • • 發佈:2021-01-11
知識點:
https://github.com/ververica/flink-cdc-connectors //官網地址
1、處理類
import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.source.SourceFunction; import com.alibaba.ververica.cdc.debezium.StringDebeziumDeserializationSchema;import com.alibaba.ververica.cdc.connectors.mysql.MySQLSource; /** * @program: Flink1.11 * @description: * @author: yang * @create: 2021-01-11 17:41 */ public class MySqlBinlogSourceExample { public static void main(String[] args) throws Exception { SourceFunction<String> sourceFunction = MySQLSource.<String>builder() .hostname("localhost") .port(3306) .databaseList("test") // monitor all tables under inventory database .username("root") .password("root") .deserializer(new StringDebeziumDeserializationSchema()) // converts SourceRecord to String.build(); StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration()); env.addSource(sourceFunction).print().setParallelism(1); // use parallelism 1 for sink to keep message ordering env.execute("test"); } }
2、列印結果
SourceRecord{sourcePartition={server=mysql-binlog-source}, sourceOffset={ts_sec=1610361995, file=mysql-bin.000004, pos=233436324, row=1, server_id=1, event=2}} ConnectRecord{topic='mysql-binlog-source.test.weblog', kafkaPartition=null, key=Struct{id=6}, keySchema=Schema{mysql_binlog_source.test.weblog.Key:STRUCT}, value=Struct{after=Struct{id=6,url=66666,method=66666,ip=666,args=6666,create_time=1610390788000},source=Struct{version=1.2.0.Final,connector=mysql,name=mysql-binlog-source,ts_ms=1610361995000,db=test,table=weblog,server_id=1,file=mysql-bin.000004,pos=233436459,row=0,thread=944986},op=c,ts_ms=1610361996498}, valueSchema=Schema{mysql_binlog_source.test.weblog.Envelope:STRUCT}, timestamp=null, headers=ConnectHeaders(headers=)}
3、如果需要將資料進行etl,可以自定義實現sink