Упрощение конвейеров данных: интеграция Kafka с Airflow
Упрощение конвейеров данных: интеграция Kafka с Airflow

Apache Kafka

Apache Kafka — это распределенная платформа потоковой передачи событий, которая отличается масштабируемостью, надежностью и отказоустойчивостью. Он действует как брокер сообщений и поддерживает публикацию в реальном времени и подписку на потоки записей. Его архитектура обеспечивает передачу данных с высокой пропускной способностью и низкой задержкой, что делает его лучшим выбором для обработки больших объемов данных в реальном времени в нескольких приложениях.

Apache Airflow

Apache Airflow — это платформа с открытым исходным кодом, предназначенная для организации сложных рабочих процессов. Это облегчает планирование, мониторинг и управление рабочими процессами с помощью направленных ациклических графов (DAG). Модульная архитектура Airflow поддерживает различные интеграции, что делает ее популярной в отрасли для работы с конвейерами данных.

Интегрируйте Kafka с Airflow

KafkaProducerOperator и KafkaConsumerOperator

Давайте углубимся в использование специального оператора Интегрируйте Kafka с Airflow.

Пример KafkaProducerOperator:

Рассмотрим сценарий, в котором данные датчиков необходимо опубликовать в теме Kafka. Расход воздуха

KafkaProducerOperatorЭтого можно достичь:

Язык кода:javascript
копировать
from airflow.providers.apache.kafka.operators.kafka import KafkaProducerOperator

publish_sensor_data = KafkaProducerOperator(
    task_id='publish_sensor_data',
    topic='sensor_data_topic',
    bootstrap_servers='kafka_broker:9092',
    messages=[
        {'sensor_id': 1, 'temperature': 25.4},
        {'sensor_id': 2, 'temperature': 28.9},
        # More data to be published
    ],
    # Add more configurations as needed
)
Язык кода:javascript
копировать
KafkaConsumerOperator Пример:

Допустим, мы хотим получить данные из темы Kafka и выполнить анализ:

Язык кода:javascript
копировать
from airflow.providers.apache.kafka.operators.kafka import KafkaConsumerOperator

consume_and_analyze_data = KafkaConsumerOperator(
    task_id='consume_and_analyze_data',
    topic='sensor_data_topic',
    bootstrap_servers='kafka_broker:9092',
    group_id='airflow-consumer',
    # Add configurations and analytics logic
)

Создайте конвейер данных

Демонстрирует упрощенный конвейер данных с использованием Airflow DAG и интеграцией в него Kafka.

Язык кода:javascript
копировать
from airflow import DAG
from airflow.providers.apache.kafka.operators.kafka import KafkaProducerOperator, KafkaConsumerOperator
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2023, 12, 1),
    # Add more necessary arguments
}

with DAG('kafka_airflow_integration', default_args=default_args, schedule_interval='@daily') as dag:
    publish_sensor_data = KafkaProducerOperator(
        task_id='publish_sensor_data',
        topic='sensor_data_topic',
        bootstrap_servers='kafka_broker:9092',
        messages=[
            {'sensor_id': 1, 'temperature': 25.4},
            {'sensor_id': 2, 'temperature': 28.9},
            # More data to be published
        ],
        # Add more configurations as needed
    )

    consume_and_analyze_data = KafkaConsumerOperator(
        task_id='consume_and_analyze_data',
        topic='sensor_data_topic',
        bootstrap_servers='kafka_broker:9092',
        group_id='airflow-consumer',
        # Add configurations and analytics logic
    )

    publish_sensor_data >> consume_and_analyze_data  # Define task dependencies

Лучшие практики и соображения

  • обратная сериализация: убедитесь, что обратная сериализация данных соответствует ожиданиям Кафки в отношении бесперебойной связи между производителями и потребителями.
  • Мониторинг и ведение журнала. Внедрите мощный механизм мониторинга и ведения журнала для отслеживания потоков данных и решения потенциальных проблем существования в конвейере.
  • Меры безопасности: расставьте приоритеты безопасности, внедрив протоколы шифрования и аутентификации для защиты сквозной передачи данных. Kafka существовать Airflow Передача в данных.

в заключение

