Apache Iceberg как фундамент современной Data-инфраструктуры. Подробное руководство по внедрению, лучшим практикам и управлению рисками
В современном мире данных, где объемы информации растут в геометрической прогрессии, а требования к скорости, надежности и гибкости анализа становятся все более строгими, критически важным становится выбор правильной технологической основы. Устаревшие подходы к управлению данными, такие как Hive Metastore (HMS), все чаще показывают свою несостоятельность перед лицом новых вызовов. Они создают барьеры для развития, связывая бизнес жесткими рамками, сложностями масштабирования и высокими операционными затратами.
Apache Iceberg не просто еще один мощный инструмент, а стратегическая технология, кардинально меняющая правила игры. Iceberg — это мощный фундамент, большая часть возможностей которого скрыта от первого взгляда. Это открытый, высокопроизводительный формат табличных метаданных для огромных аналитических наборов данных. В данном подробном обзоре мы рассмотрим, зачем вашей организации нужен Apache Iceberg, как его правильно внедрить, с какими рисками вы можете столкнуться и как наши эксперты помогают клиентам извлекать из него максимальную пользу.
Почему Apache Iceberg — это эволюция в управлении данными?
Apache Iceberg был разработан как ответ на ключевые недостатки устаревших систем, в частности Hive Metastore. Основные проблемы HMS, которые тормозят развитие data-продуктов, включают жесткую привязку к файловой системе HDFS, крайне сложные и дорогостоящие процессы изменения схемы данных (schema evolution), отсутствие гарантий согласованности данных при параллельной записи и отсутствие механизмов для эффективного управления версиями и аудита.
Iceberg предлагает принципиально иной подход, предоставляя ряд революционных преимуществ:
- Мгновенное изменение схемы данных «на лету»: добавление, удаление или переименование столбцов не требует перезаписи самих данных или длительных простоев. Это ускоряет разработку и позволяет бизнесу быстро адаптироваться к меняющимся требованиям.
- Полноценное версионирование данных и «путешествие во времени»: каждое изменение данных создает новый снимок (snapshot). Вы можете легко запросить данные в том виде, в каком они существовали на любой момент времени в прошлом, что незаменимо для аудита, отладки пайплайнов и воспроизведения результатов отчетности.
- Управление ветками для данных: эта функция, аналогичная Git, позволяет создавать изолированные среды для тестирования и экспериментов с данными, не затрагивая продовольственные данные. Вы можете безопасно применять сложные изменения, а затем «мержить» их в основную ветку.
- Независимость от вычислительного движка: Iceberg не навязывает свой вычислительный движок. Вы можете использовать для запросов Spark, Trino, Flink, Presto, Dremio, Snowflake и многие другие инструменты, выбирая оптимальный для каждой задачи.
- Эффективное управление данными: Встроенные механизмы компaction файлов, экспирации снимков и удаления метаданных помогают оптимизировать затраты на хранение и поддерживать высокую производительность запросов.
Теоретические преимущества бессмысленны без надежной практической реализации. Наши специалисты выделили ключевые этапы и подводные камни развертывания Apache Iceberg.
Поднятие инфраструктуры: надежность с первого шага
Для быстрого старта и тестирования мы используем docker-compose, как это рекомендовано в официальном Quickstart от Apache Iceberg. Однако наша команда всегда модифицирует стандартные конфигурации для повышения отказоустойчивости и управляемости в реальных условиях (подробности – ниже).
Весь код доступен в репозитории.
docker-compose.yaml
version: '3.8'
services:
postgres:
image: postgres:13
environment:
POSTGRES_USER: postgres
POSTGRES_PASSWORD: postgres
POSTGRES_DB: iceberg
ports:
- "5432:5432"
networks:
- iceberg_net
rest:
image: tabulario/iceberg-rest:1.6.0
container_name: iceberg-rest
networks:
iceberg_net:
ports:
- "8181:8181"
environment:
- AWS_ACCESS_KEY_ID=minioadmin
- AWS_SECRET_ACCESS_KEY=minioadmin
- AWS_REGION=us-east-1
- CATALOG_WAREHOUSE=s3://warehouse/
- CATALOG_IO__IMPL=org.apache.iceberg.aws.s3.S3FileIO
- CATALOG_S3_ENDPOINT=http://minio:9000
- CATALOG_URI=jdbc:postgresql://postgres:5432/iceberg?user=postgres&password=postgres
- CATALOG_JDBC_DRIVER=org.postgresql.Driver
- CATALOG_JDBC_USER=postgres
- CATALOG_JDBC_PASSWORD=postgres
depends_on:
- postgres
minio:
image: minio/minio:RELEASE.2024-07-04T14-25-45Z
restart: always
command: server /data --console-address ":9001"
volumes:
- ./data:/data
environment:
- MINIO_ROOT_USER=minioadmin
- MINIO_ROOT_PASSWORD=minioadmin
- MINIO_DOMAIN=minio
ports:
- "9000:9000" # MinIO S3 API
- "9001:9001" # MinIO Console
networks:
iceberg_net:
aliases:
- warehouse.minio
mc:
depends_on:
- minio
image: minio/minio:RELEASE.2024-07-04T14-25-45Z
container_name: mc
networks:
iceberg_net:
environment:
- AWS_ACCESS_KEY_ID=minioadmin
- AWS_SECRET_ACCESS_KEY=minioadmin
- AWS_REGION=us-east-1
entrypoint: >
/bin/sh -c "
until (/usr/bin/mc config host add minio http://minio:9000 minioadmin minioadmin) do echo '...waiting...' && sleep 1; done;
/usr/bin/mc mb minio/warehouse;
/usr/bin/mc policy set public minio/warehouse;
tail -f /dev/null
"
spark-iceberg:
image: tabulario/spark-iceberg:3.5.1_1.5.0
container_name: spark-iceberg
build: spark/
networks:
iceberg_net:
depends_on:
- rest
- minio
volumes:
- ./warehouse:/home/iceberg/warehouse
- ./notebooks:/home/iceberg/notebooks/notebooks
environment:
- AWS_ACCESS_KEY_ID=minioadmin
- AWS_SECRET_ACCESS_KEY=minioadmin
- AWS_REGION=us-east-1
ports:
- "8888:8888"
- "8080:8080"
- "10000:10000"
- "10001:10001"
networks:
iceberg_net:
Стандартная конфигурация использует in-memory метастор на базе SQLite, что недопустимо для продакшн-среды, так как все метаданные теряются при перезапуске. Наши инженеры немедленно заменяют его на промышленную СУБД, такую как PostgreSQL. Это гарантирует сохранность и целостность метаданных — сердцевины вашего Data Lake.
... rest: image: tabulario/iceberg-rest container_name: iceberg-rest networks: iceberg_net: ports: - 8181:8181 environment: - AWS_ACCESS_KEY_ID=admin - AWS_SECRET_ACCESS_KEY=password - AWS_REGION=us-east-1 - CATALOG_WAREHOUSE=s3://warehouse/ - CATALOG_IO__IMPL=org.apache.iceberg.aws.s3.S3FileIO - CATALOG_S3_ENDPOINT=http://minio:9000 ...
Если создать с environment, то ваш meta-store будет создаваться в формате memory.
Данная команда будет выполнена следующим образом:
CATALOG_URI=jdbc:sqlite:file:/tmp/iceberg_rest_mode=memory
И ваш meta-store будет доступен внутри контейнера в папке tmp/
iceberg@e2d60003b684:/tmp$ ls -lahtotal 1.1Mdrwxrwxrwt 1 root root 4.0K Oct 7 10:58 .drwxr-xr-x 1 root root 4.0K Oct 7 06:43 ..drwxr-xr-x 2 iceberg iceberg 4.0K Oct 7 06:43 hsperfdata_iceberg-rw-r--r-- 1 iceberg iceberg 20K Oct 7 10:58 'iceberg_rest_mode=memory'-rwxr--r-- 1 iceberg iceberg 1.1M Oct 7 06:43 sqlite-3.46.0.0-be26ebff-c4de-43c2-bd44-71c205b8c5bd-libsqlitejdbc.so-rw-r--r-- 1 iceberg iceberg 0 Oct 7 06:43 sqlite-3.46.0.0-be26ebff-c4de-43c2-bd44-71c205b8c5bd-libsqlitejdbc.so.lck
Я хотел контролировать meta-store и беспрепятственно его просматривать, поэтому изменил сборку образа на такую:
... postgres: image: postgres:13 environment: POSTGRES_USER: postgres POSTGRES_PASSWORD: postgres POSTGRES_DB: iceberg ports: - "5432:5432" networks: - iceberg_net rest: image: tabulario/iceberg-rest:1.6.0 container_name: iceberg-rest networks: iceberg_net: ports: - "8181:8181" environment: - AWS_ACCESS_KEY_ID=minioadmin - AWS_SECRET_ACCESS_KEY=minioadmin - AWS_REGION=us-east-1 - CATALOG_WAREHOUSE=s3://warehouse/ - CATALOG_IO__IMPL=org.apache.iceberg.aws.s3.S3FileIO - CATALOG_S3_ENDPOINT=http://minio:9000 - CATALOG_URI=jdbc:postgresql://postgres:5432/iceberg?user=postgres&password=postgres - CATALOG_JDBC_DRIVER=org.postgresql.Driver - CATALOG_JDBC_USER=postgres - CATALOG_JDBC_PASSWORD=postgres depends_on: - postgres ...
Пример нашей доработанной конфигурации для сервиса Iceberg REST включает явное указание JDBC-подключения к PostgreSQL. Мы настоятельно рекомендуем нашим клиентам никогда не использовать стандартную in-memory конфигурацию для каких-либо задач, кроме моментального ознакомления.
Взаимодействие с данными: выбор правильного инструмента
Одним из ключевых преимуществ Iceberg является поддержка множества инструментов, но каждый из них имеет свои нюансы и степень зрелости. Наши эксперты провели глубокий анализ и готовы поделиться выводами.
Iceberg предоставляет полноценный REST API для всех операций с каталогом. Это открывает возможности для автоматизации, интеграции с внутренними системами и созданию custom-инструментов управления.
Спецификация API доступна по ссылке – Iceberg REST Open API specification.
Список namespaces:
curl http://localhost:8181/v1/namespaces
Список таблиц:
curl http://localhost:8181/v1/namespaces/default/tables
Теперь поговорим о DuckDB. На текущий момент (октябрь 2024) расширение DuckDB для Iceberg находится на ранней стадии развития. Оно поддерживает только базовые операции SELECT и содержит множество известных ошибок. Наша рекомендация — использовать его с крайней осторожностью только для простейших задач чтения. Промышленная эксплуатация нецелесообразна. Мы активно следим за развитием проекта и готовы предоставить обновления нашим клиентам по мере его стабилизации.
В начале предлагаем сконфигурировать сессию:
INSTALL iceberg;LOAD iceberg;INSTALL httpfs;LOAD httpfs;SETs3_url_style= 'path';SETs3_endpoint= 'localhost:9000';SETs3_access_key_id= 'minioadmin';SETs3_secret_access_key= 'minioadmin';SETs3_use_ssl= FALSE;
Теперь, если выполнить запрос как указано в документации:
SELECT * FROMiceberg_metadata('s3://warehouse/default/animals/',allow_moved_paths= TRUE);
Мы получим ошибку:
SQL Error: java.sql.SQLException: HTTP Error: HTTP GET error on 'http://localhost:9000/warehouse/default/animals//metadata/version-hint.text' (HTTP 400)
Решение данной проблемы можно найти здесь. Но есть одно НО!
Можно читать определённый кусок мета-данных:
SELECT * FROMiceberg_scan('s3://warehouse/default/animals/metadata/00001-f0d7c171-b179-4857-8eec-078f0108c1a9.metadata.json');
Но если выполнить INSERT/UPDATE/DELETE, то прошлые мета-даннные будут ссылаться на прошлый снэпшот. Чтобы получить "актуальные" данные, необходимо найти новые мета-данные:
SELECT * FROMiceberg_scan('s3://warehouse/default/animals/metadata/00002-ff3b9eaf-3952-4494-a2c0-b29a78cc6bb7.metadata.json');
Примечание: существует "маска" для создания новых снепшотов: 00001-*, 00002-*, etc
И теперь НО, про которое было сказано выше. Если читать первый файл метаданных, то получим ошибку:
SQL Error: java.sql.SQLException: IO Error: No snapshots found
Чтобы корректно работать с Iceberg через DuckDB необходимо ждать исправления багов.
Как и DuckDB, ClickHouse поддерживает работу с Iceberg только на чтение через табличные функции. Это мощное решение для выполнения высокоскоростных аналитических запросов к данным в Iceberg, но не для их управления.
Теперь перейдем к PyIceberg - это основной Python-фреймворк для работы с Iceberg. Он зрелый, полнофункциональный и идеально подходит для построения пайплайнов, скриптов управления и автоматизации. Наши разработчики активно используют его в проектах. Мы настраиваем централизованный конфигурационный файл .pyiceberg.yaml для всех инструментов, что устраняет дублирование кода и повышает безопасность управления учетными данными.
Рассмотрим несколько практических примеров, которые иллюстрируют мощь PyIceberg.
В начале рекомендуем настроить доступ к каталогу.
Создание файла в .pyiceberg.yaml в корне учетной записи:
catalog: s3_warehouse: uri: http://127.0.0.1:8181 py-io-impl: pyiceberg.io.pyarrow.PyArrowFileIO s3.endpoint: http://127.0.0.1:9000 s3.access-key-id: minioadmin s3.secret-access-key: minioadmin
Вызываем код:
from pyiceberg.catalog import load_catalog catalog = load_catalog("s3_warehouse")catalog.create_namespace("default")
Если не создать .pyiceberg.yaml, то необходимо будет каждый раз инициализировать конфиг для подключения к каталогу таким образом:
from pyiceberg.catalog import load_catalog catalog = load_catalog( name="warehouse", **{ "uri": "http://127.0.0.1:8181", "s3.endpoint": "http://127.0.0.1:9000", "py-io-impl": "pyiceberg.io.pyarrow.PyArrowFileIO", "s3.access-key-id": "minioadmin", "s3.secret-access-key": "minioadmin", }, ) catalog.create_namespace("default")На наш взгляд, лучше создать .pyiceberg.yaml, чем дублировать весь код.
Для дальнейшей работы необходимо создать namespace . Для этого выполним код:
from pyiceberg.catalog import load_catalog catalog = load_catalog("s3_warehouse") catalog.create_namespace("default")Далее переходим к созданию таблицы. Оно может осуществляться двумя способами: с использованием нативных типов PyIceberg или через более привычный PyArrow Schema. Мы обычно рекомендуем подход с PyArrow, так как он интуитивно понятен для широкого круга дата-инженеров и легко интегрируется с существующими пайплайнами на Pandas и PySpark.
Для создания таблицы с использованием pyiceberg.types необходимо:
-
Импортировать нужные типы из пакета
pyiceberg.types. - Создать схему.
- Создать таблицу с созданной схемой. Пример:
from pyiceberg.catalog import load_catalog from pyiceberg.schema import Schema from pyiceberg.types import BinaryType, LongType, NestedField, StringType, TimestamptzType schema = Schema( fields=[ NestedField(field_id=1, name="id", field_type=LongType(), required=False), NestedField(field_id=2, name="uuid", field_type=BinaryType(), required=False, doc="Binary -> UUID"), NestedField(field_id=3, name="name", field_type=StringType(), required=False), NestedField(field_id=4, name="created_at", field_type=TimestamptzType(), required=False), ], ) catalog = load_catalog("s3_warehouse") # Comment if Table does not exist catalog.drop_table("default.custom_table_pyiceberg_fields") catalog.create_table( identifier="default.custom_table_pyiceberg_fields", schema=schema, )Для того чтобы создать таблицу с использованием pa.field необходимо:
-
Импортировать
pyarrow. - Создать схему.
- Создать таблицу. Пример:
import pyarrow as pa from pyiceberg.catalog import load_catalog schema = pa.schema( [ pa.field(name="id", type=pa.int64(), nullable=True), pa.field(name="uuid", type=pa.binary(), nullable=True), pa.field(name="name", type=pa.string(), nullable=True), pa.field(name="created_at", type=pa.timestamp(unit="s", tz="UTC"), nullable=True), ], ) catalog = load_catalog("s3_warehouse") # Comment if Table does not exist catalog.drop_table("default.custom_table_pyarrow_fields") catalog.create_table( identifier="default.custom_table_pyarrow_fields", schema=schema, )Следующий пункт - вставка данных в таблицу, которая всегда требует преобразования данных в объект pyarrow.Table. Наша команда разрабатывает универсальные утилиты-врапперы, которые автоматически конвертируют распространенные форматы (Pandas DataFrame, словари) в требуемый вид, что ускоряет разработку.
import datetime import uuid from random import randint import pyarrow as pa from faker import Faker from pyiceberg.catalog import load_catalog fake = Faker(locale="ru_RU") catalog = load_catalog("s3_warehouse") tbl = catalog.load_table("default.custom_table_pyiceberg_fields") pa_table = pa.table( { "id": [randint(a=1, b=9223372036854775806)], # noqa: S311 "uuid": [uuid.uuid4().bytes], "name": [fake.first_name()], "created_at": [fake.date_time_ad(tzinfo=datetime.UTC)], }, ) tbl.append(pa_table) print( tbl.scan().to_arrow(),Также, если вы привыкли работать с pandas.DataFrame, то его также необходимо трансформировать в pyarrow.lib.Table.
Пример:
import pandas as pd import pyarrow as pa from pyiceberg.catalog import load_catalog catalog = load_catalog("s3_warehouse") catalog.drop_table("default.yellow_taxi") df = pd.read_parquet("https://d37ci6vzurychx.cloudfront.net/trip-data/yellow_tripdata_2023-01.parquet") table = pa.Table.from_pandas(df=df) tbl = catalog.create_table( identifier="default.yellow_taxi", schema=table.schema, ) tbl.append(table)Для удаления данных из таблицы необходимо воспользоваться методом tbl.delete() предварительно инициализировав tbl:
rom pyiceberg.catalog import load_catalog catalog = load_catalog("s3_warehouse")tbl = catalog.load_table("default.animals")tbl.delete(delete_filter="animal == 'Bear'")
Для чтения необходимо вызвать метод tbl.scan() предварительно инициализировав catalog и tbl:
from pyiceberg.catalog import load_catalog catalog = load_catalog("s3_warehouse")tbl = catalog.load_table("default.animals",)print(tbl.scan().to_arrow(),)
Здесь стоит обратить внимание, что после вызова метода tbl.scan() мы получаем объект типа <class 'pyiceberg.table.DataScan'> и можем преобразовать в удобный формат:
-
.to_pandas() -
.to_arrow() - etc
Для удаления таблицы необходимо вызвать catalog.drop_table(), предварительно инициализировав catalog.
from pyiceberg.catalog import load_catalog catalog = load_catalog("s3_warehouse")catalog.drop_table("default.animals")
Примечание: При удалении таблицы происходит удаление из мета-стора, а файлы остаются на месте. Таблица в данном случае – это просто ссылка на объект в S3.
Важнейшая функция — Time Travel. После каждой модификации данных Iceberg создает снапшот.
Для того, чтобы получить текущий список snapshots для таблицы необходимо выполнить следующий код:
from pyiceberg.catalog import load_catalog import pandas as pd pd.set_option("display.max_columns",None)catalog = load_catalog("s3_warehouse")table = catalog.load_table("default.custom_table_pyarrow_fields")print(table.inspect.snapshots())
В выводе мы увидим все id , которые сможем считать следующим образом:
from pyiceberg.catalog import load_catalog catalog = load_catalog("s3_warehouse")table = catalog.load_table("default.custom_table_pyarrow_fields")print(table.scan(snapshot_id=3826389439422852561))
Наконец, мы дошли до Apache Spark.
Несмотря на обилие инструментов, Apache Spark остается самым мощным и зрелым движком для работы с Iceberg, поддерживающим полный цикл операций DDL (Data Definition Language) и DML (Data Manipulation Language).
Во время сборки контейнера уже произошла инициализация подключения к хранилищу, поэтому мы можем сразу же обращаться запросами, исключая инициализацию.
Для начала посмотрим на то, какие таблицы нам доступны в namespace default, созданном ранее. Для этого выполним скрипт:
%%sql
show tables in default;
Для создания namespace необходимо выполнить команду:
%%sql
CREATE DATABASE IF NOT EXISTS foo;
После создания namespace foo мы его сможем увидеть его в мета-сторе, но физически его пока нет, поэтому давайте создадим таблицу:
%%sql CREATE TABLE foo.bar(id bigint);
Также мы можем создать таблицу с партиционированием по выбранному полю:
CREATE TABLEfoo.bar_partition(idbigint,created_attimestamp,)USING icebergPARTITIONEDBY(days(created_at))
Примечание: Правильное партиционирование — залог высокой производительности. Spark позволяет создавать партиции по дням, часам, а также по конкретным значениям (например, по региону). Наша команда проводит аудит ваших данных и запросов, чтобы предложить оптимальную стратегию партиционирования. Важно отметить, что стратегию можно менять «на лету» без пересоздания таблицы.
Данный код создаст таблицу и при каждой вставке данных в неё он будет создавать партицию за нужный день:
.├── data│ ├── created_at=YYYY-MM-DD│ ├── created_at=YYYY-MM-DD│ ├── created_at=YYYY-MM-DD│ ├── ...│ ├── created_at=YYYY-MM-DD
Если вы некорректно создали партицию ,то можете изменить ее через удаление и создание новой партиции:
%%sql ALTER TABLE foo.bar_partition DROP PARTITION FIELD days(created_at)
Создание новой партиции:
%%sql ALTER TABLE foo.bar_partition ADD PARTITION FIELD hours(created_at)
После этого необходимо вызвать метод rewrite_data_files, чтобы перезаписать партиции:
%%sql CALL system.rewrite_data_files('foo.bar_partition')
Что касается вставки данных в таблицу, то поскольку мы работаем через Spark, то у нас есть два варианта вставки данных:
-
через
dataframe:
df = spark.read.parquet(".parquet")df.writeTo("foo.bar").append()
-
через привычный
INSERT:
%%sql INSERT INTO foo.bar(id)VALUES(1)
Для удаления данных из таблицы необходимо выполнить следующий код:
%%sql
DELETE FROM foo.bar WHERE id = 1
Чтобы прочитать данные из таблицы foo.bar, необходимо выполнить код:
%%sql
SELECT * FROM foo.bar
Для того чтобы прочитать данные для дальнейшей работы, необходимо выполнить следующий код:
df = spark.sql('SELECT * FROM foo.bar')df.show()
Удаление таблицы происходит стандартной командой:
%%sql DROP TABLE IF EXISTS foo.bar;
Как я говорил ранее, отличительной способность Iceberg является «путешествие времени». Поэтому каждый раз, когда мы изменяли таблицу foo.bar, при помощи INSERT, DELETE, UPDATE, у нас создавался snapshot, к которому мы можем откатиться.
В начале посмотрим список доступных нам snapshots:
%%sql SELECT * FROM foo.bar.snapshots ORDER BY committed_at DESC
Теперь мы знаем все snapshot_id для нашей таблицы и можем прочитать любое состояние таблицы по snapshot_id:
%%sql SELECT count(*)as c FROM foo.bar FOR VERSION AS OF 3331575308018494635
Также мы можем это сделать по полю committed_at:
%%sql SELECT count(*)as c FROM foo.bar FOR TIMESTAMP AS OF TIMESTAMP 'YYYY-MM-DD HH:MM:SS.000000'
Ну и конечно же мы можем откатиться на прошлое состояние:
CALL system.rollback_to_snapshot('foo.bar',3331575308018494635)
Еще одним важным свойством является Branching - это одна из самых передовых функций. Мы помогаем клиентам внедрить Git-like workflow для данных. Например, вы можете создать ветку experiment, провести в ней агрессивные тесты и преобразования, а затем безопасно «вмержить» проверенные изменения в основную ветку main с помощью операции cherrypick_snapshot. Это кардинально меняет процессы тестирования и разработки в дата-инженерии.
Рассмотрим это на примере. Мы ранее создали таблицу foo.bar, теперь если вставить в неё несколько значений при помощи следующего кода:
%%sql INSERT INTO foo.bar(id)VALUES(1)
Теперь нам необходимо установить возможность версионирования через ветки для таблицы, выполним следующий код:
%%sql ALTER TABLE foo.bar SET TBLPROPERTIES('write.wap.enabled'='true')
Теперь мы можем создавать ветки для нашей таблицы, давайте создадим ветку delete_one_value:
%%sql
ALTER TABLE foo.bar CREATE BRANCH delete_one_value
После этого нам нужно переключиться на неё:
spark.conf.set('spark.wap.branch','delete_one_value')
Теперь выполним задачу в нужной ветке:
%%sql
DELETE FROM foo.bar WHERE id = 1
Для проверки того, что мы всё сделали корректно, выполним следующий код:
%%sql SELECT count(*)AS cnt -- FROM foo.bar VERSION AS OF 'main' FROM foo.bar.branch_main
Здесь мы получим то количество строк, который у нас находятся в main.
При выполнении следующего кода мы сможем получить то количество строк, которое соответствует ветке delete_one_value.
%%sql SELECT count(*)AS cnt -- FROM foo.bar VERSION AS OF 'delete_one_value' FROM foo.bar.branch_delete_one_value
Теперь самое главное –публикация изменений в main.
Для начала нам нужно узнать какой snapshot_id использовать для публикации.
%%sql SELECT * FROM foo.bar.refs
Чтобы сделать публикацию, необходимо выполнить код:
%%sql CALL system.cherrypick_snapshot('foo.bar',5235792610289419748)
В конце нужно удалить ветку delete_one_value:
%%sql
ALTER TABLE foo.bar DROP BRANCH delete_one_value
Чтобы проверить, что всё удалено корректно, выполним данный код:
%%sql SELECT * FROM foo.bar.refs
Управление рисками и предотвращение ошибок
Внедрение любой новой технологии сопряжено с рисками. На основе нашего опыта мы выделили ключевые зоны внимания:
- Потеря метаданных. Самая критичная точка отказа — это каталог (метастор). Использование ненадежной БД (например, SQLite) или отсутствие режима репликации ведет к катастрофическим последствиям. Для того, чтобы избежать чего-то подобного, мы развертываем метастор только на отказоустойчивых кластерах PostgreSQL с настроенным бэкапом и репликацией. Регулярное тестирование процедуры восстановления является обязательным.
- Неконтролируемый рост затрат на хранение. При удалении таблиц через DROP TABLE физические данные в S3 не удаляются. Со временем это приведет к накоплению «мертвых» данных и росту счетов. Чтобы справиться с этой ситуацией, мы внедряем автоматизированные пайплайны жизненного цикла данных (Data Lifecycle Management), которые помечают и удаляют orphan-файлы, а также настраиваем политики экспирации старых снапшотов.
- Низкая производительность запросов. Неправильное партиционирование, большое количество мелких файлов или отсутствие компакции убивают производительность. Возможное решение заключается в том, что инженеры проводят постоянный мониторинг и оптимизацию таблиц: запускают процедуры rewrite_data_files для компакции и перераспределения партиций в фоновом режиме.
- Сложность управления доступом и документацией. Iceberg не заменяет собой Data Catalog. Без должного описания наборов данных они быстро превращаются в «болото». Поэтому мы интегрируем Iceberg с корпоративными каталогами данных (например, DataHub или Amundsen), что обеспечивает прозрачность, управление доступом и поиск данных для бизнес-пользователей.
Итак, Apache Iceberg — это не просто технологический тренд, это качественный скачок в построении надежных, масштабируемых и управляемых Data Lake. Он позволяет четко разделить вычислительный слой (Compute) и слой хранения (Storage), предоставляя при этом мощные возможности для управления данными, которые раньше были доступны только в дорогих коммерческих базах данных.
Однако его сила раскрывается только при грамотном внедрении и эксплуатации. Неправильная архитектура, недооценка важности метаданных и отсутствие процессов мониторинга и оптимизации могут свести на нет все преимущества.




