BI Consult Desktop Logo BI Consult Mobile Logo
  • Russian BI Исследование российских bi
  • Перейти на Fine BI
  • Контакты
  • +7 812 334-08-01
    +7 499 608-13-06
  • Отправить сообщение
  • Главная
  • Продукты Эксперт-BI
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Сельское хозяйство
    • Энергетика
    • FMCG
    • Девелоперы
    • Маркетплейсы
    • Пищевая промышленность
    • Фармацевтика
    • Построение Data Platform
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и FP&A
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • IBP
    • ИТ (CIO)
    • Закупки
  • Платформы
    • Системы бизнес-анализа (BI)
    • Интегрированное бизнес-планирование (IBP)
    • Хранилища данных (DWH / Lakehouse)
    • Каталоги данных (Data Catalog)
    • Системы ETL и ELT
    • AI / Исскуственный интеллект
    • Шина данных (ESB)
    • Система управления мастер-данными (MDM)
    • Семантический слой
  • Услуги
    • Переход на отечественные BI и DWH системы
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений и DWH
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Курсы
    • Учебный курс Информационная грамотность (Data Literacy)
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Greenplum
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt (Data Build Tool)
  • Компания
    • Руководство
    • Новости
    • Клиенты
    • Карьера
    • Скачать
    • Контакты

BI

  • FineBI
  • FineReport
  • FineDataLink
  • FineChatBI (FineAI)
  • Коннекторы данных из 1С в BI
  • Airflow / Nifi
  • Visiology
  • PIX BI
  • Modus BI
  • Yandex.DataLens
  • Open-source BI: Superset/Metabase
  • Luxms BI
  • AW BI + Alpha BI
  • FlyBI + Форсайт. Аналитическая Платформа
  • Loginom
  • Триафлай
  • AI / Исскуственный интеллект
  • Optimacros
  • Навигатор BI
  • Семантический слой

СУБД

  • Arenadata
  • ClickHouse
  • Greenplum
  • Postgres Professional
  • TData

Другое

  • Построение Data Platform
    • Аналитическое хранилище данных
    • Data Lake и Data Engineering
    • Подробнее про Data Lake
    • Внедрение Lakehouse
      • Apache Doris
      • StarRocks
      • Trino
    • Миграция витрин из пропиетарных DWH на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс Современная архитектура хранилища данных » Apache Iceberg как фундамент современной Data-инфраструктуры. Подробное руководство по внедрению, лучшим практикам и управлению рисками

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 -lah
total 1.1M
drwxrwxrwt 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; SET s3_url_style = 'path';
 SET s3_endpoint = 'localhost:9000';
 SET s3_access_key_id = 'minioadmin';
 SET s3_secret_access_key = 'minioadmin';
 SET s3_use_ssl = FALSE;

 

Теперь, если выполнить запрос как указано в документации:

SELECT
         * FROM
         iceberg_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
         * FROM
  iceberg_scan('s3://warehouse/default/animals/metadata/00001-f0d7c171-b179-4857-8eec-078f0108c1a9.metadata.json');

 

Но если выполнить INSERT/UPDATE/DELETE, то прошлые мета-даннные будут ссылаться на прошлый снэпшот. Чтобы получить "актуальные" данные, необходимо найти новые мета-данные:

SELECT
         * FROM
  iceberg_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 TABLE foo.bar_partition (
     id bigint,
     created_at timestamp,
 ) USING iceberg PARTITIONED BY (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), предоставляя при этом мощные возможности для управления данными, которые раньше были доступны только в дорогих коммерческих базах данных.

Однако его сила раскрывается только при грамотном внедрении и эксплуатации. Неправильная архитектура, недооценка важности метаданных и отсутствие процессов мониторинга и оптимизации могут свести на нет все преимущества.

Узнать стоимость решенияЗапросить видео презентацию

← Предыдущая статья
Не потеряться в данных: как DataHub спасает аналитику в эру Big Data
Следующая статья →
Подробное руководство по ETL-процессам: слои, маппинги ключевых полей и управление обновлением данных

Решения

Анализировать ФинансыУвеличивайте ПродажиОптимальный Склад и ЛогистикаМаркетинговые Метрики

Клиенты
  • Розничный и интернет-магазин 12 Storeez один из лидеров на рынке женской одежды. С географией рынка не только на территории России, своя продукция представлена еще и в таких странах как Казахстан и Дубай.

  • КАМИ – компания-лидер по поставкам тяжёлых станков в России, занимающаяся продажей и обслуживанием оборудования для обработки металла и дерева, изготовления мебели и не только. На сегодняшний день в компании работают более 1300 человек, запущено 10 обучающих центров, в продаже более 7000 единиц техники. 

  • "Холодильник.ру" - крупнейший в России интернет-магазин бытовой техники и электроники. Компания была основана в 2003 году и за почти 20 лет работы завоевала лидирующие позиции на рынке онлайн ритейла. По данным исследовательского агентства Data Insight, "Холодильник.ру" входит в top-10 крупнейших интернет-магазинов России в категории "электроника и бытовая техника". Компания имеет развитую логистическую инфраструктуру и ежедневно осуществляет более 3500 доставок заказов по всей стране.

  • "Уральский банк реконструкции и развития" входит в топ-25 крупнейших банков России и список значимых кредитных организаций на рынке платежных услуг по версии ЦБ РФ.

  • Решения
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Энергетика
    • Фармацевтика
  • Услуги
    • Переход на отечественные BI и DWH
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Техническая поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Платформы
    • FineBI
    • FineReport
    • FineDataLink
    • Коннекторы данных из 1С в BI
    • Airflow + NiFi
    • Visiology
    • Luxms BI
    • Modus BI
    • PIX BI
    • Arenadata
    • ClickHouse
    • Greenplum
    • Postgres Professional
    • Open-source BI: Superset/Metabase
    • Loginom
    • Yandex.DataLens
    • AI / Исскуственный интеллект
    • Optimacros
    • Шины данных
  • Курсы
    • Учебный курс Информационная грамотность
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt
  • Функциональные решения
    • Создание Data Lake
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и прогнозная аналитика
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • Сквозная аналитика
  • Компания
    • О нас
    • Руководство
    • Новости
    • Клиенты
    • Скачать
    • Контакты
    • Политика конфиденциальности
RutubeVkontakteLinkedInYouTube
ООО "Би Ай Консалт",
ИНН: 7811437757,
ОГРН: 1097847154184
199178, Россия,
Санкт-Петербург,
6-ая линия В.О., Д. 63, 4 этаж
Тел: +7 (812) 334-08-01
Тел: +7 (499) 608-13-06
E-mail: info@biconsult.ru

 

 

 

 

 

×

Пользуясь сайтом, вы соглашаетесь с использованием cookies и политикой конфиденциальности.