SQL Streaming Data
Этот проект представляет собой систему для анализа потоковых данных событий веб-страницы в реальном времени. Система моделирует интернет-магазин, где пользователи генерируют события (просмотры страниц — view, клики — click), отправляемые в Apache Kafka, обрабатываемые Apache Flink с использованием SQL и анализируемые альтернативно с помощью Python. Цель — подсчёт количества событий по типам за 5-минутные временные окна для сравнения производительности подходов Flink SQL и Python.
Apache Kafka (3.6.1): Распределённая система для передачи и хранения событий.
Apache Flink (1.19.1, Scala 2.12): Платформа для потоковой обработки данных с использованием Flink SQL.
Flask-приложение: Веб-сайт для генерации событий.
Python-скрипты: Альтернативный анализ с использованием kafka-python.
Kafka: 3.6.1
Flink: 1.19.1 (Scala 2.12)
Flink Kafka Connector: 3.2.0-1.19
Python: 3.x (с библиотеками flask и kafka-python)
Операционная система: Linux (WSL 2 на Windows 10/11, Ubuntu).
Java: Oracle JDK 11 или OpenJDK 11.
Python: 3.8 или новее.
Установленные пакеты:
wget, unzip, python3, python3-pip, apt.
cd ~
wget https://kafka.apache.org/downloads/kafka_2.13-3.6.1.tgz
tar -xzf kafka_2.13-3.6.1.tgz
mv kafka_2.13-3.6.1 kafka
Убедитесь, что путь установлен (например, /home/test/kafka/).
Требуются два терминала WSL:
~/kafka/bin/zookeeper-server-start.sh ~/kafka/config/zookeeper.properties &
~/kafka/bin/kafka-server-start.sh ~/kafka/config/server.properties &
ps aux | grep kafka
~/kafka/bin/kafka-topics.sh --create --topic website_events --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
Скачайте Flink 1.19.1:
wget https://archive.apache.org/dist/flink/flink-1.19.1/flink-1.19.1-bin-scala_2.12.tgz
tar -xzf flink-1.19.1-bin-scala_2.12.tgz
mv flink-1.19.1 flink
Убедитесь, что путь установлен (например, /home/test/flink/).
Скачайте и установите flink-connector-kafka-3.2.0-1.19.jar:
wget https://repo.maven.apache.org/maven2/org/apache/flink/flink-connector-kafka/3.2.0-1.19/flink-connector-kafka-3.2.0-1.19.jar
mv flink-connector-kafka-3.2.0-1.19.jar ~/flink/lib/
Убедитесь, что версия совместима с Scala 2.12 (имя файла должно содержать _2.12).
~/flink/bin/start-cluster.sh
Проверьте веб-интерфейс: http://localhost:8081.
Установите Python и необходимые библиотеки:
sudo apt update
sudo apt install python3 python3-pip -y
pip3 install flask kafka-python
Файл ~/kafka/config/server.properties:
Убедитесь, что listeners=PLAINTEXT://:9092 настроен.
num.io.threads и num.network.threads могут быть увеличены для повышения производительности.
Файл ~/kafka/config/zookeeper.properties:
Убедитесь, что dataDir=/tmp/zookeeper настроен (или измените на другой путь).
Файл ~/flink/conf/flink-conf.yaml:
Настройте jobmanager.rpc.address: localhost.
Установите taskmanager.numberOfTaskSlots: 1 для локального тестирования.
Убедитесь, что parallelism.default: 1 для начального тестирования.
Все файлы размещены в /home/test/. Убедитесь, что пути к шаблонам (templates/) корректны в файле website.py. Код доступен в отдельных файлах: website.py, templates/home.html, templates/product.html, analyze.py (см. инструкции по использованию).
Запустите Kafka (см. раздел "Установка").
Запустите Flink (см. раздел "Установка").
Создайте или проверьте файлы в /home/test/:
Скопируйте или создайте файлы website.py, templates/home.html, templates/product.html, analyze.py (см. инструкции по установке и примеры кода в документации проекта).
Запустите веб-приложение:
python3 ~/website.py
Откройте http://localhost:5000 в браузере и взаимодействуйте с сайтом (переходите по ссылкам и кликайте).
~/flink/bin/sql-client.sh
Выполните SQL-запросы, описанные в документации проекта (доступны в отдельном файле или проекте).
python3 ~/analyze.py
Генерируйте события (1000 событий/мин) в течение 10 минут, взаимодействуя с сайтом.
Измерьте производительность:
Латентность: Используйте time для измерения времени обработки запросов в Flink SQL и Python-скриптах.
CPU: Используйте top в WSL для мониторинга использования ресурсов.
Сравните результаты Flink SQL и Python, фиксируя точность подсчётов событий и время обработки (см. примеры метрик в документации проекта).
Kafka не запускается: Проверьте, запущен ли Zookeeper, и убедитесь, что порты (2181 для Zookeeper, 9092 для Kafka) свободны. Проверьте логи в /home/test/kafka/logs/.
Flink не видит Kafka: Убедитесь, что flink-connector-kafka-3.2.0-1.19.jar находится в ~/flink/lib/ и перезапустите Flink. Проверьте логи в /home/test/flink/log/.
Ошибки JSON в Flink: Проверьте данные в Kafka через kafka-console-consumer.sh --topic website_events --bootstrap-server localhost:9092 --from-beginning. Убедитесь, что website.py отправляет корректный JSON (например, {"event_type": "view", "page": "/home", "ts": 1677050132000}).
Python не подключается к Kafka: Проверьте настройки bootstrap_servers в website.py и analyze.py. Убедитесь, что Kafka запущена и доступна на localhost:9092.
Kafka: Логи доступны в /home/test/kafka/logs/.
Flink: Логи доступны в /home/test/flink/log/.
Python: Выводятся в консоль при запуске website.py и analyze.py.