
# Flink向Kafka生产并消费数据应用开发思路
假定某个Flink业务每秒就会收到1个消息记录。
基于某些业务要求，开发的Flink应用程序实现功能：实时输出带有前缀的消息内容。
 #### 数据规划
Flink样例工程的数据存储在Kafka组件中。Flink向Kafka组件发送数据（需要有kafka权限用户），并从Kafka组件获取数据。
1. 确保集群安装完成，包括HDFS、Yarn、Flink和Kafka。
2. 创建Topic。 
   1. 在服务端配置用户创建topic的权限。 开启Kerberos认证的安全集群将Kafka的Broker配置参数"allow.everyone.if.no.acl.found"的值修改为"true"。配置完后重启kafka服务。未开启Kerberos认证的普通集群无需此配置。
      
   
   2. 用户使用Linux命令创建topic，如果是安全集群，用户执行命令前需要使用kinit命令进行人机认证，如：***kinit flinkuser*** 。
      ![](https://support.huaweicloud.com/devg-mrs/public_sys-resources/note_3.0-zh-cn.png)
      flinkuser需要用户自己创建，并拥有创建Kafka的topic权限。具体操作请参考[准备Flink应用开发用户](https://support.huaweicloud.com/devg-mrs/mrs_06_0389.html)。
      创建topic的命令格式：
      ***bin/kafka-topics.sh --create --zookeeper {zkQuorum}/kafka --partitions {partitionNum} --replication-factor {replicationNum} --topic {Topic}***
      表1参数说明 
      | 参数名              | 说明                        |
      |:---|:---|
      | {zkQuorum}       | ZooKeeper集群信息，格式为IP:port。 |
      | {partitionNum}   | topic的分区数。                |
      | {replicationNum} | topic中每个partition数据的副本数。  |
      | {Topic}          | Topic名称。                  |
         
      示例：在Kafka的客户端路径下执行命令，此处以ZooKeeper集群的IP:port是10.96.101.32:2181,10.96.101.251:2181,10.96.101.177:2181，Topic名称为topic1的数据为例。
      ```
      bin/kafka-topics.sh --create --zookeeper 10.96.101.32:2181,10.96.101.251:2181,10.96.101.177:2181/kafka --partitions 5 --replication-factor 1 --topic topic1
      ```
      
   
   
   
   
3. 如果集群开启了kerberos， 执行该步骤进行安全认证，否则跳过该步骤。 
   - **Kerberos** **认证配置**
     1. 客户端配置。 在Flink配置文件"flink-conf.yaml"中，增加kerberos认证相关配置（主要在"contexts"项中增加"KafkaClient"），示例如下：
        ```
        security.kerberos.login.keytab: /home/demo/flink/release/flink-1.2.1/keytab/admin.keytab
        security.kerberos.login.principal: admin
        security.kerberos.login.contexts: Client,KafkaClient
        security.kerberos.login.use-ticket-cache: false
        ```
        
     
     2. 运行参数。 关于"SASL_PLAINTEXT"协议的运行参数示例如下：
        ```
        --topic topic1 --bootstrap.servers 10.96.101.32:21007 --security.protocol SASL_PLAINTEXT  --sasl.kerberos.service.name kafka //10.96.101.32:21007表示kafka服务器的IP:port
        ```
        
      
   
   
   
   
 
#### 开发思路
1. 启动Flink Kafka Producer应用向Kafka发送数据。
2. 启动Flink Kafka Consumer应用从Kafka接收数据，保证topic与producer一致。
3. 在数据内容中增加前缀并进行打印。
 
