Apache Kafka на Python
Когда речь идет о потоковой передаче данных, Apache Kafka является стандартом де-факто. Это распределенная система с открытым исходным кодом, состоящая из серверов и клиентов. Apache Kafka используется в основном для построения конвейеров потоковой передачи данных в реальном времени.
Apache Kafka используется тысячами ведущими мировыми организациями для построения высокопроизводительных конвейеров данных, потоковой аналитики, интеграции данных и многих других жизненно важных приложений.
Предполагаю, что Apache Kafka уже настроен в Вашей системе.
В рамках этой статьи вам необходимо знать 4 основные концепции Kafka.
- Топик: топик в Kafka - это название категории или потока, в который публикуются сообщения в системе обмена сообщениями Kafka. Топики в Kafka похожи на таблицы в базе данных, где каждый топик состоит из одного или нескольких разделов, а каждый раздел можно представить как журнал записей. Все сообщения Kafka проходят через топики. Топик - это простой способ организовать и сгруппировать коллекцию сообщений.
- Консюмеры: консюмер в Kafka - это клиентское приложение, которое читает данные из топиков Kafka. Консюмеры подписываются на один или несколько топиков Kafka и потребляют сообщения из одного или нескольких разделов этих самых топиков.
- Продюсеры: Продюсер в Kafka - это клиентское приложение, которое записывает данные в топики Kafka. Продюсеры публикуют сообщения в одном или нескольких топиках Kafka, которые затем хранятся в разделах этих топиков. Продюсеры в Kafka могут быть реализованы с помощью низкоуровневого Kafka Producer API или высокоуровневого Kafka Streams API. Kafka Producer API обеспечивает тонкий контроль над тем, как создаются сообщения, в то время как Kafka Streams API предоставляет высокоуровневые абстракции для создания приложений потоковой обработки поверх Kafka.
- Группы консюмеров: группа консюмеров в Kafka - это группа из одного или нескольких консюмеров, которые работают вместе для того, чтобы потреблять сообщения из одного или нескольких разделов топика. Когда несколько консюмеров входят в группу консюмеров, Kafka автоматически назначает разделы для каждого консюмера в группе, гарантируя, что каждый раздел потребляется только одним консюмером. Это позволяет масштабировать консюмеров горизонтально и обеспечивать высокую доступность обработки данных.
Обратите внимание, что другие терминалы для серверов Zookeeper и Kafka должны быть запущены в фоновом режиме.
Продюсер Kafka на Python
Вот мой код продюсера с библиотекой kafka-python. Откройте файл producer.py, теперь Вы готовы к работе.
Продюсеру Kafka нужно знать, где именно запущена Kafka. Ваш экземпляр, вероятно, находится на localhost:9092, если только на этапе настройки Вы не изменили порт. Кроме того, класс KafkaProducerclass должен знать, как будут сериализовываться значения. Вы знаете ответ на оба вопроса.
Вот и все для продюсера Kafka. Вот полный исходный код:
(Я получаю потоковые данные из API с частотой 2 секунды и переформировываю ответные данные, затем передаю их в поток kafka методом producer.send.)
import json as json
import time
import requests
from concurrent.futures import ThreadPoolExecutor
from time import sleep
from kafka import KafkaProducer
# Set up Kafka producer
def serializer(message):
return json.dumps(message).encode('utf-8')
# Kafka Producer
producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
value_serializer=serializer
)
def get_push_reords():
# Define function to flatten client records
def flatten_record(server, client):
client.pop('mac', None)
client.pop('deviceHash', None)
client.pop('lon', None)
client.pop('lat', None)
client.pop('ip4', None)
if 'ip' in client:
client['client_ip'] = client.pop('ip')
client['application'] = ip_to_name.get(server['ip'], '')
client['timestamp'] = int(time.time() * 1000)
flat_record = json.dumps({"timestamp": client['timestamp'], "application": client['application'], "client_ip": client['client_ip'], "server_ip": server['ip'], "server_latlon":server['latlon'],**client})
return flat_record
# Fetch nested JSON record from Dropbox link
response = requests.get("https://www.dropbox.com/s/iwcpg1oo59i4yrn/exfoCustosTest%20%283%29.json?dl=1")
json_data = json.loads(response.content)
# Extract servers and services arrays
servers = json_data['servers']
services = json_data['services']
# Create dictionary mapping IP addresses to service names
ip_to_name = {server['ip']: service['name'] for service in services for server in service['servers']}
# Output flattened JSON record for each client record using threads
debug = True # Set to True to enable debug output
with ThreadPoolExecutor(max_workers=16) as executor:
future_results = []
for server in servers:
for client in server['clients']:
future = executor.submit(flatten_record, server, client)
future_results.append(future)
for future in future_results:
flat_record = future.result()
if debug:
print(f"Sending record to kafka Stream: {json.loads(flat_record)}")
data = json.loads(flat_record)
print(data["timestamp"])
producer.send('mec-xdr', flat_record)
producer.flush() # Wait for messages to be delivered
def main():
while True:
get_push_reords()
sleep(2)
if __name__ == "__main__":
main()
run python3 producer_from_api.py inti the terminal
Вывод:
Консюмер Kafka на Python
Вот мой consumer.py.
Консюмер Kafka будет гораздо проще. При его запуске Вы получите все сообщения из топика 'mec-xdr' и распечатаете их. Конечно, печатью сообщений дело не ограничивается - Вы можете делать все, что захотите. В моем случае я в реальном времени заношу все данные потока в БД Postgres в схему app_data в raw table.
Параметр auto_offset_reset гарантирует, что самые старые сообщения будут первыми.
Как только Вы прочитали данные в Вашем консюмере, Вы сможете делать с ними все, что Вашей душе угодно. Например, Вы можете извлечь из этих данных некоторые переменные и передать их в другую функцию для дальнейшего преобразования.
Можно написать логику и отправить оповещение при выполнении определенных условий. Можно отобразить данные в виде графика на веб-странице для того, чтобы получить график в режиме реального времени. Все Ваши дальнейшие действия ограничиваются исключительно Вашим воображением.
import json
from kafka import KafkaConsumer
import psycopg2
from sqlalchemy import create_engine
if __name__ == '__main__':
conn = psycopg2.connect(
host="localhost",
database="postgres",
user="robert",
password="Admin123"
)
# Create a cursor object
cur = conn.cursor()
# Kafka Consumer
consumer = KafkaConsumer(
'mec-xdr',
bootstrap_servers='localhost:9092',
max_poll_records = 100,
value_deserializer=lambda m: json.loads(m.decode('ascii')))
auto_offset_reset='earliest'#,'smallest'
)
for message in consumer:
print(json.loads(message.value))
cur.execute("INSERT INTO app_data.raw_data (raw) VALUES (%s)", (str(json.loads(message.value)),))
conn.commit()
run python3 consumer.py from terminal
Чтобы остановить выполнение во всех терминалах, можно использовать ctrl + c.







