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 на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс Современная архитектура хранилища данных » Потоковая обработка данных с помощью Python и Apache Kafka: руководство для новичков

Потоковая обработка данных с помощью Python и Apache Kafka: руководство для новичков

Введение

В современном мире, основанном на данных, обработка больших объемов данных в режиме реального времени становится все более и более важной.

Одной из самых популярных технологий для потоковой обработки данных является Apache Kafka.

В этой статье мы будем использовать Apache Kafka и Python для создания простого и эффективного приложения для потоковой обработки данных.

 

Давайте сначала рассмотрим, что именно представляет собой Apache Kafka на самом высоком уровне.

 

Введение в Apache Kafka

Apache Kafka - это распределенная и масштабируемая платформа потоковой передачи данных, способная обрабатывать миллиарды событий в день.

Она предназначена для обработки больших объемов данных в режиме реального времени с высокой пропускной способностью и низкой задержкой.

Kafka часто используется для анализа потоковых данных в режиме реального времени и считается одной из самых быстрых и надежных потоковых платформ.

Используется организациями любого размера для создания конвейеров данных, работающих в режиме  реального времени, и приложений потоковой обработки данных.

 

Начало работы с  Kafka

Для успешного начала работы с Apache Kafka ее необходимо правильно установить.

Процесс установки очень прост, инструкцию по установке Вы найдете на сайте Apache Kafka.

После установки Вы сможете начать работать с Kafka, создавая топики, а также создавая и потребляя сообщения.

 

Создание и потребление сообщений в Kafka

Kafka использует модель pub-sub, в которой продюсеры публикуют сообщения в топики, а консюмеры подписываются на топики для того, чтобы получать сообщения.

В этом руководстве мы сосредоточимся на потреблении сообщений из Kafka с помощью Python.

Для этого мы будем использовать библиотеку python-kafka, которая предоставляет высокоуровневый API для работы с Apache Kafka.

Сначала Вам нужно установить брокер Apache Kafka и клиентскую библиотеку Apache Kafka.

Существует несколько клиентских библиотек для Apache Kafka, включая Java, Python и другие, однако в этом посте мы будем использовать только клиентскую библиотеку Python.

 

Установка библиотека python-kafka

Первым шагом к получению сообщений из Kafka является установка библиотеки python-kafka. Вы можете установить библиотеку, выполнив следующую команду:

pip install kafka-python

 

Потребление сообщений из Kafka

После установки библиотеки python-kafka Вы можете начать потреблять сообщения из Kafka.

Следующий код представляет собой простой пример того, как можно потреблять сообщения из топика Kafka с помощью Python:

from kafka import KafkaConsumer

consumer = KafkaConsumer('my_topic',
                         bootstrap_servers=['localhost:9092'],
                         auto_offset_reset='earliest',
                         enable_auto_commit=False,
                         group_id='my_group_id',
                         value_deserializer=lambda x: x.decode('utf-8')
                        )

 

для сообщения в консюмере:

print("Received message: ", message.value)

 

В приведенном выше коде мы создаем объект KafkaConsumer и подписываемся на топик«my_topic».

Цикл for используется для перебора сообщений в топике, значение каждого сообщения выводится в консоль.

Работа с данными в режиме реального времени с помощью Python и Apache Kafka

В реальном мире вам придется работать с большими объемами данных, которые генерируются в режиме реального времени, и это один из самых популярных вариантов использования Apache Kafka!

Сейчас мы рассмотрим способ использования Python и Apache Kafka для обработки данных в режиме реального времени.

Для того, чтобы написать консюмера Kafka на Python, нам сначала нужно создать объект KafkaConsumer и указать необходимые параметры, такие как топик, брокер и ID группы.

Затем мы можем использовать метод subscribe для подписки на один или сразу несколько топиков.

Наконец, для опроса новых данных и их обработки по мере необходимости мы можем использовать метод poll.

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

from kafka import KafkaConsumer

# Create a KafkaConsumer instance
consumer = KafkaConsumer(
    bootstrap_servers=['localhost:9092'],
    auto_offset_reset='earliest',
    enable_auto_commit=False
)

# Subscribe to a specific topic
consumer.subscribe(topics=['my-topic'])

