标签:pac 消费者 def 生产者和消费者 循环 主题 pip 规模 ges
Kafka 是一个开源的分布式流处理平台,其简化了不同数据系统的集成。流指的是一个数据管道,应用能够通过流不断地接收数据。Kafka 作为流处理系统主要有两个用处:
相比于其它技术,Kafka 拥有更高的吞吐量,内置分区,副本和容错率。这些使得 Kafka 成为大规模消息处理应用的良好解决方案。
Kafka 系统有三个主要的部分:
见 Kafka 简单实验一
我们的项目将包括:
生产者:将字符串发送给 Kafka 消费者: 获取数据并展示在终端窗口中 Kafka: 作为中间人
安装需要的依赖包
pip install kafka-python
生产者是给 Kafka 中间人发送消息的服务。值得注意的是,生产者并不关注最终消费或加载数据的消费者。 创建生产者: 创建一个 producer.py 文件并添加如下代码:
import time from kafka import SimpleProducer, KafkaClient # connect to Kafka kafka = KafkaClient(‘localhost:9092‘) producer = SimpleProducer(kafka) # Assign a topic topic = ‘my-topic‘
创建消息:
循环生成1到100之间的数字
发送消息:
Kafka 消息是二进制字符串格式(byte)
以下是完整的生产者代码:
import time from kafka import SimpleProducer, KafkaClient # connect to Kafka kafka = KafkaClient(‘localhost:9092‘) producer = SimpleProducer(kafka) # Assign a topic topic = ‘my-topic‘ def test(): print(‘begin‘) n = 1 while (n<=100): producer.send_messages(topic, str(n)) print "send" + str(n) n += 1 time.sleep(0.5) print(‘done‘) if __name__ == ‘__main__‘: test()
消费者监听并消费来自 Kafka 中间人的消息。我们的消费者应该监听 my-topic 主题的消息并将消息展示。
以下是消费者代码(consumer.py):
from kafka import KafkaConsumer #connect to Kafka server and pass the topic we want to consume consumer = KafkaConsumer(‘my-topic‘) print "begin" for msg in consumer: print msg
确保 Kafka 在运行
打开两个终端,在第一个终端中运行消费者:
$ python consumer.py
在第二个终端运行生产者:
$ python producer.py
标签:pac 消费者 def 生产者和消费者 循环 主题 pip 规模 ges
原文地址:http://www.cnblogs.com/zhangtianyuan/p/7655904.html