ApacheKafka:在一段时间后将消息发送到另一个主题
创始人
2024-09-06 05:01:05
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)  

相关内容

热门资讯

安装apache-beam==... 出现此错误可能是因为用户的Python版本太低,而apache-beam==2.34.0需要更高的P...
避免在粘贴双引号时向VS 20... 在粘贴双引号时向VS 2022添加反斜杠的问题通常是由于编辑器的自动转义功能引起的。为了避免这个问题...
Android Recycle... 要在Android RecyclerView中实现滑动卡片效果,可以按照以下步骤进行操作:首先,在项...
omi系统和安卓系统哪个好,揭... OMI系统和安卓系统哪个好?这个问题就像是在问“苹果和橘子哪个更甜”,每个人都有自己的答案。今天,我...
原生ios和安卓系统,原生对比... 亲爱的读者们,你是否曾好奇过,为什么你的iPhone和安卓手机在操作体验上有着天壤之别?今天,就让我...
Android - 无法确定任... 这个错误通常发生在Android项目中,表示编译Debug版本的Java代码时出现了依赖关系问题。下...
Android - NDK 预... 在Android NDK的构建过程中,LOCAL_SRC_FILES只能包含一个项目。如果需要在ND...
Akka生成Actor问题 在Akka框架中,可以使用ActorSystem对象生成Actor。但是,当我们在Actor类中尝试...
Agora-RTC-React... 出现这个错误原因是因为在 React 组件中使用,import AgoraRTC from “ago...
Alertmanager在pr... 首先,在Prometheus配置文件中,确保Alertmanager URL已正确配置。例如:ale...