feat: fix
This commit is contained in:
7
external/kafka/producer.py
vendored
7
external/kafka/producer.py
vendored
@@ -1,8 +1,15 @@
|
|||||||
|
from typing import Optional
|
||||||
|
|
||||||
from aiokafka import AIOKafkaProducer
|
from aiokafka import AIOKafkaProducer
|
||||||
|
|
||||||
from backend.config import KAFKA_URL
|
from backend.config import KAFKA_URL
|
||||||
from external.kafka.context import context
|
from external.kafka.context import context
|
||||||
|
|
||||||
|
producer: Optional[AIOKafkaProducer] = None
|
||||||
|
|
||||||
|
|
||||||
|
async def init_producer():
|
||||||
|
global producer
|
||||||
producer = AIOKafkaProducer(
|
producer = AIOKafkaProducer(
|
||||||
bootstrap_servers=KAFKA_URL,
|
bootstrap_servers=KAFKA_URL,
|
||||||
security_protocol='SSL',
|
security_protocol='SSL',
|
||||||
|
|||||||
Reference in New Issue
Block a user