добавив Apache Kafka и Apache Airflow Интегрированные инженеры по обработке данных имеют доступ к мощной экосистеме для создания эффективных конвейеров данных в режиме реального времени. Кафка Высокая пропускная способность и Airflow Комбинация оркестрации рабочих процессов позволяет предприятиям создавать сложные конвейеры для удовлетворения современных потребностей в обработке данных.

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

Автор оригинала: Лукас Фонсека

boy illustration
Неразрушающее увеличение изображений одним щелчком мыши, чтобы сделать их более четкими артефактами искусственного интеллекта, включая руководства по установке и использованию.
boy illustration
Копикодер: этот инструмент отлично работает с Cursor, Bolt и V0! Предоставьте более качественные подсказки для разработки интерфейса (создание навигационного веб-сайта с использованием искусственного интеллекта).
boy illustration
Новый бесплатный RooCline превосходит Cline v3.1? ! Быстрее, умнее и лучше вилка Cline! (Независимое программирование AI, порог 0)
boy illustration
Разработав более 10 проектов с помощью Cursor, я собрал 10 примеров и 60 подсказок.
boy illustration
Я потратил 72 часа на изучение курсорных агентов, и вот неоспоримые факты, которыми я должен поделиться!
boy illustration
Идеальная интеграция Cursor и DeepSeek API
boy illustration
DeepSeek V3 снижает затраты на обучение больших моделей
boy illustration
Артефакт, увеличивающий количество очков: на основе улучшения характеристик препятствия малым целям Yolov8 (SEAM, MultiSEAM).
boy illustration
DeepSeek V3 раскручивался уже три дня. Сегодня я попробовал самопровозглашенную модель «ChatGPT».
boy illustration
Open Devin — инженер-программист искусственного интеллекта с открытым исходным кодом, который меньше программирует и больше создает.
boy illustration
Эксклюзивное оригинальное улучшение YOLOv8: собственная разработка SPPF | SPPF сочетается с воспринимаемой большой сверткой ядра UniRepLK, а свертка с большим ядром + без расширения улучшает восприимчивое поле
boy illustration
Популярное и подробное объяснение DeepSeek-V3: от его появления до преимуществ и сравнения с GPT-4o.
boy illustration
9 основных словесных инструкций по доработке академических работ с помощью ChatGPT, эффективных и практичных, которые стоит собрать
boy illustration
Вызовите deepseek в vscode для реализации программирования с помощью искусственного интеллекта.
boy illustration
Познакомьтесь с принципами сверточных нейронных сетей (CNN) в одной статье (суперподробно)
boy illustration
50,3 тыс. звезд! Immich: автономное решение для резервного копирования фотографий и видео, которое экономит деньги и избавляет от беспокойства.
boy illustration
Cloud Native|Практика: установка Dashbaord для K8s, графика неплохая
boy illustration
Краткий обзор статьи — использование синтетических данных при обучении больших моделей и оптимизации производительности
boy illustration
MiniPerplx: новая поисковая система искусственного интеллекта с открытым исходным кодом, спонсируемая xAI и Vercel.
boy illustration
Конструкция сервиса Synology Drive сочетает проникновение в интрасеть и синхронизацию папок заметок Obsidian в облаке.
boy illustration
Центр конфигурации————Накос
boy illustration
Начинаем с нуля при разработке в облаке Copilot: начать разработку с минимальным использованием кода стало проще
boy illustration
[Серия Docker] Docker создает мультиплатформенные образы: практика архитектуры Arm64
boy illustration
Обновление новых возможностей coze | Я использовал coze для создания апплета помощника по исправлению домашних заданий по математике
boy illustration
Советы по развертыванию Nginx: практическое создание статических веб-сайтов на облачных серверах
boy illustration
Feiniu fnos использует Docker для развертывания личного блокнота Notepad
boy illustration
Сверточная нейронная сеть VGG реализует классификацию изображений Cifar10 — практический опыт Pytorch
boy illustration
Начало работы с EdgeonePages — новым недорогим решением для хостинга веб-сайтов
boy illustration
[Зона легкого облачного игрового сервера] Управление игровыми архивами
boy illustration
Развертывание SpringCloud-проекта на базе Docker и Docker-Compose