python?實(shí)時(shí)獲取kafka消費(fèi)隊(duì)列信息示例詳解
安裝 pykafka
pip install pykafka
一、消費(fèi)kafka消息
#!/usr/bin/env python
# -*- coding: utf-8 -*-
from pykafka import KafkaClient
from pykafka.common import OffsetType
from vpn_data_handler import handler_data
bootstrap_servers = '10.*.**.**:9092'
group_id = 'test1'
class KConsumer(object):
"""kafka 消費(fèi)者; 動(dòng)態(tài)傳參,非配置文件傳入;
kafka 的消費(fèi)者應(yīng)該盡量和生產(chǎn)者保持在不同的節(jié)點(diǎn)上;否則容易將程序陷入死循環(huán)中;
"""
_encode = "UTF-8"
def __init__(self, topics, bootstrap_server=None, group_id=group_id, partitions=None):
""" 初始化kafka的消費(fèi)者;
1. 設(shè)置默認(rèn) kafka 的主題, 節(jié)點(diǎn)地址, 消費(fèi)者組 id(不傳入的時(shí)候使用默認(rèn)的值)
2. 當(dāng)需要設(shè)置特定參數(shù)的時(shí)候可以直接在 kwargs 直接傳入,進(jìn)行解包傳入原始函數(shù);
Args:
topics: str; kafka 的消費(fèi)主題;
bootstrap_server: list; kafka 的消費(fèi)者地址;
group_id: str; kafka 的消費(fèi)者分組 id,默認(rèn)是 start_task 主要是接收并啟動(dòng)任務(wù)的消費(fèi)者,僅此一個(gè)消費(fèi)者組id;
"""
if bootstrap_server is None:
bootstrap_server = bootstrap_servers
self.client = KafkaClient(hosts=bootstrap_server)
# 選擇要消費(fèi)的topic
vpn_topic = self.client.topics[topics]
self.consumer = vpn_topic.get_simple_consumer(consumer_group=group_id,
consumer_timeout_ms=200,
auto_commit_enable=True,# 自動(dòng)提交偏移量
auto_offset_reset=OffsetType.LATEST) #LATEST 獲取當(dāng)前偏移量最新消息 EARLIEST從頭開始獲取信息
def recv(self):
"""
接收消費(fèi)中的數(shù)據(jù)
Returns:
"""
return self.consumer
def main():
"""
kafka消費(fèi)隊(duì)列入口
:param topic:
:return:
"""
obj = KConsumer(topics="topics_name")
while True:
for message in obj.recv():
data = eval(message.value.decode('utf-8'))
handler_data(data)
if __name__ == '__main__':
main()二、生產(chǎn)者推送消息
#!/usr/bin/python
# -*- coding:utf-8 -*-
from pykafka import KafkaClient
client = KafkaClient(hosts="10.XX0.XX0.XX4:9092") # 可接受多個(gè)client
# 查看所有的topic
# print(client.topics)
topic = client.topics['test_78'] # 選擇一個(gè)topic
message = "test message2 test message2"
with topic.get_sync_producer() as producer:
producer.produce(bytes(message, encoding='utf8')) #python3需要編碼
print(message)到此這篇關(guān)于python 實(shí)時(shí)獲取kafka消費(fèi)隊(duì)列信息的文章就介紹到這了,更多相關(guān)python kafka消費(fèi)隊(duì)列信息內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Python微服務(wù)開發(fā)之使用FastAPI構(gòu)建高效API
微服務(wù)架構(gòu)在現(xiàn)代軟件開發(fā)中日益普及,它將復(fù)雜的應(yīng)用程序拆分成多個(gè)可獨(dú)立部署的小型服務(wù)。本文將介紹如何使用 Python 的 FastAPI 庫快速構(gòu)建和部署微服務(wù),感興趣的可以了解一下2023-05-05
Python格式化處理JSON數(shù)據(jù)的完整指南
在Python中,我們經(jīng)常需要處理JSON數(shù)據(jù),而格式化JSON數(shù)據(jù)是開發(fā)過程中的常見需求,本文將詳細(xì)介紹如何在Python中對(duì)JSON數(shù)據(jù)進(jìn)行格式化處理,感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下2026-04-04
pandas計(jì)數(shù) value_counts()的使用
這篇文章主要介紹了pandas計(jì)數(shù) value_counts()的使用,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2019-06-06
Python二進(jìn)制文件讀取并轉(zhuǎn)換為浮點(diǎn)數(shù)詳解
這篇文章主要介紹了Python二進(jìn)制文件讀取并轉(zhuǎn)換為浮點(diǎn)數(shù)詳解,用python讀取二進(jìn)制文件,這里主要用到struct包,而這個(gè)包里面的方法主要是unpack、pack、calcsize。,需要的朋友可以參考下2019-06-06
Python實(shí)現(xiàn)線程池之線程安全隊(duì)列
這篇文章主要為大家詳細(xì)介紹了Python實(shí)現(xiàn)線程池之線程安全隊(duì)列,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2022-05-05
Python實(shí)例方法與類方法和靜態(tài)方法介紹與區(qū)別分析
在 Python 中,實(shí)例方法(instance method),類方法(class method)與靜態(tài)方法(static method)經(jīng)常容易混淆。本文通過代碼例子來說明它們的區(qū)別2022-10-10

