
# 配置ClickHouse对接普通模式Kafka
#### 操作场景
本章节主要介绍ClickHouse连接普通模式的Kafka，消费Kafka的数据。
#### 约束与限制
本章节仅适用于MRS 3.3.0-LTS及之后版本。
#### 前提条件
- 已创建Kafka集群，且为普通模式（关闭Kerberos认证）。
- 已创建ClickHouse集群，并且ClickHouse集群和Kafka集群网络可以互通，并安装ClickHouse客户端。
 
#### 操作步骤
1. 登录ClickHouse服务所在集群的Manager页面，选择"集群 \> 服务 \> ClickHouse \> 配置 \> 全部配置 \> ClickHouseServer（角色） \> 引擎"，修改如下参数： 
   
   | 参数              | 参数说明                                |
   |:---|:---|
   | kafka_auth_mode | ClickHouse连接Kafka的认证方式，参数值选择NoAuth。 |
      
   ![](https://support.huaweicloud.com/cmpntguide-lts-mrs/zh-cn_image_0000001724990777.png "点击放大")
   
   
2. 选择"集群 \> 服务 \> ClickHouse \> 配置 \> 全部配置 \> ClickHouseServer（角色） \> 自定义"，在"clickhouse-config-customize"中添加如下参数： 
   
   | 名称                      | 值         |
   |:---|:---|
   | kafka.security_protocol | plaintext |
      
   ![](https://support.huaweicloud.com/cmpntguide-lts-mrs/zh-cn_image_0000001725962850.png "点击放大")
   
   
3. 单击"保存"，在弹窗页面中单击"确定"，保存配置。单击"实例"，勾选ClickHouseServer实例，选择"更多 \> 滚动重启实例"，重启ClickHouseServer实例。
4. 参考[Kafka客户端使用实践](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_1767.html)，登录到Kafka客户端安装目录。
   
   1. 以Kafka客户端安装用户，登录Kafka安装客户端的节点。
   
   2. 执行以下命令，切换到客户端安装目录。
      ```
      cd /opt/client
      ```
      
   
   3. 执行以下命令配置环境变量。
      ```
      source bigdata_env
      ```
      
   
   4. 如果当前集群已启用Kerberos认证，执行以下命令认证当前用户。如果当前集群未启用Kerberos认证，则无需执行此命令。
      ```
      kinit 组件业务用户
      ```
      
   
   
   
   
5. 执行以下命令，创建Kafka的Topic。详细的命令使用可以参考[创建Kafka Topic](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_24136.html)。
   
   ```
   kafka-topics.sh --topic topic1 --create --bootstrap-server <Kafka集群IP:21005> --command-config Kafka/kafka/config/client.properties --partitions 2 --replication-factor 1
   ```
   表1参数解释 
   | 参数                  | 参数说明                                                                                                                           |
   |:---|:---|
   | --topic             | 参数值为要创建的Topic名称，本示例创建的名称为topic1。                                                                                               |
   | --bootstrap-server  | Kafka集群节点IP。获取方式如下： 登录FusionInsight Manager页面，选择"集群 \> 服务 \> Kafka \> 实例"，查看Broker角色实例的业务IP地址。 |
   | --partitions        | 主题分区数，不能大于Kafka角色实例数量。                                                                                                         |
   | -replication-factor | 主题备份数，不能大于Kafka角色实例数量。                                                                                                         |
      
   
   
6. 登录ClickHouse客户端节点，连接ClickHouse服务端，具体请参考[ClickHouse客户端使用实践](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_2345.html)章节。
7. 创建Kafka的表引擎，示例如下： 
   ```
   CREATE TABLE queue1 (
   key String,
   value String,
   event_date DateTime
   ) ENGINE = Kafka()
   SETTINGS kafka_broker_list = 'kafka_ip1:21005,kafka_ip2:21005,kafka_ip3:21005',
   kafka_topic_list = 'topic1',
   kafka_group_name = 'group2',
   kafka_format = 'CSV',
   kafka_row_delimiter = '\n',
   kafka_handle_error_mode='stream';
   ```
   相关参数说明如下表：
   
   | 参数                         | 参数说明                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                |
   |:---|:---|
   | kafka_broker_list          | Kafka集群Broker实例的IP和端口列表。例如*：kafka集群broker实例IP1* :9092,*kafka集群broker实例IP2* :9092,*kafka集群broker实例IP3*:9092。 Kafka集群broker实例IP获取方法如下： 登录FusionInsight Manager页面，选择"集群 \> 服务 \> Kafka"。单击"实例"，查看Kafka角色实例的IP地址。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                     |
   | kafka_topic_list           | 消费Kafka的Topic。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      |
   | kafka_group_name           | Kafka消费组。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           |
   | kafka_format               | 消费数据的格式化类型，JSONEachRow表示每行一条数据的json格式，CSV格式表示逗号分隔的一行数据。更多请参考：<https://clickhouse.com/docs/en/interfaces/formats/>。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                  |
   | kafka_row_delimiter        | 每个消息体（记录）之间的分隔符。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                    |
   | kafka_handle_error_mode    | 设置为stream，会把每条消息处理的异常打印出来。需要创建视图，通过视图查询异常数据的具体处理异常。 创建视图语句，示例如下： ``` CREATE MATERIALIZED VIEW default.kafka_errors2 ( `topic` String, `key` String, `partition` Int64, `offset` Int64, `timestamp` Date, `timestamp_ms` Int64, `raw` String, `error` String ) ENGINE = MergeTree ORDER BY (topic, partition, offset) SETTINGS index_granularity = 8192 AS SELECT _topic AS topic, _key  AS key, _partition AS partition, _offset AS offset, _timestamp AS timestamp, _timestamp_ms AS timestamp_ms, _raw_message AS raw, _error AS error FROM default.queue1; ``` 查询视图，示例如下： ``` host1 :) select * from kafka_errors2;  SELECT * FROM kafka_errors2  Query id: bf4d788f-bcb9-44f5-95d0-a6c83c591ddb  ┌─topic──┬─key─┬─partition─┬─offset─┬──timestamp─┬─timestamp_ms─┬─raw─┬─error────────────────────────────────────────────────────────────────────────────────────────────────────────────────────┐ │ topic1 │     │         1 │      8 │ 2023-06-20 │   1687252213 │ 456 │ Cannot parse date: value is too short: (at row 1) Buffer has gone, cannot extract information about what has been parsed. │ └────────┴─────┴───────────┴────────┴────────────┴──────────────┴─────┴──────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────┘  1 rows in set. Elapsed: 0.003 sec.   host1 :) ``` |
   | kafka_skip_broken_messages | （可选）表示忽略解析异常的Kafka数据的条数。如果出现了N条异常后，后台线程结束，Materialized View会被重新安排后台线程去监测数据。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                         |
   | kafka_num_consumers        | （可选）单个Kafka Engine的消费者数量，通过增加该参数，可以提高消费数据吞吐，但总数不应超过对应topic的partitions总数。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                            |
      
   其他配置可参考<https://clickhouse.com/docs/zh/engines/table-engines/integrations/kafka>。
   
   
8. 通过客户端连接ClickHouse创建本地表，示例如下： 
   ```
   CREATE TABLE daily1(
   key String,
   value String,
   event_date DateTime
   )ENGINE = MergeTree()
   ORDER BY key;
   ```
   
   
9. 通过客户端连接ClickHouse创建物化视图，示例如下： 
   ```
   CREATE MATERIALIZED VIEW default.consumer1 TO default.daily1 (
   `event_date` DateTime,
   `key` String,
   `value` String
   ) AS
   SELECT
   event_date,
   key,
   value
   FROM default.queue1;
   ```
   
   
10. 再次执行[4]，进入Kafka客户端安装目录。
11. 执行以下命令，在Kafka的Topic中产生消息。例如，如下命令向[5]中创建的Topic发送消息。
    
    ```
    kafka-console-producer.sh --broker-list kafka集群broker实例IP1:9092,kafka集群broker实例IP2:9092,kafka集群broker实例IP3:9092  --topic topic1
    ```
    结果如下：
    ```
    >a1,b1,'2020-08-01 10:00:00'
    >a2,b2,'2020-08-02 10:00:00'
    >a3,b3,'2020-08-02 10:00:00'
    >a4,b4,'2023-09-02 10:00:00'
    ```
    
    
12. 查询消费到的Kafka数据，查询上述的物化视图，示例如下： 
    ```
    select * from daily;
    ```
    ![](https://support.huaweicloud.com/cmpntguide-lts-mrs/zh-cn_image_0000001723893864.png)
    
    
 
