Yandex Cloud
Поиск
Связаться с экспертомПопробовать бесплатно
  • Кейсы
  • Документация
  • Блог
  • Все сервисы
    • Cloud Interconnect
    • Cloud Backup
    • Compute Cloud
    • Object Storage
    • Managed Service for Kubernetes®
    • Managed Service for PostgreSQL
    • Managed Service for ClickHouse®
    • Monium
    • Cloud CDN
    • Network Load Balancer
    • Virtual Private Cloud
    • Cloud DNS
    • Application Load Balancer
    • Yandex Cloud Router
    • Managed Service for MySQL®
    • Managed Service for Valkey™
    • Managed Service for Apache Spark™
    • Managed Service for OpenSearch
    • Managed Service for Apache Kafka®
    • Data Transfer
    • Yandex MPP Analytics Engine for PostgreSQL
    • Managed Service for YDB
    • SpeechKit
    • Yandex Identity Hub
    • Key Management Service
    • Certificate Manager
    • Yandex Lockbox
    • Audit Trails
    • Container Registry
    • Managed Service for Prometheus®
    • Message Queue
    • Identity and Access Management
    • Yandex Cloud Console
    • Resource Manager
    • Yandex Cloud Billing
    • Cloud Apps
    • Yandex AI Studio
    • Yandex BareMetal
    • Smart Web Security
    • Security Deck
    • Yandex Cloud Video
    • Stackland
    • Yandex Managed Service for Apache Airflow®
    • Data Processing
    • Yandex MetaData Hub
    • Yandex WebSQL
    • DataLens
    • Yandex Search API
    • SpeechSense
    • DataSphere
    • Vision OCR
    • Translate
    • Cloud Registry
    • SmartCaptcha
    • Cloud Desktop
    • Yandex SIEM
    • SourceCraft Code Assistant
    • Managed Service for GitLab
    • Cloud Functions
    • API Gateway
    • Yandex Cloud Postbox
    • Serverless Integrations
    • IoT Core
    • Serverless Containers
    • Cloud Notification Service
    • Yandex Query
  • Статус работы сервисов
  • Marketplace
    • Доступны в регионе
    • Инфраструктура и сеть
    • Платформа данных
    • Искусственный интеллект
    • Безопасность
    • Инструменты DevOps
    • Бессерверные вычисления
    • Управление ресурсами
  • Все решения
    • По отраслям
    • По типу задач
    • Экономика платформы
    • Безопасность
    • Техническая поддержка
    • Каталог партнёров
    • Обучение и сертификация
    • Облако для стартапов
    • Облако для крупного бизнеса
    • Центр технологий для общества
    • Партнёрская программа
    • Поддержка IT-бизнеса
    • Облако для фрилансеров
    • Обучение и сертификация
    • Блог
    • Документация
    • Мероприятия и вебинары
    • Контакты, чаты и сообщества
    • Идеи
    • Калькулятор цен
    • Тарифы
    • Акции и free tier
  • Кейсы
  • Документация
  • Блог
Создавайте контент и получайте гранты!Готовы написать своё руководство? Участвуйте в контент-программе и получайте гранты на работу с облачными сервисами!
Подробнее о программе
Проект Яндекса
© 2026 ТОО «Облачные Сервисы Казахстан»
Yandex Managed Service for Apache Kafka®
KZ
  • Начало работы
    • Все руководства
    • Самостоятельное развертывание веб-интерфейса Apache Kafka®
    • Обновление кластера Managed Service for Apache Kafka® с ZooKeeper на KRaft
    • Работа с топиками Apache Kafka® с помощью Yandex Data Processing
    • Отслеживание потери сообщений в топике Apache Kafka®
    • Передача данных из Managed Service for Apache Kafka® в таблицу Apache Iceberg™ в Object Storage
  • Управление доступом
  • Правила тарификации
  • Справочник Terraform
  • Метрики Yandex Monitoring
  • Аудитные логи Audit Trails
  • Публичные материалы
  • История изменений
  • Вопросы и ответы
  • Обучающие курсы

В этой статье:

  • Перед началом работы
  • Необходимые платные ресурсы
  • Подготовьте инфраструктуру
  • Создайте таблицу Apache Iceberg™
  • Создайте коннектор и проверьте его работу
  • Удалите созданные ресурсы
  1. Практические руководства
  2. Передача данных из Managed Service for Apache Kafka® в таблицу Apache Iceberg™ в Object Storage

Передача данных из топика Yandex Managed Service for Apache Kafka® в таблицу Apache Iceberg™ в Yandex Object Storage

Статья создана
Yandex Cloud
Обновлена 30 сентября 2026 г.
Открыть в Markdown
  • Перед началом работы
    • Необходимые платные ресурсы
  • Подготовьте инфраструктуру
  • Создайте таблицу Apache Iceberg™
  • Создайте коннектор и проверьте его работу
  • Удалите созданные ресурсы

