Data Lake с помощью Debezium, Kafka Connect и Apache Iceberg sink
Что такое Apache Iceberg?
Apache Iceberg - это открытый формат таблиц, предназначенный для огромных аналитических наборов данных, который можно использовать с многими мощнейшими движками обработки Big Data, такими как Apache Spark, Trino, PrestoDB, Flink и Hive. Эта технология может использоваться не только в пакетной обработке, но и быть отличным инструментом для сбора данных в режиме реального времени, которые поступают из действий пользователей, метрик, журналов, из CDC или других источников. Apache Iceberg предоставляет механизмы для изоляции чтения-записи и уплотнения данных из коробки, что позволяет избежать проблем при работе с маленькими файлами.
Также стоит отметить и то, что Apache Iceberg можно использовать с любым облачным провайдером или собственным решением, поддерживающим метахранилище Apache Hive и blob-хранилище.
Kafka Connect , Apache Iceberg sink
Компания GetInData создала Apache Iceberg sink, который можно развернуть на экземпляре Kafka Connect. Вы можете найти репозиторий и пакет на нашем GitHub.
Apache Iceberg sink был создан на базе memiiso/debezium-server-iceberg, который в свою очередь был создан для автономного использования наряду с Debezium Server.
Формат данных, используемый Apache Iceberg, должен представлять табличные данные и их схему, поэтому для сбора данных об изменениях мы использовали формат, созданный Debezium.
Пример CDC
Давайте попробуем использовать Apache Kafka sink для репликации базы данных PostgreSQL с помощью Debezium для захвата всех изменений и передачи их в таблицу Apache Iceberg.
Мы запустим экземпляр Kafka Connect, на котором развернем источник Debezium, а также наш Apache Iceberg sink. Для связи между ними будет использоваться топик Kafka, а sink будет записывать данные в бакет S3 и метаданные в Amazon Glue. Позже для чтения и отображения данных мы будем использовать Amazon Athena.
Шаг 1: запуск Kafka Connect
Сначала пройдите аутентификацию и сохраните учетные данные AWS в файле, например ~/.aws/confi
[default] region = eu-west-1 aws_access_key_id=\*\** aws_secret_access_key=\*\**
Загрузите sink отсюда. Например,~/Downloads/kafka-connect-iceberg-sink-0.1.3-shaded.jar
Для подключения Kafka мы будем использовать образ докера от Debezium, который поставляется с исходными пакетами Debezium. Мы смонтируем наш Apache Iceberg sink в каталог плагинов Kafka Connect и добавим файл с учетными данными AWS.
docker run -it --name connect --net=host -p 8083:8083 \ -e GROUP_ID=1 \ -e CONFIG_STORAGE_TOPIC=my-connect-configs \ -e OFFSET_STORAGE_TOPIC=my-connect-offsets \ -e BOOTSTRAP_SERVERS=localhost:9092 \ -e CONNECT_TOPIC_CREATION_ENABLE=true \ -v ~/.aws/config:/kafka/.aws/config \ -v ~/Downloads/kafka-connect-iceberg-sink-0.1.3-shaded.jar:/kafka/connect/kafka-connect-iceberg-sink-0.1.3-shaded.jar \ debezium/connect
Шаг 2: Чтение данных из источника PostgreSQL
Одна из возможностей Debezium по считыванию данных из PostgreSQL - это работа в качестве реплики базы данных. Чтобы Debezium работал корректно, нам необходимо увеличить объем информации, хранящейся в журнале опережающей записи. Для этого нам нужно настроить уровень wal_level на логический.
Запустим PostgreSQL на Docker:
docker run -d --name postgres -e POSTGRES_PASSWORD=postgres \ -p 5432:5432 postgres -c wal_level=logical
Нам также понадобится экземпляр Kafka.
Если функция автоматического создания топиков не включена, нужно создать топик, который будет использоваться для связи между источником Debezium и нашим Apache Iceberg sink. Имя топика состоит из логического имени, которое мы присвоим источнику Debezium, имени схемы базы данных и имени таблицы.
kafka-topics.sh --bootstrap-server localhost:9092 --create --topic postgres.public.dbz_test --partitions 1 --replication-factor 1
Теперь нам нужно развернуть источник на Kafka Connect. Мы можем сделать это с помощью POST-запроса, содержащего его конфигурацию.
curl -X POST -H "Content-Type: application/json" \
-d '{
"name": "postgres-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "localhost",
"database.port": "5432",
"database.user": "postgres",
"topic.prefix": "postgres",
"database.password": "postgres",
"database.dbname" : "postgres",
"database.server.name": "postgres",
"slot.name": "debezium",
"plugin.name": "pgoutput",
"table.include.list": "public.dbz_test",
"transforms" : "unwrap",
"transforms.unwrap.type":"io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.add.fields":"op,table,lsn,source.ts_ms,db",
"transforms.unwrap.drop.tombstones":"true",
"transforms.unwrap.delete.handling.mode":"rewrite",
"drop.tombstones": "true"
}
}' \
<http://localhost:8083/connectors>
Шаг 3: Apache Iceberg Sink
Для нашего Apache Iceberg sink нам понадобится бакет S3, например gid-streaminglabs-eu-west-1, и база данных в Amazon Glue, например gid_streaminglabs_eu_west_1_dbz.
Поскольку у нас уже есть готовый экземпляр Kafka Connect, включая учетные данные AWS, и пакет с sink, осталось только развернуть его. По аналогии с источником PostgreSQL мы сделаем это с помощью POST-запроса.
curl -X POST -H "Content-Type: application/json" \
-d '{
"name": "iceberg-sink",
"config": {
"connector.class": "com.getindata.kafka.connect.iceberg.sink.IcebergSink",
"topics": "postgres.public.dbz_test",
"upsert": true,
"upsert.keep-deletes": true,
"table.auto-create": true,
"table.write-format": "parquet",
"table.namespace": "gid_streaminglabs_eu_west_1_dbz",
"table.prefix": "debeziumcdc_",
"iceberg.catalog-impl": "org.apache.iceberg.aws.glue.GlueCatalog",
"iceberg.warehouse": "s3a://gid-streaminglabs-eu-west-1/dbz_iceberg/gl_test",
"iceberg.fs.defaultFS": "s3a://gid-streaminglabs-eu-west-1/dbz_iceberg/gl_test",
"iceberg.com.amazonaws.services.s3.enableV4": true,
"iceberg.com.amazonaws.services.s3a.enableV4": true,
"iceberg.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
"iceberg.fs.s3a.path.style.access": true,
"iceberg.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem"
}
}' \
<http://localhost:8083/connectors>
Шаг 4: Проверка
Теперь мы можем открыть клиент psql и создать несколько таблиц:
psql -U postgres -h localhost
create table dbz_test (timestamp bigint, id int PRIMARY KEY, value int); insert into dbz_test values(1, 1, 1); insert into dbz_test values(2, 2, 2); alter table dbz_test add test varchar(30); insert into dbz_test values(3, 3, 3, 'aaa'); delete from dbz_test where id = 1; update dbz_test set value = 1 where id = 2;
Затем перейдите в Amazon Athena и выполните следующий запрос:
select * from debeziumcdc_postgres_public_dbz_test order by timestamp desc;
Вот, собственно говоря, и все.





