Python向Kafka發消息
后端研發可以提供一個向kafka發消息的接口,用requests向接口post消息就行:
import requests
import json
import time
now = int(time.time())
n = 10
while n > 0:
tt = now - n * 60
data = {
"queue": "alarm-dog-alarm-dog-test",
"payload": "{\"test\":80,\"notice_time\":%d}" % tt
}
header = {"Content-Type": "application/json"}
res = requests.post(url="http://10.90.100.130:8088/v1/kafka/send", headers=header, data=json.dumps(data))
print(res.status_code)
print(res.content)
n -= 1
如果沒有提供接口,可以借助python-kafka庫連接kafka,模擬生產者向kafka發消息:
同步發送消息:
from kafka import KafkaProducer
import json
# 創建一個KafkaProducer實例,指定Kafka服務器地址
producer = KafkaProducer(bootstrap_servers='http://10.90.100.130:8088')
# 要發送的消息內容
message = {'test': 80, 'notice_time': 5}
# 將消息轉換為JSON字符串格式(也可以是其他格式,如純文本)
message_json = json.dumps(message)
# 發送消息到指定的Kafka主題,這里主題名稱是'my_topic'
producer.send('alarm-dog-alarm-dog-test', value=message_json.encode('utf - 8'))
# 確保所有消息都已發送
producer.flush()
# 關閉生產者連接
producer.close()
異步發送消息
from kafka import KafkaProducer
import json
import time
# 創建一個KafkaProducer實例,設置異步發送和回調函數
producer = KafkaProducer(bootstrap_servers='http://10.90.100.130:8088',
acks='all',
retries=3,
value_deliver_callback=lambda m: print(f"消息已發送到主題{m.topic()},分區{m.partition()}"))
# 要發送的消息內容
message = {'test': 80, 'notice_time': 6}
message_json = json.dumps(message)
# 異步發送消息到'my_topic'主題
future = producer.send('alarm-dog-alarm-dog-test', value=message_json.encode('utf - 8'))
try:
record_metadata = future.get(timeout=10)
print(f"消息已發送到主題{record_metadata.topic()},分區{record_metadata.partition()},偏移量{record_metadata.offset()}")
except Exception as e:
print(f"發送消息時出錯: {e}")
# 關閉生產者連接
producer.close()

浙公網安備 33010602011771號