使用pykafka進(jìn)行生產(chǎn)和消費(fèi)

生產(chǎn)側(cè)

from pykafka import KafkaClient
import sys
import json
client = KafkaClient(hosts='192.168.1.1')
topics = client.topics
topic = topics['test_topic']
print(topic)
written_msgs = 0
with topic.get_producer() as producer:
  for line in sys.stdin:
    line = line.strip('\n')
    items = line.split('\t')
    user_id = items[0]
    device_id = items[1]
   
    producer.produce(json.dumps({'user_id': user_id,
                                'device_id':device_id
                              }))
    written_msgs += 1
    if written_msgs % 1000 == 0:
      print("written_msgs: %d" % written_msgs)
print("written_msgs: %d" % written_msgs)

消費(fèi)側(cè)

#coding:utf-8
from pykafka import KafkaClient
import time
from pykafka.common import OffsetType
#使用和生產(chǎn)者一樣的host以及topic
client = KafkaClient(hosts='192.168.1.1')
topics = client.topics
topic = topics['test_topic']
print(topic)

consumer = topic.get_simple_consumer(
    consumer_group="mygroup",
    auto_commit_enable=True #加這一行是為了避免重復(fù)消費(fèi)
)
## 獲取被消費(fèi)數(shù)據(jù)的偏移量和消費(fèi)內(nèi)容
for message in consumer:
    if message is not None:
    #打印消費(fèi)信息
        print(message.offset, message.value)
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
【社區(qū)內(nèi)容提示】社區(qū)部分內(nèi)容疑似由AI輔助生成,瀏覽時(shí)請(qǐng)結(jié)合常識(shí)與多方信息審慎甄別。
平臺(tái)聲明:文章內(nèi)容(如有圖片或視頻亦包括在內(nèi))由作者上傳并發(fā)布,文章內(nèi)容僅代表作者本人觀點(diǎn),簡(jiǎn)書(shū)系信息發(fā)布平臺(tái),僅提供信息存儲(chǔ)服務(wù)。

相關(guān)閱讀更多精彩內(nèi)容

友情鏈接更多精彩內(nèi)容