Вы можете настроить передачу данных из топика Managed Service for Apache Kafka® в таблицу Apache Iceberg™. Схема и таблица создаются через Yandex Managed Service for Trino. Данные хранятся в бакете Object Storage, а метаданные — в Apache Hive™ Metastore. В Managed Service for Apache Kafka® настраивается коннектор Iceberg Sink для передачи данных из топика в таблицу.

Таблицу Apache Iceberg™ также можно создавать через Yandex Managed Service for Apache Spark™ с помощью PySpark-задания. Подробнее в руководстве Работа с таблицей формата Apache Iceberg™ из PySpark-задания.

Чтобы настроить передачу данных:

  1. Подготовьте инфраструктуру.
  2. Создайте таблицу Apache Iceberg™.
  3. Создайте коннектор и проверьте его работу.

Если созданные ресурсы вам больше не нужны, удалите их.

Перед началом работыПеред началом работы

Зарегистрируйтесь в Yandex Cloud и создайте платежный аккаунт:

  1. Перейдите в консоль управления, затем войдите в Yandex Cloud или зарегистрируйтесь.
  2. На странице Yandex Cloud Billing убедитесь, что у вас подключен платежный аккаунт, и он находится в статусе ACTIVE или TRIAL_ACTIVE. Если платежного аккаунта нет, создайте его и привяжите к нему облако.

Если у вас есть активный платежный аккаунт, вы можете создать или выбрать каталог, в котором будет работать ваша инфраструктура, на странице облака.

Подробнее об облаках и каталогах.

Необходимые платные ресурсыНеобходимые платные ресурсы

  • Кластер Managed Service for Apache Kafka®: использование выделенных хостам вычислительных ресурсов и объем хранилища (тарифы Managed Service for Apache Kafka®).
  • Публичные IP-адреса, если для хостов кластера включен публичный доступ (тарифы Yandex Virtual Private Cloud).
  • Кластер Apache Hive™ Metastore: вычислительные ресурсы компонентов кластера (тарифы Yandex MetaData Hub).
  • Кластер Managed Service for Trino: вычислительные ресурсы компонентов кластера и объем исходящего трафика из Yandex Cloud в интернет (тарифы Managed Service for Trino).
  • Бакет Object Storage: использование хранилища и выполнение операций с данными (тарифы Object Storage).
  • NAT-шлюз: почасовое использование шлюза и исходящий через него трафик (тарифы Virtual Private Cloud).

Подготовьте инфраструктуруПодготовьте инфраструктуру

  1. Создайте сервисный аккаунт sa-metastore и назначьте ему роли:

    • storage.editor — для работы с бакетом Object Storage.
    • managed-metastore.integrationProvider — для взаимодействия кластера Apache Hive™ Metastore с сервисами Yandex Cloud.
  2. Создайте сервисный аккаунт sa-trino и назначьте ему роли:

    • storage.editor — для работы с бакетом Object Storage.
    • managed-trino.integrationProvider — для взаимодействия кластера Managed Service for Trino с сервисами Yandex Cloud.
  3. Создайте бакет Object Storage.

  4. Создайте статический ключ доступа для сервисного аккаунта sa-metastore.

    Сохраните идентификатор ключа и секретный ключ, они понадобятся при создании коннектора.

  5. Создайте облачную сеть с именем demo-network.

    Вместе с ней будут автоматически созданы три подсети в разных зонах доступности.

  6. Настройте NAT-шлюз для подсети demo-network-kz1-a.

    NAT-шлюз нужен для взаимодействия кластера Apache Hive™ Metastore с сервисами Yandex Cloud.

  7. В сети demo-network создайте группу безопасности metastore-sg для кластера Apache Hive™ Metastore и добавьте в группу правила, необходимые для работы кластера.

  8. В сети demo-network создайте группу безопасности mkf-sg для кластера Managed Service for Apache Kafka® и добавьте в нее следующие правила:

    • Правило для входящего трафика, которое разрешает подключения к кластеру через интернет:

      • Диапазон портов — 9091.
      • Протокол — TCP.
      • Источник — Диапазон адресов.
      • IPv4 CIDR — 0.0.0.0/0.
    • Правило для исходящего трафика, которое разрешает доступ к Apache Hive™ Metastore:

      • Диапазон портов — 9083.
      • Протокол — TCP.
      • Назначение — Диапазон адресов.
      • IPv4 CIDR — 0.0.0.0/0.
  9. Создайте кластер Apache Hive™ Metastore со следующими настройками:

    • Сервисный аккаунт — sa-metastore.
    • Имя бакета — имя созданного ранее бакета.
    • Сеть — demo-network.
    • Подсеть — demo-network-kz1-a.
    • Группы безопасности — metastore-sg.
  10. Создайте кластер Managed Service for Trino со следующими настройками:

    • Сеть — demo-network.

    • Сервисный аккаунт — sa-trino.

    • Параметры каталога:

      • Имя каталога — iceberg.

      • Тип коннектора — Iceberg.

      • Тип Metastore — Hive Metastore.

      • URI — thrift://<IP-адрес_кластера_Metastore>:9083.

        IP-адрес кластера Apache Hive™ Metastore можно получить с информацией о кластере.

      • Файловое хранилище — Yandex Object Storage.

  11. Создайте кластер Managed Service for Apache Kafka® со следующими настройками:

    • Сеть — demo-network.

    • Группы безопасности — mkf-sg.

    • Публичный доступ — включен.

      Примечание

      Публичный доступ к хостам кластера нужен, если вы планируете подключаться к кластеру через интернет. Этот вариант подключения более простой, и его рекомендуется использовать для прохождения руководства. К хостам без публичного доступа тоже можно подключиться, но только с виртуальных машин Yandex Cloud, расположенных в той же облачной сети, что и кластер.

  12. В кластере Managed Service for Apache Kafka® создайте топики:

    • iceberg_control_topic — для управления коннектором;
    • my_topic — для обмена сообщениями.
  13. В кластере Managed Service for Apache Kafka® создайте пользователя kafka-producer с ролью производителя и предоставьте ему доступ к топику my_topic.

    Этот пользователь используется для отправки сообщений в my_topic.

