BI Consult Desktop Logo BI Consult Mobile Logo
  • Russian BI Исследование российских bi
  • Перейти на Fine BI
  • Контакты
  • +7 812 334-08-01
    +7 499 608-13-06
  • Отправить сообщение
  • Главная
  • Продукты Эксперт-BI
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Сельское хозяйство
    • Энергетика
    • FMCG
    • Девелоперы
    • Маркетплейсы
    • Пищевая промышленность
    • Фармацевтика
    • Построение Data Platform
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и FP&A
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • IBP
    • ИТ (CIO)
    • Закупки
  • Платформы
    • Системы бизнес-анализа (BI)
    • Интегрированное бизнес-планирование (IBP)
    • Хранилища данных (DWH / Lakehouse)
    • Каталоги данных (Data Catalog)
    • Системы ETL и ELT
    • AI / Исскуственный интеллект
    • Шина данных (ESB)
    • Система управления мастер-данными (MDM)
    • Семантический слой
  • Услуги
    • Переход на отечественные BI и DWH системы
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений и DWH
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Курсы
    • Учебный курс Информационная грамотность (Data Literacy)
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Greenplum
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt (Data Build Tool)
  • Компания
    • Руководство
    • Новости
    • Клиенты
    • Карьера
    • Скачать
    • Контакты

BI

  • FineBI
  • FineReport
  • FineDataLink
  • FineChatBI (FineAI)
  • Коннекторы данных из 1С в BI
  • Airflow / Nifi
  • Visiology
  • PIX BI
  • Modus BI
  • Yandex.DataLens
  • Open-source BI: Superset/Metabase
  • Luxms BI
  • AW BI + Alpha BI
  • FlyBI + Форсайт. Аналитическая Платформа
  • Loginom
  • Триафлай
  • AI / Исскуственный интеллект
  • Optimacros
  • Навигатор BI
  • Семантический слой

СУБД

  • Arenadata
  • ClickHouse
  • Greenplum
  • Postgres Professional
  • TData

Другое

  • Построение Data Platform
    • Аналитическое хранилище данных
    • Data Lake и Data Engineering
    • Подробнее про Data Lake
    • Внедрение Lakehouse
      • Apache Doris
      • StarRocks
      • Trino
    • Миграция витрин из пропиетарных DWH на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс Современная архитектура хранилища данных » Apache Kafka на Python

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.

 

Узнать стоимость решенияЗапросить видео презентацию

← Предыдущая статья
Выбор клиента Python Kafka: сравнительный анализ
Следующая статья →
Потоковая обработка данных с помощью Python и Apache Kafka: руководство для новичков

Решения

Анализировать ФинансыУвеличивайте ПродажиОптимальный Склад и ЛогистикаМаркетинговые Метрики

Клиенты
  • ПАО «Ростелеком» — российский провайдер цифровых услуг и сервисов. Предоставляет услуги широкополосного доступа в Интернет, интерактивного телевидения, сотовой связи, местной и дальней телефонной связи и др. Занимает лидирующие позиции на российском рынке высокоскоростного доступа в интернет, платного ТВ, хранения и обработки данных, а также кибербезопасности

  • KERAMA MARAZZI — международный бренд, входящий в число лидеров глобального рынка керамики. Бизнес компании охватывает весь процесс создания керамических изделий, от глиняных карьеров до фирменной розницы во всех крупных городах РФ и за рубежом.

  • Торгово-производственному холдингу ТБМ, специализирующемуся на поставке комплектующих и фурнитуры для производства окон, дверей, стеклопакетов и мебели, был необходим аналитический инструмент для выявления узким мест и поиска зон роста бизнеса и, как результат, оптимизации процессов. Добиться этого можно было, только внедрив data-driven подход.

  • Банк "Санкт-Петербург" - это универсальный коммерческий банк, предоставляющий полный спектр финансовых услуг для частных и корпоративных клиентов. Банк основан в 1990 году и имеет генеральную лицензию Банка России на осуществление банковских операций. Сеть банка включает более 170 офисов и отделений, а также свыше 1000 банкоматов и терминалов в Санкт-Петербурге, Москве и других регионах.

  • Решения
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Энергетика
    • Фармацевтика
  • Услуги
    • Переход на отечественные BI и DWH
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Техническая поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Платформы
    • FineBI
    • FineReport
    • FineDataLink
    • Коннекторы данных из 1С в BI
    • Airflow + NiFi
    • Visiology
    • Luxms BI
    • Modus BI
    • PIX BI
    • Arenadata
    • ClickHouse
    • Greenplum
    • Postgres Professional
    • Open-source BI: Superset/Metabase
    • Loginom
    • Yandex.DataLens
    • AI / Исскуственный интеллект
    • Optimacros
    • Шины данных
  • Курсы
    • Учебный курс Информационная грамотность
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt
  • Функциональные решения
    • Создание Data Lake
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и прогнозная аналитика
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • Сквозная аналитика
  • Компания
    • О нас
    • Руководство
    • Новости
    • Клиенты
    • Скачать
    • Контакты
    • Политика конфиденциальности
RutubeVkontakteLinkedInYouTube
ООО "Би Ай Консалт",
ИНН: 7811437757,
ОГРН: 1097847154184
199178, Россия,
Санкт-Петербург,
6-ая линия В.О., Д. 63, 4 этаж
Тел: +7 (812) 334-08-01
Тел: +7 (499) 608-13-06
E-mail: info@biconsult.ru

 

 

 

 

 

×

Пользуясь сайтом, вы соглашаетесь с использованием cookies и политикой конфиденциальности.