ApacheKafka:在一段时间后将消息发送到另一个主题
创始人
2024-09-06 05:00:25
0

您可以使用Kafka的时间管理器(TimeBasedUUID)函数将消息的延迟时间与新消息的主题一起存储在Kafka中。 在此之后,将启动一个Kafka Consumer Group,该组将,根据您设置的时间,在指定时间后消耗并发送该消息。以下是Kafka消息延迟发送的代码示例。

在这个示例中,使用UUID来生成具有唯一ID的随机字符串。消息格式将包含(topic,message,delay)。在示例中,我们设置了一个延迟时间为3秒钟。

producer.py:

import time
import json    
from kafka import KafkaProducer    
from uuid import uuid4

producer = KafkaProducer(bootstrap_servers='localhost:9092')

def send_to_kafka(topic, data, delay):
    obj = {"data": data, "topic": topic, "delay": delay}
    future = producer.send("delayed-messages", json.dumps(obj).encode('utf-8'), key=str(uuid4()).encode('utf-8'))
    result = future.get(timeout=60)
    print(result)

if __name__ == "__main__":
    send_to_kafka("my-topic", {"id": 1, "message": "Hello World!"}, 3)

consumer.py:

import time    
import json    
from kafka import KafkaConsumer, KafkaProducer    
from datetime import datetime, timedelta    
from uuid import UUID, uuid4  
  
KAFKA_TOPIC = "delayed-messages"  
KAFKA_CONSUMER_GROUP = "delayed-consumer-group"  
KAFKA_SERVER = "localhost:9092"

consumer = KafkaConsumer(KAFKA_TOPIC, bootstrap_servers=KAFKA_SERVER, groupId=KAFKA_CONSUMER_GROUP)    
producer = KafkaProducer(bootstrap_servers='localhost:9092')

for message in consumer:  
    value = message.value  
    
    data = json.loads(value.decode('utf-8'))  
    currtime = datetime.now()  
    timestamp = (currtime + timedelta(seconds=data["delay"])).strftime('%Y-%m-%d %H:%M:%S')

    if currtime >= datetime.strptime(timestamp, '%Y-%m-%d %H:%M:%S'):  
        print("Message sent: "  + str(value))  
        future = producer.send(data["topic"], json.dumps(data["data"]).encode('utf-8'))  
        result = future.get(timeout=60)  
        print(result)  

相关内容

热门资讯

八分钟了解!一起宁德游戏钓蟹输... 八分钟了解!一起宁德游戏钓蟹输赢规律,白金岛跑得快辅助工具,黑科技教程(有挂脚本)小薇(透视辅助)致...
8分钟了解!衡阳丫丫字牌外 挂... 8分钟了解!衡阳丫丫字牌外 挂,拱趴大菠萝切牌规律,wpk教程(有挂普及)1、完成拱趴大菠萝切牌规律...
八分钟了解!广西跑得快助赢神器... 八分钟了解!广西跑得快助赢神器购买,赣牌圈开挂是真的吗,AA德州教程(有挂方法);1、超多福利:超高...
一分钟了解!大凉山生活号跑得快... 一分钟了解!大凉山生活号跑得快有挂吗,哈局十三张安卓辅助,玩家教你(有挂教学)在进入大凉山生活号跑得...
六分钟了解!闽悦麻将是不是有挂... 六分钟了解!闽悦麻将是不是有挂,花花生活圈怎么老是输,教你攻略(有挂工具)花花生活圈怎么老是输辅助器...
一分钟了解!小程序的雀神麻将怎... 一分钟了解!小程序的雀神麻将怎么玩才会赢,中至窝龙如何提高自己的胜率,专业教程(有挂神器)1、在小程...
3分钟了解!皮皮斗地主外 挂,... 3分钟了解!皮皮斗地主外 挂,兴动棋牌麻将有挂吗,解密教程(有挂插件)兴动棋牌麻将有挂吗辅助器中分为...
七分钟了解!胡乐辅助器免费版,... 七分钟了解!胡乐辅助器免费版,掌心圈麻将有挂是真的吗,详细教程(有挂解说)一、掌心圈麻将有挂是真的吗...
8分钟了解!随意玩拼三张能破解... 8分钟了解!随意玩拼三张能破解吗,中至麻将发牌规律,攻略方法(有挂科普)1、玩家可以在随意玩拼三张能...
二分钟了解!蜂娱棋牌2有挂吗,... 二分钟了解!蜂娱棋牌2有挂吗,拱趴十三水输赢规律,德州教程(有挂辅助)1.拱趴十三水输赢规律 ai辅...