
# 从Kafka读取数据写入到Elasticsearch
![](https://support.huaweicloud.com/devg-dli/public_sys-resources/notice_3.0-zh-cn.png)
本指导仅适用于Flink 1.12版本。
#### 场景描述
本示例场景对用户购买商品的数据信息进行分析，将满足特定条件的数据结果进行汇总输出。购买商品数据信息为数据源发送到Kafka中，再将Kafka数据的分析结果输出到Elasticsearch中。
例如，输入如下样例数据：
```
{"order_id":"202103241000000001", "order_channel":"webShop", "order_time":"2021-03-24 10:00:00", "pay_amount":"100.00", "real_pay":"100.00", "pay_time":"2021-03-24 10:02:03", "user_id":"0001", "user_name":"Alice", "area_id":"330106"}
{"order_id":"202103241606060001", "order_channel":"appShop", "order_time":"2021-03-24 16:06:06", "pay_amount":"200.00", "real_pay":"180.00", "pay_time":"2021-03-24 16:10:06", "user_id":"0002", "user_name":"Jason", "area_id":"330106"}
```
DLI从Kafka读取数据写入Elasticsearch，在Elasticsearch集群的Kibana中查看相应结果。
#### 前提条件
1. 已创建DMS Kafka实例。 具体步骤可参考：[DMS Kafka入门指引](https://support.huaweicloud.com/qs-kafka/kafka-qs-0409001.html)。
   ![](https://support.huaweicloud.com/devg-dli/public_sys-resources/caution_3.0-zh-cn.png)
   创建DMS Kafka实例时，**不能开启Kafka SASL_SSL**。
   
2. 已创建Elasticsearch类型的CSS集群。 具体创建CSS集群的操作可以参考[创建CSS集群](https://support.huaweicloud.com/usermanual-css/css_01_0011.html)。
   本示例创建的CSS集群版本为：7.6.2，集群为非安全集群。
   
 
#### 整体作业开发流程
整体作业开发流程参考[图1]。
图1作业开发流程   
![](https://support.huaweicloud.com/devg-dli/zh-cn_image_0000002430392568.png "点击放大")
[步骤1：创建弹性资源池并添加队列]：创建DLI作业运行的队列。
[步骤2：创建Kafka的Topic]：创建Kafka生产消费数据的Topic。
[步骤3：创建Elasticsearch搜索索引]：创建Elasticsearch搜索索引用于接收结果数据。
[步骤4：创建增强型跨源连接]：DLI上创建连接Kafka和CSS的跨源连接，打通网络。
[步骤5：运行作业]：DLI上创建和运行Flink OpenSource作业。
[步骤6：发送数据和查询结果]：Kafka上发送流数据，在CSS上查看运行结果。
 #### 步骤1：创建弹性资源池并添加队列
新建队列的网段不能和DMS Kafka、RDS for MySQL实例的子网网段有重合，否则后续创建跨源连接会失败。
1. 登录DLI管理控制台。
2. 在左侧导航栏单击"资源管理 \> 弹性资源池"，可进入弹性资源池管理页面。
3. 在弹性资源池管理界面，单击界面右上角的"购买弹性资源池"。
4. 在"购买弹性资源池"界面，填写具体的弹性资源池参数。
   本例在华东-上海二区域购买按需计费的弹性资源池。相关参数说明如[表1]所示。
    表1参数说明 
   | 参数名称 | 参数说明                                                             | 配置样例              |
   |:---|:---|:---|
   | 计费模式 | 选择弹性资源池计费模式。                                                     | 按需计费              |
   | 区域   | 选择弹性资源池所在区域。                                                     | 华东-上海二            |
   | 项目   | 每个区域默认对应一个项目，由系统预置。                                              | 系统默认项目            |
   | 名称   | 弹性资源池名称。                                                         | dli_resource_pool |
   | 规格   | 选择弹性资源池规格。                                                       | 标准版               |
   | CU范围 | 弹性资源池最大最小CU范围。                                                   | 64-64             |
   | 网段   | 规划弹性资源池所属的网段。如需使用DLI增强型跨源，弹性资源池网段与数据源网段不能重合。**弹性资源池网段设置后不支持更改**。 | 172.16.0.0/19     |
   | 企业项目 | 选择对应的企业项目。                                                       | default           |
      
   
5. 参数填写完成后，单击"立即购买"，在界面上确认当前配置是否正确。
6. 单击"提交"完成弹性资源池的创建。
7. 在弹性资源池的列表页，选择要操作的弹性资源池，单击操作列的"添加队列"。
8. 配置队列的基础配置，具体参数信息如下。
   表2弹性资源池添加队列基础配置 
   | 参数名称 | 参数说明                                                                                                                                                                                                                                                                                                                                                                              | 配置样例                                                                                   |
   |:---|:---|:---|
   | 名称   | 弹性资源池添加的队列名称。                                                                                                                                                                                                                                                                                                                                                                     | dli_queue_01                                                                           |
   | 类型   | 选择创建的队列类型。 - 执行SQL作业请选择SQL队列。  - 执行Flink或Spark作业请选择通用队列。   | SQL作业场景请选择"SQL队列"。 其他场景请选择"通用队列"。 |
   | 执行引擎 | SQL队列可以选择队列引擎为Spark。                                                                                                                                                                                                                                                                                                                                                              | Spark                                                                                  |
   | 企业项目 | 选择对应的企业项目。                                                                                                                                                                                                                                                                                                                                                                        | default                                                                                |
      
   
9. 单击"下一步"，配置队列的扩缩容策略。 单击"新增"，可以添加不同优先级、时间段、"最小CU"和"最大CU"扩缩容策略。
   本例配置的扩缩容策略如[图2]所示。
   图2添加队列时配置扩缩容策略   
   ![](https://support.huaweicloud.com/devg-dli/zh-cn_image_0000002000460541.png "点击放大")
   表3扩缩容策略参数说明 
   | 参数名称 | 参数说明                                                                                              | 配置样例  |
   |:---|:---|:---|
   | 优先级  | 当前弹性资源池中的优先级数字越大表示优先级越高。本例设置一条扩缩容策略，默认优先级为1。                                                      | 1     |
   | 时间段  | 首条扩缩容策略是默认策略，不能删除和修改时间段配置。 即设置00-24点的扩缩容策略。 | 00-24 |
   | 最小CU | 设置扩缩容策略支持的最小CU数。                                                                                  | 16    |
   | 最大CU | 当前扩缩容策略支持的最大CU数。                                                                                  | 64    |
      
   
10. 单击"确定"完成添加队列配置。
 
 #### 步骤2：创建Kafka的Topic
1. 在Kafka管理控制台，选择"Kafka专享版"，单击对应的Kafka名称，进入到Kafka的基本信息页面。
2. 单击"Topic管理 \> 创建Topic"，创建一个Topic。Topic配置参数如下：
   - Topic名称。本示例输入为：testkafkatopic。
   
   - 分区数：1。
   
   - 副本数：1。
   
   
   其他参数保持默认即可。
   
 
 #### 步骤3：创建Elasticsearch搜索索引
1. 登录CSS管理控制台，选择"集群管理 \> Elasticsearch"。
2. 在集群管理界面，在已创建的CSS集群的"操作"列，单击"Kibana"访问集群。
3. 在Kibana的左侧导航中选择"Dev Tools"，进入到Console界面。
4. 在Console界面，执行如下命令创建索引"shoporders"。
   ```
   PUT /shoporders
   {
     "settings": {
       "number_of_shards": 1
     },
   "mappings": {
     "properties": {
       "order_id": {
         "type": "text"
       },
       "order_channel": {
         "type": "text"
       },
       "order_time": {
         "type": "text"
       },
       "pay_amount": {
         "type": "double"
       },
       "real_pay": {
         "type": "double"
       },
       "pay_time": {
         "type": "text"
       },
       "user_id": {
         "type": "text"
       },
       "user_name": {
         "type": "text"
       },
       "area_id": {
         "type": "text"
       }
     }
   }
   }
   ```
   
 
 #### 步骤4：创建增强型跨源连接
- **创建DLI连接Kafka的增强型跨源连接**
  1. 在Kafka管理控制台，选择"Kafka专享版"，单击对应的Kafka名称，进入到Kafka的基本信息页面。
  
  2. 在"连接信息"中获取该Kafka的"内网连接地址"，在"基本信息"的"网络"中获取获取该实例的"虚拟私有云"和"子网"信息，方便后续操作步骤使用。
  
  3. 单击"网络"中的安全组名称，在"入方向规则"中添加放通队列网段的规则。例如，本示例队列网段为"10.0.0.0/16"，则规则添加为：优先级选择：1，策略选择：允许，协议选择：TCP，端口值不填，类型：IPv4，源地址为：10.0.0.0/16，单击"确定"完成安全组规则添加。
  
  4. 登录DLI管理控制台，在左侧导航栏单击"跨源管理"，在跨源管理界面，单击"增强型跨源"，单击"创建"。
  
  5. 在增强型跨源创建界面，配置具体的跨源连接参数。具体参考如下。
     - 连接名称：设置具体的增强型跨源名称。本示例输入为：dli_kafka。
     
     - 弹性资源池：选择[步骤1：创建弹性资源池并添加队列]中已经创建的弹性资源池。
     
     - 虚拟私有云：选择Kafka的虚拟私有云。
     
     - 子网：选择Kafka的子网。
     
     - 其他参数可以根据需要选择配置。
     
     
     参数配置完成后，单击"确定"完成增强型跨源配置。单击创建的跨源连接名称，查看跨源连接的连接状态，等待连接状态为："已激活"后可以进行后续步骤。
     
  
  6. 单击"队列管理"，选择操作的队列，本示例为[步骤1：创建弹性资源池并添加队列]中添加的队列，在操作列，单击"更多 \> 测试地址连通性"。
  
  7. 在"测试连通性"界面，根据中获取的Kafka连接信息，地址栏输入"Kafka内网地址:Kafka数据库端口"，单击"测试"测试DLI到Kafka网络是否可达。
   
- **创建DLI连接CSS的增强型跨源连接**
  1. 在CSS管理控制台，选择"集群管理"，单击已创建的CSS集群名称，进入到CSS的基本信息页面。
  
  2. 在"基本信息"中获取CSS的"内网访问地址"、"虚拟私有云"和"子网"信息，方便后续操作步骤使用。
  
  3. 单击"连接信息"中的安全组名称，在"入方向规则"中添加放通队列网段的规则。例如，本示例队列网段为"10.0.0.0/16"，则规则添加为：优先级选择：1，策略选择：允许，协议选择：TCP，端口值不填，类型：IPv4，源地址为：10.0.0.0/16，单击"确定"完成安全组规则添加。
  
  4. Kafka和CSS实例属于同一VPC和子网下？
     1. 是，执行[7]。Kafka和CSS实例在同一VPC和子网，不用再重复创建增强型跨源连接。
     
     2. 否，执行[5]。Kafka和CSS实例分别在两个VPC和子网下，则要分别创建增强型跨源连接打通网络。
      
  
  5. 登录DLI管理控制台，在左侧导航栏单击"跨源管理"，在跨源管理界面，单击"增强型跨源"，单击"创建"。
  
  6. 在增强型跨源创建界面，配置具体的跨源连接参数。具体参考如下。
     - 连接名称：设置具体的增强型跨源名称。本示例输入为：dli_css。
     
     - 弹性资源池：选择[步骤1：创建弹性资源池并添加队列]中已经创建的弹性资源池。
     
     - 虚拟私有云：选择CSS的虚拟私有云。
     
     - 子网：选择CSS的子网。
     
     - 其他参数可以根据需要选择配置。
     
     
     参数配置完成后，单击"确定"完成增强型跨源配置。单击创建的跨源连接名称，查看跨源连接的连接状态，等待连接状态为："已激活"后可以进行后续步骤。
     
  
  7. 单击"队列管理"，选择操作的队列，本示例为[步骤1：创建弹性资源池并添加队列]中添加的队列，在操作列，单击"更多 \> 测试地址连通性"。
  
  8. 在"测试连通性"界面，根据[2]获取的CSS连接信息，地址栏输入"CSS内网地址:CSS内网端口"，单击"测试"测试DLI到CSS网络是否可达。
   
 
 #### 步骤5：运行作业
1. 在DLI管理控制台，单击"作业管理 \> Flink作业"，在Flink作业管理界面，单击"创建作业"。
2. 在创建队列界面，类型选择"Flink OpenSource SQL"，名称填写为：FlinkKafkaES。单击"确定"，跳转到Flink作业编辑界面。
3. 在Flink OpenSource SQL作业编辑界面，配置如下参数，其他参数默认即可。
   - 所属队列：选择[步骤1：创建弹性资源池并添加队列]中创建的队列。
   
   - Flink版本：选择1.12。
   
   - 保存作业日志：勾选。
   
   - OBS桶：选择保存作业日志的OBS桶，根据提示进行OBS桶权限授权。
   
   - 开启Checkpoint：勾选。
   
   - Flink作业编辑框中输入具体的作业SQL，本示例作业参考如下。SQL中加粗的参数需要根据实际情况修改。
     ![](https://support.huaweicloud.com/devg-dli/public_sys-resources/note_3.0-zh-cn.png)
     本示例使用的Flink版本为1.12，故Flink OpenSource SQL语法也是1.12。本示例数据源是Kafka，写入结果数据到Elasticsearch。
     请参考[Flink OpenSource SQL 1.12创建Kafka源表](https://support.huaweicloud.com/sqlref-flink-dli/dli_08_0386.html)和[Flink OpenSource SQL 1.12创建Elasticsearch结果表](https://support.huaweicloud.com/sqlref-flink-dli/dli_08_0395.html)。
     
   
   - 创建Kafka源表，将DLI和Kafka数据源进行连接。
     ```
     CREATE TABLE kafkaSource (
       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",
       "properties.bootstrap.servers" = "10.128.0.120:9092,10.128.0.89:9092,10.128.0.83:9092",--替换为kafka的内网连接地址和端口
       "properties.group.id" = "click",
       "topic" = "testkafkatopic", --创建的Kafka Topic
       "format" = "json",
       "scan.startup.mode" = "latest-offset"
     );
     ```
     
   
   - 创建Elasticsearch结果表，将DLI分析后的数据的结果展示在Elasticsearch结果表上。
     ```
     CREATE TABLE elasticsearchSink (
       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' = 'elasticsearch-7',
       'hosts' = '192.168.168.125:9200', --替换为CSS集群的内网地址和端口
       'index' = 'shoporders' --创建的Elasticsearch搜索引擎
     );
     --将Kafka数据写入到Elasticsearch索引中
     insert into
       elasticsearchSink
     select
       *
     from
       kafkaSource;
     ```
     
    
4. 单击"语义校验"确保SQL语义校验成功。单击"保存"，保存作业。单击"启动"，启动作业，确认作业参数信息，单击"立即启动"开始执行作业。等待作业运行状态变为"运行中"。
 
 #### 步骤6：发送数据和查询结果
1. Kafka端发送数据。 使用Kafka客户端向[步骤2：创建Kafka的Topic]中的Topic发送数据，模拟实时数据流。
   Kafka生产和发送数据的方法请参考：[DMS - 连接实例生产消费信息](https://support.huaweicloud.com/usermanual-kafka/kafka-ug-180604020.html)。
   发送样例数据如下：
   ```
   {"order_id":"202103241000000001", "order_channel":"webShop", "order_time":"2021-03-24 10:00:00", "pay_amount":"100.00", "real_pay":"100.00", "pay_time":"2021-03-24 10:02:03", "user_id":"0001", "user_name":"Alice", "area_id":"330106"}
   {"order_id":"202103241606060001", "order_channel":"appShop", "order_time":"2021-03-24 16:06:06", "pay_amount":"200.00", "real_pay":"180.00", "pay_time":"2021-03-24 16:10:06", "user_id":"0002", "user_name":"Jason", "area_id":"330106"}
   ```
   
2. 查看Elasticsearch端数据处理后的相应结果。
   发送成功后，在CSS集群的Kibana中执行下述语句并查看相应结果：
   ```
   GET shoporders/_search
   ```
   查询结果返回如下：
   ```
   {
     "took" : 0,
     "timed_out" : false,
     "_shards" : {
       "total" : 1,
       "successful" : 1,
       "skipped" : 0,
       "failed" : 0
     },
     "hits" : {
       "total" : {
         "value" : 2,
         "relation" : "eq"
       },
       "max_score" : 1.0,
       "hits" : [
         {
           "_index" : "shoporders",
           "_type" : "_doc",
           "_id" : "6fswzIAByVjqg3_qAyM1",
           "_score" : 1.0,
           "_source" : {
             "order_id" : "202103241000000001",
             "order_channel" : "webShop",
             "order_time" : "2021-03-24 10:00:00",
             "pay_amount" : 100.0,
             "real_pay" : 100.0,
             "pay_time" : "2021-03-24 10:02:03",
             "user_id" : "0001",
             "user_name" : "Alice",
             "area_id" : "330106"
           }
         },
         {
           "_index" : "shoporders",
           "_type" : "_doc",
           "_id" : "6vs1zIAByVjqg3_qyyPp",
           "_score" : 1.0,
           "_source" : {
             "order_id" : "202103241606060001",
             "order_channel" : "appShop",
             "order_time" : "2021-03-24 16:06:06",
             "pay_amount" : 200.0,
             "real_pay" : 180.0,
             "pay_time" : "2021-03-24 16:10:06",
             "user_id" : "0002",
             "user_name" : "Jason",
             "area_id" : "330106"
           }
         }
       ]
     }
   }
   ```
   
   
 
