В настоящее время все работает, и в докере не появляется никаких ошибок, но при запуске кода у меня появляется ошибка NoBrokersAvailable.
Вот код:
Код: Выделить всё
from kafka import KafkaProducer
import requests
import time
import json
API_KEY = 'my_api_key'
URL = f'https://api.openweathermap.org/data/2.5/weather?lat=44.34&lon=10.99&appid={API_KEY}'
KAFKA_TOPIC = 'real-time-weather'
producer = KafkaProducer(
bootstrap_servers='localhost:9092',
value_serializer=lambda v: json.dumps(v).encode('utf-8'),
)
def get_weather_data():
response = requests.get(URL)
return response.json()
while True:
weather_data = get_weather_data()
producer.send(KAFKA_TOPIC, value=weather_data)
print(f'Send to Kafka : {weather_data}')
time.sleep(60)
Код: Выделить всё
NoBrokersAvailable Traceback (most recent call last)
Cell In[1], line 10
7 URL = f'https://api.openweathermap.org/data/2.5/weather?lat=44.34&lon=10.99&appid={API_KEY}'
8 KAFKA_TOPIC = 'real-time-weather'
---> 10 producer = KafkaProducer(
11 bootstrap_servers='localhost:9092',
12 value_serializer=lambda v: json.dumps(v).encode('utf-8'),
13 )
16 def get_weather_data():
17 response = requests.get(URL)
File /opt/conda/lib/python3.11/site-packages/kafka/producer/kafka.py:381, in KafkaProducer.__init__(self, **configs)
378 reporters = [reporter() for reporter in self.config['metric_reporters']]
379 self._metrics = Metrics(metric_config, reporters)
--> 381 client = KafkaClient(metrics=self._metrics, metric_group_prefix='producer',
382 wakeup_timeout_ms=self.config['max_block_ms'],
383 **self.config)
385 # Get auto-discovered version from client if necessary
386 if self.config['api_version'] is None:
File /opt/conda/lib/python3.11/site-packages/kafka/client_async.py:244, in KafkaClient.__init__(self, **configs)
242 if self.config['api_version'] is None:
243 check_timeout = self.config['api_version_auto_timeout_ms'] / 1000
--> 244 self.config['api_version'] = self.check_version(timeout=check_timeout)
File /opt/conda/lib/python3.11/site-packages/kafka/client_async.py:900, in KafkaClient.check_version(self, node_id, timeout, strict)
898 if try_node is None:
899 self._lock.release()
--> 900 raise Errors.NoBrokersAvailable()
901 self._maybe_connect(try_node)
902 conn = self._conns[try_node]
NoBrokersAvailable: NoBrokersAvailable
Подробнее здесь: https://stackoverflow.com/questions/790 ... rk-project