文章详情

短信预约-IT技能 免费直播动态提醒

请输入下面的图形验证码

提交验证

短信预约提醒成功

kafka-3python生产者和消费者

2023-01-31 07:46

关注

程序分为productor.py是发送消息端,consumer为消费消息端,

启动的时候先启动product再启动consumer,毕竟只有发了消息,消费端才有消息可以消费,

productor.py

#!/usr/bin/env python2.7
#_*_coding: utf-8 _*_
from kafka import KafkaProducer


kafka_host = '192.168.1.200'  # kafka服务器地址
kafka_port = 9092  # kafka服务器的端口


producer = KafkaProducer(bootstrap_servers=['{kafka_host}:{kafka_port}'.format(
    kafka_host = kafka_host,
    kafka_port = kafka_port
)])


#简单for循环10次,发送10条消息
for i in range(1,10):
    message_string = 'some message'.format(i)

    #调用send方法,发送名字为'topic1'的topicid ,发送的消息为message_string
    response = producer.send('topic1', message_string.encode('utf-8'))
    print response


consumer.py

#!/usr/bin/env python
#_*_coding: utf-8 _*_
import json
from kafka import *

kafka_host = '192.168.1.200'  # kafka服务器地址
kafka_port = 9092  # kafka服务器端口


#消费topic1的topic,并指定group_id(自定义),多个机器或进程想顺序消费,可以指定同一个group_id,
# 如果想一条消费多次消费,可以换一个group_id,会从头开始消费
consumer = KafkaConsumer(
    'topic1',
    group_id = 'my-group',
    bootstrap_servers = ['{kafka_host}:{kafka_port}'.format(kafka_host=kafka_host, kafka_port=kafka_port)]
)
for message in consumer:
    #json读取kafka的消息
    content = json.loads(message.value)
    print content


阅读原文内容投诉

免责声明:

① 本站未注明“稿件来源”的信息均来自网络整理。其文字、图片和音视频稿件的所属权归原作者所有。本站收集整理出于非商业性的教育和科研之目的,并不意味着本站赞同其观点或证实其内容的真实性。仅作为临时的测试数据,供内部测试之用。本站并未授权任何人以任何方式主动获取本站任何信息。

② 本站未注明“稿件来源”的临时测试数据将在测试完成后最终做删除处理。有问题或投稿请发送至: 邮箱/279061341@qq.com QQ/279061341

软考中级精品资料免费领

  • 历年真题答案解析
  • 备考技巧名师总结
  • 高频考点精准押题
  • 2024年上半年信息系统项目管理师第二批次真题及答案解析(完整版)

    难度     813人已做
    查看
  • 【考后总结】2024年5月26日信息系统项目管理师第2批次考情分析

    难度     354人已做
    查看
  • 【考后总结】2024年5月25日信息系统项目管理师第1批次考情分析

    难度     318人已做
    查看
  • 2024年上半年软考高项第一、二批次真题考点汇总(完整版)

    难度     435人已做
    查看
  • 2024年上半年系统架构设计师考试综合知识真题

    难度     224人已做
    查看

相关文章

发现更多好内容

猜你喜欢

AI推送时光机
位置:首页-资讯-后端开发
咦!没有更多了?去看看其它编程学习网 内容吧
首页课程
资料下载
问答资讯