SSL协议生产和消费消息
更新时间:2024-06-18
在 Kafka 集群所在 VPC 网络下或者公网环境访问,使用 SSL 协议接入,接入点可以在集群详情页面查看。
环境准备
                Shell
                
            
            1pip install confluent-kafka
            集群准备
1. 购买专享版消息服务for Kafka集群
开通消息服务 for Kafka服务后,在控制台页面点击『创建集群』,即可进行购买。

2. 为购买的集群创建主题
在控制台页面点击集群名称,进入集群详情页面。
在左侧的边栏中点击『主题管理』,进入主题管理页面。

在主题管理页面点击『创建主题』,进行主题的创建。
使用步骤:
步骤一:获取集群接入点
具体请参考:接入点查看。
步骤二:下载证书文件
下载证书文件:如何下载证书?
步骤三:编写测试代码
- 需要关注并自行修改的参数
 
| 参数名 | 含义 | 
|---|---|
| bootstrap.servers | 接入点信息 | 
| topic_name | 主题名称 | 
| message | 消息的具体内容 | 
| ssl.ca.location | ca.pem文件路径 | 
| ssl.certificate.location | client.pem文件路径 | 
| ssl.key.location | client.key文件路径 | 
| group.id | 消费组id | 
生产者代码示例
创建KafkaProducerDemo.py文件,具体代码示例如下:
                Python
                
            
            1from confluent_kafka import Producer
2
3producer = Producer({
4    # 接入点
5    'bootstrap.servers':'xxx.xx.xx.xx:9099',
6    # 接入协议
7    'security.protocol':'SSL',
8    'ssl.endpoint.identification.algorithm': 'none',
9    # 证书文件路径
10    'ssl.ca.location':'ca.pem',
11    'ssl.certificate.location':'client.pem',
12    'ssl.key.location':'client.key'
13})
14
15def callback_msg(err, msg):
16    if err is not None:
17        print('send failed:{}'.format(err))
18    else:
19        print('send success:{}'.format(msg.topic(),msg.partition()))
20
21for _ in range(100):
22    # 第一个参数topic_name填写创建的主题名称,第二个参数message写需要发送的消息内容
23    producer.produce('topic_name', "message".encode('utf-8'), callback = callback_msg)
24    producer.poll(0)
25    
26producer.flush()
            消费者代码示例
创建KafkaConsumerDemo.py文件,具体代码示例如下:
                Python
                
            
            1from confluent_kafka import Consumer
2
3consumer = Consumer({
4    # 接入点
5    'bootstrap.servers':'xxx.xx.xx.xx:9099',
6    # 接入协议
7    'security.protocol':'SSL',
8    'ssl.endpoint.identification.algorithm':'none',
9    # 证书文件路径
10    'ssl.ca.location':'ca.pem',
11    'ssl.certificate.location':'client.pem',
12    'ssl.key.location':'client.key',
13    # 消费组id
14    'group.id':'test_group',
15    'auto.offset.reset':'latest',
16    'fetch.message.max.bytes':'1024*512',
17})
18
19# 订阅的主题名称
20consumer.subscribe(['topic_name'])
21
22while True:
23    msg = consumer.poll(1.0)
24
25    if msg is None:
26        continue
27    if msg.error():
28        print("Consumer error: {}".format(msg.error()))
29        continue
30
31    print('Received message: {}'.format(msg.value().decode('utf-8')))
32
33consumer.close()
            步骤四:编译并运行
运行上述两个代码文件。
                Bash
                
            
            1# 启动消费者
2python KafkaConsumerDemo.py
3# 启动生产者
4python KafkaProducerDemo.py
            步骤五:查看集群监控
查看消息是否发送成功或消费成功有两种方式:
- 查看程序输出日志。
 - 在专享版消息服务 for Kafka控制台查看集群监控,获取集群生产、消息情况。
 
推荐使用第二种方式,下面介绍如何查看集群监控。
(1)在专享版消息服务 for Kafka的控制台页面找到需要连接的集群,点击集群名称进入『集群详情』页面。

(2)页面跳转后,进入左侧边中的『集群详情』页面。

(3)点击左侧边栏中的『集群监控』,进入『集群监控』页面。

(4)通过查看『集群监控』页面,提供的不同纬度的监控信息(集群监控、节点监控、主题监控、消费组监控),即可获知集群的生产和消费情况。
集群监控的具体使用请参考:集群监控

