
# FlinkSQL Kafka Connector支持消费drs-json格式数据
#### 使用场景
在数据复制服务（DRS）场景中，DRS将源端数据库的变更数据以drs-json格式（一种CDC消息格式）写入Kafka。用户需要使用FlinkSQL消费Kafka中的drs-json格式数据并进行实时处理，但FlinkSQL默认不支持直接解析drs-json格式。通过在Kafka Connector Source表中设置"format"为"drs-json"，即可直接消费Kafka中的drs-json格式数据。
#### 约束与限制
本章节仅适用于MRS 3.3.0及之后版本。
#### 使用方法
在创建的Kafka Connector Source流表中，设置 'format' = 'drs-json'，指定数据格式为drs-json。
SQL示例如下：
```
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' = 'drs-json',
  'properties.sasl.kerberos.service.name' = 'kafka',
  'properties.security.protocol' = 'SASL_PLAINTEXT',
  'properties.kerberos.domain.name' = 'hadoop.系统域名'
);
CREATE TABLE printSink(
  `user_id` VARCHAR,
  `user_name` VARCHAR,
  `age` INT
) WITH (
  'connector' = 'print'
);
Insert into
  printSink
select
  *
from
  KafkaSource;
```
![](https://support.huaweicloud.com/cmpntguide-lts-mrs/public_sys-resources/note_3.0-zh-cn.png)
- Kafka Broker实例IP地址及端口号说明：
  - 服务的实例IP地址可通过登录FusionInsight Manager后，单击"集群 \> 服务 \> Kafka \> 实例"，在实例列表页面中查询。
  
  - 集群已启用Kerberos认证（安全模式）时Broker端口为"sasl.port"参数的值，默认为"21007"。
  
  - 集群未启用Kerberos认证（普通模式）时Broker端口为"port"的值，默认为"9092"。如果配置端口号为9092，则需要配置"allow.everyone.if.no.acl.found"参数为true，具体操作如下： 登录FusionInsight Manager系统，选择"集群 \> 服务 \> Kafka \> 配置 \> 全部配置"，搜索"allow.everyone.if.no.acl.found"配置，修改参数值为true，保存配置即可。
    
   
- 系统域名：可登录FusionInsight Manager，选择"系统 \> 权限 \> 域和互信"，查看"本端域"参数，即为当前系统域名。
 