# Poll for new messages
while True:
    msg = consumer.poll(timeout_ms=1000)
    if msg:
        for topic, partition, offset, key, value in msg.items():
            print("Topic: {} | Partition: {} | Offset: {} | Key: {} | Value: {}".format(
                topic, partition, offset, key, value.decode("utf-8")
            ))
    else:
        print("No new messages")

 

В этом примере мы сначала создаем экземпляр KafkaConsumer и настраиваем свойство bootstrap_servers так, чтобы оно указывало на адрес брокера Kafka.

Свойство auto_offset_reset установлено в earliest так, чтобы консюмер начинал с самого раннего сообщения в топике.

Свойство enable_auto_commit установлено в False для того, чтобы консюмер не фиксировал смещение после обработки сообщения.

Далее мы подписываемся на определенный топик с помощью метода subscribe.

Наконец, для опроса новых сообщений мы используем метод poll в цикле while. Метод poll возвращает словарь сообщений, где ключи представляют собой топик, раздел, смещение и ключ сообщения.

Доступ к значению сообщения осуществляется с помощью свойства value. В этом примере мы декодируем значение из байтов в кодировку «utf-8» с помощью метода decode.

Это всего лишь пример использования методов subscribe и poll в консюмере Python Kafka.

В классе KafkaConsumer доступно множество дополнительных функций и опций настройки, которые Вы можете найти в официальной документации.

Здесь представлен более простой и базовый вариант написания консюмера:

from kafka import KafkaConsumer
import pandas as pd

consumer = KafkaConsumer('my_topic',
                         bootstrap_servers=['kafka_broker:9092'],
                         auto_offset_reset='earliest',
                         enable_auto_commit=False,
                         group_id='my_group_id',
                         value_deserializer=lambda x: x.decode('utf-8')
                        )

 

для сообщения  консюмере:

    data = message.value.decode("utf-8")
    df = pd.read_json(data)
    print("Received data: ", df.head())

 

Написание продюсера Kafka на Python

Хотя основное внимание в этом посте было уделено консюмеру Kafka, давайте все же вкратце затронем и продюсерскую часть.

Помимо потребления данных из Apache Kafka, мы также можем записывать данные в поток Apache Kafka с помощью продюсера Kafka.

Чтобы написать продюсера Kafka на Python, нам сначала нужно создать объект KafkaProducer и указать необходимые параметры, например, брокера.

Затем мы можем использовать метод send для отправки данных брокеру.

from kafka import KafkaProducer

producer = KafkaProducer(
    bootstrap_servers=["kafka_broker:9092"]
)

producer.send("topic_name", value="Hello, World!".encode("utf-8"))

 

Заключение

В этой статье мы рассмотрели, как использовать Python с Apache Kafka для обработки потоков.

Apache Kafka предоставляет масштабируемую и распределенную архитектуру для обработки данных в режиме реального времени, а Python - простой и удобный язык программирования для разработки консюмеров и продюсеров Kafka.

Вооружившись этими инструментами и базовыми понятиями, Вы можете приступить к созданию собственных мощных и масштабируемых приложений для обработки данных, способных обрабатывать большие потоки данных в режиме реального времени для решения Ваших задач.

 

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

← Предыдущая статья
Apache Kafka на Python
Следующая статья →
Коннекторы Kafka

Решения

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

Клиенты
  • Компания «Бизон-Трейд» является официальным дилером ведущих мировых производителей сельскохозяйственной техники (Fendt, Valtra, Lemken и др.) на Юге России. Входит в состав агрохолдинга «Бизон», основанного в 1994 году. Имеет 8 филиалов в Краснодарском и Ставропольском краях, Ростовской области.

  • Компания ООО "Комус" - один из лидеров российского рынка оптовых продаж офисных товаров и техники. Компания поставляет широкий ассортимент продукции - от канцелярских принадлежностей до компьютерной техники и офисной мебели.

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

  • СберКорус (Группа компаний Сбербанка) – это ИТ‑компания, ИТ‑интегратор, SaaS-провайдер. Является разработчиком цифровых сервисов и услуг для автоматизации широкого диапазона бизнес-процессов юридических лиц. В 2004 году компания стала первым в России оператором электронного документооборота, а в 2012 году вошла в экосистему Сбера. 

  • Решения
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • 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 и политикой конфиденциальности.