优化代码

This commit is contained in:
2026-01-29 15:10:58 +00:00
parent 0e9c7dfe85
commit 0bd4df686c
4 changed files with 17 additions and 6 deletions

View File

@@ -1,7 +1,8 @@
from aiokafka import AIOKafkaProducer
import os
kafka_nodes = os.environ.get("kafka.nodes","redpanda-0:19092")
# Kafka 配置
KAFKA_BOOTSTRAP_SERVERS = '127.0.0.1:19092'
KAFKA_BOOTSTRAP_SERVERS = kafka_nodes
# KAFKA_TOPIC = 'binance_agg_trade'
class KafkaProducerSingleton:
@@ -54,4 +55,4 @@ async def send_to_kafka(topic, data, key=None, send_now=False):
producer = await _kafka_producer_singleton.get_producer2()
else:
producer = await _kafka_producer_singleton.get_producer()
await producer.send(topic, data,key=key)
await producer.send(topic, data,key=key)