前往小程序,Get更优阅读体验!
立即前往
首页
学习
活动
专区
工具
TVP
发布
社区首页 >专栏 >FlinkCDC的探索与实践1

FlinkCDC的探索与实践1

原创
作者头像
LarkMidTable
修改2022-09-23 22:32:50
5950
修改2022-09-23 22:32:50
举报
文章被收录于专栏:FlinkCDC

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元即可加入,添加需要备注【云雀课堂知识群】,这里可以获取到上面的源码,如果遇到问题可以一起解决,同时可以一起学习和进步。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
目录
  • 1.window开启binlog
  • 2.开启mysql的服务
  • 3.查看是否开启成功
  • 4.idea中引入依赖
  • 5.示例代码
  • 6.遇到的问题
    • 问题1:
      • 解决办法:
      领券
      问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档