
# 使用DLI Flink作业实时同步MRS Kafka数据至CloudTable ClickHouse集群
此章节为您介绍数据实时同步的最佳实践，通过数据湖探索服务DLI Flink作业将MRS kafka任务制造数据实时同步给ClickHouse，实现Kafka实时入库到ClickHouse的过程。
- 了解DLI请参见[数据湖探索产品介绍](https://support.huaweicloud.com/productdesc-dli/dli_01_0378.html)。
- 了解Kafka请参见[MRS产品介绍](https://support.huaweicloud.com/productdesc-mrs/mrs_08_001301.html)。
  图1数据同步流程图   
  ![](https://support.huaweicloud.com/bestpractice-cloudtable/zh-cn_image_0000002126023954.png "点击放大") 
#### 使用限制
- MRS集群未开启Kerberos认证。
- 为了确保网络连通，MRS集群必须与CloudTable集群的安全组、区域、VPC、子网保持一致。
- MRS与CloudTable安全组入方向添加DLI队列弹性资源网段，建立跨源连接，请参见[创建增强型跨源连接](https://support.huaweicloud.com/usermanual-dli/dli_01_0006.html)。
- 必须打通DLI上下游的网络连通性，请参考[测试地址连通性](https://support.huaweicloud.com/usermanual-dli/dli_01_0489.html#section0)。
 
#### 操作流程
基本流程如下：
1. [步骤一：创建CloudTable ClickHouse集群]
2. [步骤二：MRS集群中创建Flink作业制造数据]
3. [步骤三：创建DLI Flink任务进行数据同步]
4. [步骤四：结果验证]
 
#### 准备工作
- 已注册华为账号并开通华为云，具体请参见[注册华为账号并开通华为云](https://support.huaweicloud.com/usermanual-account/account_id_001.html)，且在使用CloudTable前检查账号状态，账号不能处于欠费或冻结状态。
- 已创建虚拟私有云和子网，参见创建[虚拟私有云和子网](https://support.huaweicloud.com/usermanual-vpc/zh-cn_topic_0013935842.html)。
 
 #### 步骤一：创建CloudTable ClickHouse集群
1. 登录[表格存储服务控制台](https://console.huaweicloud.com/cloudtable)，[创建ClickHouse集群](https://support.huaweicloud.com/usermanual-cloudtable/cloudtable_01_0300.html)。
2. 下载客户端。在左侧导航树单击"帮助"，然后在页面右侧单击"客户端下载"和"客户端校验文件"。
3. [安装客户端并校验客户端](https://support.huaweicloud.com/qs-cloudtable/cloudtable_06_0003.html#section4)，连接集群。
4. 建立flink数据库。
   ```
   create database flink;
   ```
   使用flink数据库。
   ```
   use flink;
   ```
   
5. 创建flink.order表。
   ```
   create table flink.order(order_id String,order_channel String,order_time String,pay_amount Float64,real_pay Float64,pay_time String,user_id String,user_name String,area_id String) ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/flink/order', '{replica}')ORDER BY order_id;
   ```
   
6. 查看表是否创建成功。
   ```
   select * from flink.order;
   ```
   
 
 #### 步骤二：MRS集群中创建Flink作业制造数据
1. 创建[MRS集群](https://support.huaweicloud.com/qs-mrs/mrs_09_0002.html#section4)。
2. 登录Manager，选择"集群 \> Flink \> 概览"，进入概览页面。
3. 单击"Flink WebUI"右侧的链接，访问Flink WebUI。
4. 在MRS Flink WebUI中创建Flink任务产生数据。
   1. 单击作业管理中的"新建作业"，弹出新建作业页面。
   
   2. 填写参数，单击"确定"，建立Flink SQL作业。如果修改SQL，单击操作列的"开发"，进入SQL页面添加以下命令。
      ![](https://support.huaweicloud.com/bestpractice-cloudtable/public_sys-resources/note_3.0-zh-cn.png)
      获取IP地址和端口：
      - IP地址获取：进入集群的Manager页面，单击"集群 \> Kafka \> 实例 \> 管理IP（Broker）"，可获取IP地址。
      
      - port获取：单击配置，进入配置页面，搜索"port"，获取端口（该port是Broker服务监听的PLAINTEXT协议端口号）。
      
      - 建议properties.bootstrap.servers参数添加多个ip:port，防止kafka实例网络不稳定或其他原因宕机，导致作业运行失败。
       
      SQL语句示例：
      ```
      CREATE TABLE IF NOT EXISTS `lineorder_ck` (
      `order_id` string,
      `order_channel` string,
      `order_time` string,
      `pay_amount` double,
      `real_pay` double,
      `pay_time` string,
      `user_id` string,
      `user_name` string,
      `area_id` string
      ) WITH (
      'connector' = 'kafka',
      'topic' = 'test_flink',
      'properties.bootstrap.servers' = 'ip:port',
      'value.format' = 'json',
      'properties.sasl.kerberos.service.name' = 'kafka'
      );
      CREATE TABLE lineorder_datagen (
      `order_id` string,
      `order_channel` string,
      `order_time` string,
      `pay_amount` double,
      `real_pay` double,
      `pay_time` string,
      `user_id` string,
      `user_name` string,
      `area_id` string
      ) WITH (
      'connector' = 'datagen',
      'rows-per-second' = '1000'
      );
      INSERT INTO
      lineorder_ck
      SELECT
      *
      FROM
      lineorder_datagen;
      ```
      
   
   3. 回到作业管理界面，单击操作列的"启动"。作业状态为"运行中"表示作业运行成功。
    
 
 #### 步骤三：创建DLI Flink任务进行数据同步
1. 创建弹性资源池和队列，请参见"[创建弹性资源池并添加队列](https://support.huaweicloud.com/usermanual-dli/dli_01_0505.html)"章节。
2. 创建跨源连接，请参见[创建增强型跨源连接](https://support.huaweicloud.com/usermanual-dli/dli_01_0006.html)。
3. 分别测试DLI与上游MRS Kafka和下游CloudTable ClickHouse的连通性。
   1. 弹性资源池和队列创建后，单击"资源管理 \> 队列管理"，进入队列管理界面测试地址连通性，请参见[测试地址连通性](https://support.huaweicloud.com/usermanual-dli/dli_01_0489.html#section0)。
   
   2. 获取上游IP地址和端口：进入集群的Manager页面，单击"集群 \> Kafka \> 实例 \> 管理IP（Broker）"，可获取IP地址。单击配置，进入配置页面，搜索"port"，获取端口（该port是Broker服务监听的PLAINTEXT协议端口号）。
   
   3. 获取下游IP地址和端口：进入集群详情页可查看节点IP和端口。
    
4. 创建Flink作业，请参见[使用DLI提交Flink作业](https://support.huaweicloud.com/usermanual-dli/dli_01_0498.html)。
5. 选择[4]中创建的Flink作业，单击操作列的"编辑"，添加SQL进行数据同步。
   ```
   create table orders (
   order_id string,
   order_channel string,
   order_time string,
   pay_amount double,
   real_pay double,
   pay_time string,
   user_id string,
   user_name string,
   area_id string
   ) WITH (
   'connector' = 'kafka',
   'topic' = 'test_flink',
   'properties.bootstrap.servers' = 'ip:port',
   'properties.group.id' = 'testGroup_1',
   'scan.startup.mode' = 'latest-offset',
   'format' = 'json'
   );
   create table clickhouseSink(
   order_id string,
   order_channel string,
   order_time string,
   pay_amount double,
   real_pay double,
   pay_time string,
   user_id string,
   user_name string,
   area_id string
   ) with (
   'connector' = 'clickhouse',
   'url' = 'jdbc:clickhouse://ip:port/flink',
   'username' = 'admin',
   'password' = '****',
   'table-name' = 'order',
   'sink.buffer-flush.max-rows' = '10',
   'sink.buffer-flush.interval' = '3s'
   );
   insert into clickhouseSink select * from orders;
   ```
   
6. 单击"格式化"，再单击"保存"。
   ![](https://support.huaweicloud.com/bestpractice-cloudtable/public_sys-resources/notice_3.0-zh-cn.png)
   请务必先单击"格式化"将SQL代码进行格式化处理，否则可能会因为代码复制和粘贴操作过程中引入新的空字符，而导致作业执行失败。
   
7. 回到DLI控制台首页，单击左侧"作业管理 \> Flink作业"。
8. 启动[4]中创建的作业，单击操作列的"启动 \> 立即启动"。作业状态为"运行中"表示作业运行成功。
 
 #### 步骤四：结果验证
1. 待MRS Flink任务和DLI任务运行成功后，返回ClickHouse集群运行命令的窗口。
2. 查看数据库。
   ```
   show databases;
   ```
   
3. 使用数据库。
   ```
   use database;
   ```
   
4. 查看数据表。
   ```
   show tables;
   ```
   
5. 查看同步数据。
   ```
   select * from order limit 10;
   ```
   图2查看同步数据   
   ![](https://support.huaweicloud.com/bestpractice-cloudtable/zh-cn_image_0000002132036017.png "点击放大") 
 