Создайте таблицу Apache Iceberg™Создайте таблицу Apache Iceberg™

  1. Подключитесь к кластеру Managed Service for Trino.

  2. Выполните SQL-запросы:

    1. Создайте схему:

      CREATE SCHEMA iceberg.myschema
      WITH (
        location = 's3a://<имя_бакета>/iceberg/warehouse/myschema'
      );
      
    2. Создайте таблицу:

      CREATE TABLE iceberg.myschema.mytable (
        id BIGINT,
        name VARCHAR,
        created_at VARCHAR
      )
      WITH (format = 'PARQUET');
      
  3. Проверьте, что в бакете создана структура каталогов iceberg/warehouse/myschema/mytable-*/metadata/, где * — системная часть имени таблицы.

Создайте коннектор и проверьте его работуСоздайте коннектор и проверьте его работу

  1. В Managed Service for Apache Kafka® создайте коннектор типа Iceberg Sink со следующими параметрами:

    • Топик управления — iceberg_control_topic.

    • Топики — my_topic.

    • Таблицы — myschema.mytable.

    • URI каталога — thrift://<IP-адрес_кластера_Metastore>:9083.

      IP-адрес кластера Apache Hive™ Metastore можно получить с информацией о кластере.

    • Warehouse — s3a://<имя_бакета>/iceberg/warehouse.

    • Эндпоинт — storage.yandexcloud.net.

    • Идентификатор ключа доступа, Секретный ключ — полученные ранее данные о статическом ключе.

    • Интервал коммита, мс — 5000 миллисекунд.

    • Дополнительные свойства:

      • key.converter: org.apache.kafka.connect.json.JsonConverter
      • key.converter.schemas.enable: false
      • value.converter: org.apache.kafka.connect.json.JsonConverter
      • value.converter.schemas.enable: false
  2. Проверьте работу коннектора:

    1. Установите SSL-сертификат.

    2. Установите утилиту kcat (kafkacat).

    3. Отправьте сообщение в my_topic:

      echo '{"id":1,"name":"Alice","created_at":"2024-01-15T10:30:00"}' | kcat -P \
          -b <FQDN_брокера>:9091 \
          -t my_topic \
          -X security.protocol=SASL_SSL \
          -X sasl.mechanism=SCRAM-SHA-512 \
          -X sasl.username="kafka-producer" \
          -X sasl.password="<пароль>" \
          -X ssl.ca.location=/usr/local/share/ca-certificates/Yandex/YandexInternalRootCA.crt
      

      Подробнее о получении FQDN хоста-брокера читайте в разделе Получение FQDN хостов Apache Kafka®.

  3. После коммита проверьте, что в каталоге бакета iceberg/warehouse/myschema/mytable-*/ создан каталог data с данными в формате PARQUET.

    Чтобы посмотреть записанные данные, выполните SQL-запрос в Managed Service for Trino:

    SELECT * FROM iceberg.myschema.mytable;
    

Удалите созданные ресурсыУдалите созданные ресурсы

Некоторые ресурсы платные. Чтобы за них не списывалась плата, удалите ресурсы, которые вы больше не будете использовать:

  1. Кластер Managed Service for Apache Kafka®.
  2. Кластер Apache Hive™ Metastore.
  3. Кластер Managed Service for Trino.
  4. Бакет Object Storage. Перед удалением бакета удалите из него все объекты.
  5. NAT-шлюз.

Была ли статья полезна?

Предыдущая
Отслеживание потери сообщений в топике Apache Kafka®
Следующая
Взаимосвязь ресурсов сервиса
Создавайте контент и получайте гранты!Готовы написать своё руководство? Участвуйте в контент-программе и получайте гранты на работу с облачными сервисами!
Подробнее о программе
Проект Яндекса
© 2026 ТОО «Облачные Сервисы Казахстан»