1.window开启binlog
my.ini 下面 [mysqld] 加入如下内容
代码语言:javascript复制[mysqld]
log_bin=mysql-bin
binlog-format=ROW
server-id=1
2.开启mysql的服务
代码语言:javascript复制net start mysql
3.查看是否开启成功
代码语言:javascript复制show variables like 'log_bin';
log_bin ON
4.idea中引入依赖
代码语言:javascript复制 <dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>1.12.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_2.12</artifactId>
<version>1.12.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients_2.12</artifactId>
<version>1.12.0</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>3.1.3</version>
</dependency>
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>5.1.49</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-planner-blink_2.12</artifactId>
<version>1.12.0</version>
</dependency>
<dependency>
<groupId>com.ververica</groupId>
<artifactId>flink-connector-mysql-cdc</artifactId>
<version>2.0.0</version>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.75</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-assembly-plugin</artifactId>
<version>3.0.0</version>
<configuration>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
</configuration>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
5.示例代码
代码语言:javascript复制public static void main(String[] args) throws Exception {
// 1.设置流的环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// // 2.设置并行度
// env.setParallelism(1);
// // 3.1状态信息保存到CK中,进行断点续传,从checkpoint和savepoint开始
// env.enableCheckpointing(5000L);
// // 3.2 设置仅一次的语义
// env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// //3.3 设置任务关闭的时候保留最后一次 CK 数据
// env.getCheckpointConfig().enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
// //3.4 指定从 CK 自动重启策略
// env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, 2000L));
// //3.5 设置状态后端
// env.setStateBackend(new FsStateBackend("hdfs://192.168.1.204:9000/flinkCDC"));
// //3.6 设置访问 HDFS 的用户名
// System.setProperty("HADOOP_USER_NAME", "hadoop");
//4 获取mysql的数据源
DebeziumSourceFunction<String> sourceFunction = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.username("root")
.password("root")
.databaseList("test")
.serverTimeZone("UTC")
// .tableList("cdc_test.user_info")
.deserializer(new StringDebeziumDeserializationSchema())
.startupOptions(StartupOptions.initial())
.build();
DataStreamSource<String> dataStreamSource = env.addSource(sourceFunction);
//5.数据打印
dataStreamSource.print();
//6.启动任务
env.execute("FlinkCDC");
}
6.遇到的问题
问题1:
代码语言:javascript复制Caused by: com.mysql.cj.exceptions.InvalidConnectionAttributeException: The server time zone
value '�й���ʱ��' is unrecognized or represents more than one time zone. You must configure
either the server or JDBC driver (via the 'serverTimezone' configuration property)
to use a more specifc time zone value if you want to utilize time zone support.
解决办法:
代码语言:javascript复制MySqlSource.<String>builder().serverTimeZone("UTC")
本人开通付费的知识群,如果需要可以添加QQ:975863632,需要99.9元即可加入,添加需要备注【云雀课堂知识群】,这里可以获取到上面的源码,如果遇到问题可以一起解决,同时可以一起学习和进步。