Передача данных из топика Yandex Managed Service for Apache Kafka® в таблицу Apache Iceberg™ с использованием Apache Hive™ Metastore
Вы можете настроить передачу данных из топика 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-задания.
Чтобы настроить передачу данных:
- Подготовьте инфраструктуру.
- Создайте таблицу Apache Iceberg™.
- Создайте коннектор и проверьте его работу.
Если созданные ресурсы вам больше не нужны, удалите их.
Перед началом работы
Зарегистрируйтесь в Yandex Cloud и создайте платежный аккаунт:
- Перейдите в консоль управления
, затем войдите в Yandex Cloud или зарегистрируйтесь. - На странице 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).
Подготовьте инфраструктуру
-
Создайте сервисный аккаунт
sa-metastoreи назначьте ему роли:- storage.editor — для работы с бакетом Object Storage.
- managed-metastore.integrationProvider — для взаимодействия кластера Apache Hive™ Metastore с сервисами Yandex Cloud.
-
Создайте сервисный аккаунт
sa-trinoи назначьте ему роли:- storage.editor — для работы с бакетом Object Storage.
- managed-trino.integrationProvider — для взаимодействия кластера Managed Service for Trino с сервисами Yandex Cloud.
-
Создайте бакет Object Storage.
-
Создайте статический ключ доступа для сервисного аккаунта
sa-metastore.Сохраните идентификатор ключа и секретный ключ, они понадобятся при создании коннектора.
-
Создайте облачную сеть с именем
demo-network.Вместе с ней будут автоматически созданы три подсети в разных зонах доступности.
-
Настройте NAT-шлюз для подсети
demo-network-ru-central1-a.NAT-шлюз нужен для взаимодействия кластера Apache Hive™ Metastore с сервисами Yandex Cloud.
-
В сети
demo-networkсоздайте группу безопасностиmetastore-sgдля кластера Apache Hive™ Metastore и добавьте в группу правила, необходимые для работы кластера. -
В сети
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.
- Диапазон портов —
-
-
Создайте кластер Apache Hive™ Metastore со следующими настройками:
- Сервисный аккаунт —
sa-metastore. - Имя бакета — имя созданного ранее бакета.
- Сеть —
demo-network. - Подсеть —
demo-network-ru-central1-a. - Группы безопасности —
metastore-sg.
- Сервисный аккаунт —
-
Создайте кластер 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.
-
-
-
Создайте кластер Managed Service for Apache Kafka® со следующими настройками:
-
Сеть —
demo-network. -
Группы безопасности —
mkf-sg. -
Публичный доступ — включен.
Примечание
Публичный доступ к хостам кластера нужен, если вы планируете подключаться к кластеру через интернет. Этот вариант подключения более простой, и его рекомендуется использовать для прохождения руководства. К хостам без публичного доступа тоже можно подключиться, но только с виртуальных машин Yandex Cloud, расположенных в той же облачной сети, что и кластер.
-
-
В кластере Managed Service for Apache Kafka® создайте топики:
iceberg_control_topic— для управления коннектором;my_topic— для обмена сообщениями.
-
В кластере Managed Service for Apache Kafka® создайте пользователя
kafka-producerс ролью производителя и предоставьте ему доступ к топикуmy_topic.Этот пользователь используется для отправки сообщений в
my_topic.
Создайте таблицу Apache Iceberg™
-
Подключитесь к кластеру Managed Service for Trino.
-
Выполните SQL-запросы:
-
Создайте схему:
CREATE SCHEMA iceberg.myschema WITH ( location = 's3a://<имя_бакета>/iceberg/warehouse/myschema' ); -
Создайте таблицу:
CREATE TABLE iceberg.myschema.mytable ( id BIGINT, name VARCHAR, created_at VARCHAR ) WITH (format = 'PARQUET');
-
-
Проверьте, что в бакете создана структура каталогов
iceberg/warehouse/myschema/mytable-*/metadata/, где*— системная часть имени таблицы.
Создайте коннектор и проверьте его работу
-
В 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.JsonConverterkey.converter.schemas.enable:falsevalue.converter:org.apache.kafka.connect.json.JsonConvertervalue.converter.schemas.enable:false
-
-
Проверьте работу коннектора:
-
Отправьте сообщение в
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®.
-
После коммита проверьте, что в каталоге бакета
iceberg/warehouse/myschema/mytable-*/создан каталогdataс данными в форматеPARQUET.Чтобы посмотреть записанные данные, выполните SQL-запрос в Managed Service for Trino:
SELECT * FROM iceberg.myschema.mytable;
Удалите созданные ресурсы
Некоторые ресурсы платные. Чтобы за них не списывалась плата, удалите ресурсы, которые вы больше не будете использовать: