@@ -66,8 +66,7 @@ class KafkaService:
|
||||
def produce(self, job_data: JobData) -> bool:
|
||||
"""发送消息到Kafka"""
|
||||
try:
|
||||
data = job_data.to_dict()
|
||||
future = self.producer.send(self.topic, key=data.get("_id"), value=data)
|
||||
future = self.producer.send(self.topic, key=job_data.id, value=job_data.model_dump())
|
||||
future.get(timeout=10)
|
||||
return True
|
||||
except KafkaError as e:
|
||||
|
||||
Reference in New Issue
Block a user