
# FlinkServer作业对接Kafka消息队列
本章节适用于MRS 3.1.2及之后的版本。
#### 操作场景
本章节介绍Kafka作为source表或者sink表的DDL定义，以及创建表时使用的WITH参数和代码示例，并指导如何在FlinkServer作业管理页面操作。
本示例以安全模式Kafka为例。
#### 前提条件
- 集群中已安装HDFS、Yarn、Kafka和Flink服务。
- 包含Kafka服务的客户端已安装，例如安装路径为：/opt/client
- 参考[创建FlinkServer权限角色](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_24049.html)创建一个具有FlinkServer管理员权限的用户用于访问Flink WebUI，如：flink_admin。
 
#### 创建作业步骤
1. 使用**flink_admin**登录Manager，选择"集群 \> 服务 \> Flink"，在"Flink WebUI"右侧，单击链接，访问Flink的WebUI。
2. 参考[集群连接模式创建Flink SQL作业](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_24024.html)，新建Flink SQL流作业，在作业开发界面进行作业开发，配置完成后启动作业。
   
   需勾选"基础参数"中的"开启CheckPoint"，"时间间隔（ms）"可设置为"60000"，"模式"可使用默认值。
   ```
   CREATE TABLE KafkaSource (
     `user_id` VARCHAR,
     `user_name` VARCHAR,
     `age` INT
   ) WITH (
     'connector' = 'kafka',
     'topic' = 'test_source',
     'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号',
     'properties.group.id' = 'testGroup',
     'scan.startup.mode' = 'latest-offset',
     'format' = 'csv',
     'properties.sasl.kerberos.service.name' = 'kafka',--普通模式集群不需要该参数，同时删除上一行的逗号
     'properties.security.protocol' = 'SASL_PLAINTEXT',--普通模式集群不需要该参数
     'properties.kerberos.domain.name' = 'hadoop.系统域名'--普通模式集群不需要该参数
   );
   CREATE TABLE KafkaSink(
     `user_id` VARCHAR,
     `user_name` VARCHAR,
     `age` INT
   ) WITH (
     'connector' = 'kafka',
     'topic' = 'test_sink',
     'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号',
     'scan.startup.mode' = 'latest-offset',
     'value.format' = 'csv',
     'properties.sasl.kerberos.service.name' = 'kafka',--普通模式集群不需要该参数，同时删除上一行的逗号
     'properties.security.protocol' = 'SASL_PLAINTEXT',--普通模式集群不需要该参数
     'properties.kerberos.domain.name' = 'hadoop.系统域名'--普通模式集群不需要该参数
   );
   Insert into
     KafkaSink
   select
     *
   from
     KafkaSource;
   ```
   ![](https://support.huaweicloud.com/cmpntguide-lts-mrs/public_sys-resources/note_3.0-zh-cn.png)
   - Kafka Broker实例IP地址及Kafka端口号：
     - 服务的实例IP地址可通过登录FusionInsight Manager后，单击"集群 \> 服务 \> Kafka \> 实例"，在实例列表页面中查询。
     
     - 集群的"认证模式"为"安全模式"时为"sasl.port"的值，默认为"21007"。
     
     - 集群的"认证模式"为"普通模式"时为"port"的值，默认为"9092"。如果配置端口号为9092，则需要配置"allow.everyone.if.no.acl.found"参数为true，具体操作如下： 登录FusionInsight Manager系统，选择"集群 \> 服务 \> Kafka \> 配置 \> 全部配置"，搜索"allow.everyone.if.no.acl.found"配置，修改参数值为true，保存配置即可。
       
      
   
   - 系统域名：可登录FusionInsight Manager，选择"系统 \> 权限 \> 域和互信"，查看"本端域"参数，即为当前系统域名。
   
   - 使用Flink 1.15.0及以前版本对接Kafka，在扩容Kafka Topic分区后，需要重启相关的Flink作业，否则会导致新分区识别不及时漏消费数据。或在开发作业时，配置Flink动态发现Kafka Topic新分区功能。 可在作业SQL Kafka source表的WITH属性中，添加"scan.topic-partition-discovery.interval"参数，设置值为动态刷新时间，如"5min"。
     
    
   
   
3. 查看作业管理界面，作业状态为"运行中"。
4. 参考[使用Kafka生产消费数据](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_0379.html)，执行以下命令查看Sink表中是否接收到数据，即[5]执行完成后查看Kafka topic是否正常写入数据。
   
   ```
   sh kafka-console-consumer.sh --topic test_sink --bootstrap-server Kafka的Broker实例业务IP:Kafka端口号 --consumer.config /opt/client/Kafka/kafka/config/consumer.properties
   ```
   
   
5. 参考[使用Kafka生产消费数据](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_0379.html)，查看Topic并向Kafka中写入数据，输入完成后可在[4]中的窗口查看执行结果。
   
   查看Kafka主题：
   ```
   ./kafka-topics.sh --list --bootstrap-server Kafka的Broker实例业务IP:Kafka端口号 --command-config 客户端目录/Kafka/kafka/config/client.properties
   ```
   向Kafka中写入数据：
   ```
   sh kafka-console-producer.sh --broker-list Kafka角色实例所在节点的IP地址:Kafka端口号 --topic 主题名称 --producer.config 客户端目录/Kafka/kafka/config/producer.properties
   ```
   例如本示例使用主题名称为test_source：
   ```
   sh kafka-console-producer.sh --broker-list Kafka角色实例所在节点的IP地址:Kafka端口号 --topic test_source --producer.config  /opt/client/Kafka/kafka/config/producer.properties
   ```
   输入消息内容：
   ```
   1,clw,33
   ```
   输入完成后按回车发送消息。
   
   
 
#### WITH主要参数说明
| 配置项                                     | 是否必选                                                                                                                                                                                                     | 类型       | 描述                                                                                                                                                                                                                                                                                                        |
|:---|:---|:---|:---|
| connector                               | 必选                                                                                                                                                                                                       | String   | 指定要使用的连接器，Kafka使用"kafka"。                                                                                                                                                                                                                                                                                 |
| topic                                   | - kafka作为sink，必选  - kafka作为source，可选   | String   | 主题名称 - 当表用作source时，要从中读取数据的主题名称。支持主题列表，通过按分号分隔主题，如"主题-1；主题-2"。  - 当表用作sink时，主题名称为写入数据的主题。sink不支持主题列表。   |
| topic-pattern                           | kafka作为source时可选                                                                                                                                                                                         | String   | 主题模式 当表用作source时可设置该参数，主题名称需使用正则表达式。 不能同时设置"topic-pattern"和"topic"。                                                                                                                                                                         |
| properties.bootstrap.servers            | 必选                                                                                                                                                                                                       | String   | Kafka broker列表，以逗号分隔。                                                                                                                                                                                                                                                                                     |
| properties.group.id                     | kafka作为source时必选                                                                                                                                                                                         | String   | Kafka的使用者组ID。                                                                                                                                                                                                                                                                                             |
| format                                  | 必选                                                                                                                                                                                                       | String   | 用于反序列化和序列化Kafka消息的值部分的格式。                                                                                                                                                                                                                                                                                 |
| properties.\*                           | 可选                                                                                                                                                                                                       | String   | 安全模式下需增加认证相关的参数。                                                                                                                                                                                                                                                                                          |
| scan.topic-partition-discovery.interval | 可选                                                                                                                                                                                                       | Duration | 消费者定时动态发现创建的Partition的时间间隔。默认值：5min。                                                                                                                                                                                                                                                                      |
   
#### **起始消费位点**
**启动模式**
scan.startup.mode配置项决定了Kafka Consumer的启动模式。有效值为：
- group-offsets：从指定group id的已提交位点开始读取，group id通过properties.group.id指定。
- earliest-offset：从可能的最早偏移量开始。
- latest-offset：从最末尾偏移量开始。
- timestamp：从时间戳大于等于指定时间的第一条消息开始读取，时间戳通过scan.startup.timestamp-millis指定。
- specific-offsets：从指定的分区位点开始消费，位点通过scan.startup.specific-offsets指定。
**起始位点优先级**
1. Checkpoint或Savepoint中存储的位点。
2. WITH参数通过"scan.startup.mode"指定启动位点。
3. 如果不指定"scan.startup.mode"的情况下使用"group-offsets"启动消费。
   ![](https://support.huaweicloud.com/cmpntguide-lts-mrs/public_sys-resources/note_3.0-zh-cn.png)
   - 如果不指定启动位点"scan.startup.mode"，则默认会从已提交位点"group-offsets"启动。一种常见情况是使用全新的group id开始消费。首先源表会向Kafka集群查询该group的已提交位点，由于该group id是第一次使用，不会查询到有效位点，所以会通过"properties.auto.offset.reset"参数配置的策略进行重置。因此在使用全新group id进行消费时，必须配置"properties.auto.offset.reset"来指定位点重置策略。如果您未设置该配置项，则会产生异常"org.apache.kafka.clients.consumer.NoOffsetForPartitionException: Undefined offset with no reset policy for partitions"要求您介入。
   
   - 如果使用了"timestamp"，必须使用另外一个配置项"scan.startup.timestamp-millis"来指定一个从格林尼治标准时间1970年1月1日 00:00:00.000开始计算的毫秒单位时间戳作为起始时间。
   
   - 如果使用了"specific-offsets"，必须使用另外一个配置项"scan.startup.specific-offsets"来为每个partition指定起始偏移量。例如，选项值"partition:0,offset:42;partition:1,offset:300"表示"partition 0"从偏移量"42"开始，"partition 1"从偏移量"300"开始。
     
**代码示例如下**
```
CREATE TABLE kafka_source (
  ...
) WITH (
  'connector' = 'kafka',
  ...
  --启动位点示例，用户根据业务需求选择其中一种
  --从最早位点开始消费
  'scan.startup.mode' = 'earliest-offset',
  --从最末尾位点开始消费
  'scan.startup.mode' = 'latest-offset',
  --从消费者组"my-group"的已提交位点开始消费
  'properties.group.id' = 'my-group',
  'scan.startup.mode' = 'group-offsets',
  'properties.auto.offset.reset' = 'earliest', -- 如果 "my-group" 为首次使用，则从最早位点开始消费
  'properties.auto.offset.reset' = 'latest', -- 如果 "my-group" 为首次使用，则从最末尾位点开始消费
  --从指定的毫秒时间戳1655395200000开始消费
  'scan.startup.mode' = 'timestamp',
  'scan.startup.timestamp-millis' = '1655395200000',
  --从指定位点开始消费
  'scan.startup.mode' = 'specific-offsets',
  'scan.startup.specific-offsets' = 'partition:0,offset:42;partition:1,offset:300'
);
```
