Iceberg Lakehouse
Alex Merced RU
Приложение B

Python для Apache Iceberg

На этой странице
  1. B.1 PyIceberg
  2. B.2 Polars
  3. B.3 DuckDB
  4. B.4 Daft
  5. B.5 Dremio
  6. B.6 Bauplan
  7. B.7 Spice AI
  8. B.8 Итоги и лучшие практики

Apache Iceberg стал центральным стандартом современных data lakehouse, а Python даёт одну из самых гибких экосистем для работы с ним. Это приложение знакомит с практическими способами использовать Iceberg напрямую и опосредованно через ведущие библиотеки и фреймворки Python. Каждый раздел посвящён одной библиотеке, объясняет её связь с Iceberg и содержит пошаговые примеры для задач ETL и аналитики.

Цель - показать, как строить, вести и анализировать данные Iceberg целиком на Python, не завися от систем на JVM вроде Spark. Вы научитесь определять схемы, создавать таблицы, добавлять и перезаписывать данные и выполнять запросы с помощью PyIceberg, Polars, DuckDB, Daft, PyDremio, Bauplan и Spice AI.

У каждого инструмента своя роль в экосистеме Python - Iceberg:

  • PyIceberg даёт прямой низкоуровневый доступ к таблицам и каталогам Iceberg.
  • Polars и DuckDB обеспечивают высокопроизводительную аналитику в памяти поверх данных Iceberg.
  • Daft добавляет распределённые вычисления на основе Apache Arrow.
  • Dremio даёт всестороннюю поддержку SQL для Iceberg с масштабируемым выполнением запросов.
  • Bauplan расширяет модель данных Iceberg ветвлением и контролем версий в духе Git.
  • Spice AI обеспечивает федеративные запросы и интеллектуальную аналитику поверх наборов данных Iceberg.

Вместе эти библиотеки показывают, как Python может служить и плоскостью управления, и средой исполнения для lakehouse на Iceberg. Следующие разделы предполагают базовое знакомство с Python и понятиями Iceberg вроде схем, партиций и каталогов. Каждый пример написан так, чтобы его можно было запустить, он конкретен и легко адаптируется под ваши задачи.

B.1 PyIceberg

PyIceberg - официальный клиент Apache Iceberg для Python. Он даёт пользователям Python полный контроль над таблицами Iceberg без необходимости в JVM или Spark. С PyIceberg можно подключиться к каталогу Iceberg, создавать таблицы и управлять ими, выполнять добавление, перезапись и сканирование данных - всё на чистом Python.

Для начала установите библиотеку:

pip install pyiceberg[sql] pyarrow

Такая установка обеспечивает и управление каталогом, и интеграцию с Arrow для чтения и записи данных.

После установки можно подключиться к каталогу Iceberg. В следующем примере используется локальный каталог на SQL поверх SQLite, идеальный для тестирования или автономных конфигураций:

from pyiceberg.catalog import load_catalog
from pyarrow import parquet as pq

catalog = load_catalog(
   "default",
   type="sql",
   uri="sqlite:///demo.db",
   warehouse="file:///tmp/warehouse"
)  #1

catalog.create_namespace("default")  #2

table = catalog.create_table(
   "default.taxi_data",
   schema=pq.read_schema("trips.parquet")
)  #3

arrow_table = pq.read_table("trips.parquet")
table.append(arrow_table)  #4

result = table.scan().to_arrow()  #5

print(len(result))  #6
  1. 1Загружает каталог Iceberg с SQLite в качестве бэкенда. Аргумент warehouse задаёт, где хранятся файлы метаданных и данных.
  2. 2Создаёт логическое пространство имён для группировки связанных таблиц, подобное схеме в базе данных
  3. 3Создаёт в каталоге новую таблицу Iceberg, выводя её схему из существующего файла Parquet
  4. 4Добавляет таблицу PyArrow в таблицу Iceberg, что автоматически создаёт новые файлы данных и обновляет метаданные
  5. 5Читает таблицу Iceberg обратно в память как таблицу PyArrow через API сканирования Iceberg
  6. 6Выводит число полученных строк, подтверждая загрузку данных

Когда таблица создана, PyIceberg позволяет менять её структуру и содержимое. Эволюцию схемы, например, можно выполнить без переписывания старых данных:

update = table.update_schema()
update.add_column("pickup_borough", "string").commit()  #1

new_data = pq.read_table("updated_trips.parquet")
table.append(new_data)  #2
  1. 1Открывает транзакцию обновления схемы, добавляет новый столбец и фиксирует изменение атомарно. Iceberg беспрепятственно обрабатывает эволюцию метаданных.
  2. 2Записывает новые данные, включающие добавленный столбец; старые строки остаются совместимыми

Можно также выполнять перезаписи и добавления с фильтром. Следующий фрагмент заменяет данные конкретной партиции, оставляя остальные нетронутыми:

from pyiceberg.expressions import EqualTo

overwrite = table.overwrite(where=EqualTo("pickup_date", "2023-01-01"))
overwrite.add_file("new_data.parquet").commit()  #1
  1. 1Создаёт транзакционную перезапись, ограниченную предикатом. Заменяются только файлы, соответствующие условию, что обеспечивает атомарность и изолированность.

Наконец, PyIceberg поддерживает сканирование подмножеств данных для аналитических запросов с помощью фильтров:

from pyiceberg.expressions import GreaterThan

expensive_trips = table.scan(
   row_filter=GreaterThan("total_amount", 100)
).to_arrow()  #1

print(expensive_trips.to_pandas().head())  #2
  1. 1Проталкивает фильтр в операцию сканирования, позволяя Iceberg пропускать несоответствующие файлы
  2. 2Преобразует результат в pandas DataFrame для быстрого просмотра или анализа

Как видно из этих примеров, PyIceberg позволяет управлять таблицами Iceberg прямо из Python - от создания таблиц до эволюции схемы и аналитики. Он служит фундаментом для интеграций более высокого уровня вроде Polars, DuckDB и Daft, которые опираются на его API для операций с каталогами и таблицами.

B.2 Polars

Polars стал одной из самых быстрых библиотек датафреймов в Python, предлагая движок запросов, соперничающий по производительности даже с DuckDB. В недавних выпусках Polars получил нативную интеграцию с Iceberg, позволяющую читать таблицы Iceberg напрямую и (экспериментально) писать в них. Это делает Polars естественным выбором для аналитических нагрузок поверх данных Iceberg, сочетая эффективные векторизованные вычисления с ленивым выполнением запросов.

Прежде чем работать с Iceberg, убедитесь, что установлены Polars и PyIceberg, поскольку Polars делегирует работу с каталогом Iceberg библиотеке PyIceberg:

pip install polars pyiceberg pyarrow

После установки можно подключиться к существующей таблице Iceberg и выполнять запросы через ленивый API Polars. Следующий пример читает таблицу Iceberg прямо из её файла метаданных:

import polars as pl

lazy_df = pl.scan_iceberg(
   "s3://data-lake/warehouse/trips/metadata.json",
   storage_options={"AWS_ACCESS_KEY_ID": "key", "AWS_SECRET_ACCESS_KEY": "secret"}
)  #1

df = lazy_df.collect()  #2

print(df.describe())  #3
  1. 1Инициализирует ленивое сканирование таблицы Iceberg, указывая на её metadata.json. Polars по метаданным Iceberg обнаруживает все файлы данных Parquet.
  2. 2Выполняет ленивый запрос, загружая в память только нужные столбцы и партиции
  3. 3Печатает сводку по таблице, показывая, насколько беспрепятственно Iceberg встраивается в модель датафреймов Polars

После загрузки данных Polars даёт эффективные преобразования и аналитические операции. В следующем примере данные Iceberg фильтруются и агрегируются для расчёта средней дистанции поездки по числу пассажиров:

summary = (
   lazy_df
   .filter(pl.col("passenger_count") > 1)  #1
   .group_by("passenger_count")            #2
   .agg(pl.col("trip_distance").mean().alias("avg_distance"))  #3
   .sort("passenger_count")                #4
   .collect()                              #5
)

print(summary)
  1. 1Проталкивает предикат в сканирование Iceberg, обеспечивая чтение только нужных файлов
  2. 2Группирует строки по столбцу с числом пассажиров для агрегации
  3. 3Вычисляет среднюю дистанцию поездки для каждого числа пассажиров
  4. 4Упорядочивает результаты ради опрятного вывода
  5. 5Выполняет ленивый план, эффективно запуская все чтения из Iceberg и вычисления

Возможности записи в Iceberg в Polars ещё экспериментальны, но уже дают простой API для добавления или перезаписи данных. Следующий пример показывает запись DataFrame из Polars в таблицу Iceberg:

# Create a small Polars DataFrame
new_data = pl.DataFrame({
   "trip_id": [101, 102, 103],
   "passenger_count": [1, 2, 3],
   "trip_distance": [3.5, 5.1, 2.9]
})  #1

# Append this data to an existing Iceberg table
new_data.write_iceberg(
   "default.trips_table",
   mode="append"
)  #2
  1. 1Определяет новый DataFrame в Polars для записи в таблицу Iceberg
  2. 2Добавляет DataFrame в указанную таблицу. Polars делегирует запись PyIceberg, обеспечивая гарантии ACID от Iceberg.

Сочетание ленивых вычислений, поколоночной оптимизации и нативной поддержки Iceberg делает Polars мощным выбором для аналитических нагрузок. С его помощью можно исследовать данные, выполнять агрегации и прототипировать преобразования, прежде чем записывать результаты обратно в Iceberg для постоянного хранения. Операции записи ещё развиваются, а вот производительность чтения и вычислений уже готова к промышленной эксплуатации, давая быстрый и питонический опыт работы с таблицами Apache Iceberg.

B.3 DuckDB

DuckDB - аналитический движок баз данных, работающий внутри процесса и обеспечивающий молниеносные SQL-запросы прямо из Python. Расширяемая архитектура позволяет ему читать и записывать таблицы Iceberg через расширение Iceberg, обеспечивая беспрепятственный доступ к данным Iceberg и для ad hoc аналитики, и для задач ETL. С DuckDB можно обращаться к таблицам Iceberg, хранящимся локально или в облаке, вставлять новые данные и даже подключаться к каталогам Iceberg для управляемых операций с метаданными.

Прежде чем работать с Iceberg, нужно установить DuckDB и включить расширение Iceberg:

pip install duckdb

После установки можно инициализировать соединение и загрузить расширение Iceberg, чтобы функции Iceberg стали доступны в DuckDB:

import duckdb

con = duckdb.connect("iceberg_demo.db")  #1

con.execute("INSTALL iceberg;")  #2

con.execute("LOAD iceberg;")  #3
  1. 1Устанавливает постоянное соединение DuckDB, способное запрашивать таблицы Iceberg и хранить локальные метаданные
  2. 2Устанавливает расширение Iceberg, добавляющее специфичные для Iceberg функции вроде iceberg_scan
  3. 3Активирует расширение для сессии, немедленно делая операции Iceberg доступными

После загрузки расширения можно обращаться к существующей таблице Iceberg по её файлу метаданных. Следующий пример читает данные прямо из файла metadata.json таблицы:

result = con.execute("""
   SELECT vendor_id, trip_distance, total_amount
   FROM iceberg_scan('/data/warehouse/trips/metadata.json')
   WHERE trip_distance > 5
   ORDER BY total_amount DESC
   LIMIT 10
""").df()  #1

print(result.head())  #2
  1. 1Выполняет SQL-запрос к таблице Iceberg через iceberg_scan, которая читает метаданные таблицы и извлекает подходящие под фильтр файлы Parquet
  2. 2Показывает верхние строки результата, уже загруженные как pandas DataFrame

DuckDB также поддерживает интеграцию с каталогами Iceberg, включая REST-каталоги вроде Apache Polaris или AWS Glue. После подключения DuckDB можно запрашивать таблицы Iceberg и управлять ими привычными командами SQL DDL и DML:

con.execute("""
   ATTACH 'my_catalog' AS iceberg_catalog (
       TYPE iceberg,
       ENDPOINT 'https://catalog-host.com/v1',
       SECRET 'my_secret_key'
   );
""")  #1

tables = con.execute("SHOW TABLES FROM iceberg_catalog.default;").fetchall()  #2

con.execute("""
   CREATE TABLE iceberg_catalog.default.daily_trips AS
   SELECT * FROM iceberg_catalog.default.trips WHERE trip_date = '2023-01-01';
""")  #3
  1. 1Подключает каталог, чтобы DuckDB разрешал имена таблиц без прямых путей
  2. 2Перечисляет таблицы, зарегистрированные в пространстве имён default подключённого каталога
  3. 3Создаёт новую таблицу Iceberg прямо через каталог выражением CREATE TABLE AS SELECT

Можно также вставлять новые данные в существующие таблицы Iceberg. DuckDB прозрачно записывает данные Parquet и обновляет метаданные Iceberg через расширение:

con.execute("""
   INSERT INTO iceberg_catalog.default.daily_trips
   VALUES (1001, 3.2, 15.75), (1002, 7.1, 28.10);
""")  #1
  1. 1Добавляет новые строки в таблицу Iceberg средствами SQL. DuckDB фиксирует данные атомарно, обновляя манифест и файлы метаданных таблицы.

Для аналитических задач можно выполнять соединения, агрегации и запросы с путешествием во времени к таблицам Iceberg, как если бы это были обычные SQL-таблицы. Расширение Iceberg в DuckDB поддерживает чтение снимков по идентификатору или метке времени, обеспечивая воспроизводимые запросы к историческим состояниям таблицы:

# Query a specific snapshot using a snapshot ID
historical = con.execute("""
   SELECT COUNT(*) AS trips_count
   FROM iceberg_scan(
       '/data/warehouse/trips/metadata.json',
       snapshot_id := '9182736455'
   );
""").df()  #1

print(historical)
  1. 1Читает данные из предыдущего снимка таблицы Iceberg, обеспечивая путешествия во времени ради воспроизводимой аналитики

Интеграция Iceberg в DuckDB делает его одним из самых удобных движков и для интерактивных запросов, и для скриптовых ETL-процессов. Можно подключать каталоги, выполнять преобразования и выгружать аналитические результаты с простотой SQL, сохраняя при этом транзакционную согласованность Iceberg и гарантии эволюции схемы.

B.4 Daft

Daft - библиотека распределённых датафреймов, построенная на Apache Arrow и оптимизированная под крупномасштабную обработку данных. Она даёт привычный API в духе pandas, обеспечивая при этом высокопроизводительный ввод-вывод и выполнение запросов. Интеграция Daft с Apache Iceberg работает через PyIceberg, позволяя читать таблицы Iceberg прямо в датафреймы Daft для анализа. Сейчас Daft ориентирован скорее на чтение данных Iceberg, чем на запись, но отлично справляется с аналитическими нагрузками, требующими сканирования и преобразования больших наборов данных Iceberg.

Для начала установите Daft вместе с PyIceberg и нужными зависимостями для хранилища:

pip install daft pyiceberg pyarrow s3fs

После установки можно читать таблицы Iceberg через каталог PyIceberg. Каталог предоставляет метаданные таблиц и пути в объектном хранилище, а Daft берёт на себя извлечение данных и вычисления:

import daft
from pyiceberg.catalog import load_catalog

catalog = load_catalog(
   "default",
   type="sql",
   uri="sqlite:///demo.db",
   warehouse="file:///tmp/warehouse"
)  #1

table = catalog.load_table("default.taxi_data")  #2

df = daft.read_iceberg(table)  #3

df.show(5)  #4
  1. 1Подключается к каталогу Iceberg на SQL, хранящему определения таблиц и метаданные
  2. 2Загружает существующую таблицу Iceberg из каталога
  3. 3Читает таблицу в DataFrame Daft через интеграцию с PyIceberg. Daft выполняет параллельное чтение файлов Parquet таблицы.
  4. 4Показывает первые строки для быстрого просмотра, подобно head() в pandas

После загрузки данных можно выполнять распределённые преобразования и фильтрацию через API выражений Daft. Daft проталкивает фильтры в сканирование Iceberg, минимизируя лишний ввод-вывод за счёт чтения только нужных файлов:

filtered = (
   df.filter(df["total_amount"] > 100)        #1
     .groupby("passenger_count")              #2
     .agg({"trip_distance": ["mean"]})        #3
)

filtered.show()                                #4
  1. 1Применяет выражение фильтра, которое Iceberg использует для отсечения партиций и пропуска файлов
  2. 2Группирует строки по числу пассажиров для агрегированной аналитики
  3. 3Вычисляет среднее по столбцу trip_distance внутри каждой группы
  4. 4Показывает агрегированные результаты, вычисленные распределённо

Тесная связь Daft с Arrow обеспечивает обмен данными без копирования с другими инструментами. Например, можно преобразовать DataFrame из Daft в таблицы Arrow или датафреймы Polars для дальнейшей обработки:

# Convert Daft DataFrame to Arrow Table
arrow_table = df.to_arrow_table()  #1

# Convert Arrow Table to Polars DataFrame
import polars as pl
polars_df = pl.from_arrow(arrow_table)  #2
  1. 1Экспортирует DataFrame из Daft в таблицу PyArrow ради совместимости
  2. 2Преобразует таблицу Arrow в DataFrame из Polars, обеспечивая конвейеры на смеси библиотек для аналитики и визуализации

Daft пока не поддерживает прямую запись в таблицы Iceberg, но обработанные данные можно выгрузить в Parquet и зарегистрировать их через PyIceberg. Такой двухшаговый процесс позволяет Daft выступать высокопроизводительным слоем вычислений в конвейерах данных на Iceberg:

# Write Daft DataFrame to Parquet files
df.write_parquet("s3://data-lake/cleaned/trips/")  #1

# Register exported data as a new Iceberg table
catalog.create_table(
   "default.cleaned_trips",
   schema=df.schema().to_arrow(),
   location="s3://data-lake/cleaned/trips/"
)  #2
  1. 1Записывает DataFrame из Daft в файлы Parquet в объектном хранилище
  2. 2Создаёт в каталоге новую таблицу Iceberg, ссылающуюся на эти файлы, что позволит обращаться к ним из совместимых с Iceberg движков

На практике интеграция Daft с Apache Iceberg даёт современный питонический подход к эффективному чтению больших наборов данных. Она особенно полезна там, где данные лежат в таблицах Iceberg, а аналитическая логика выполняется на Python. Сочетая возможности PyIceberg по работе с метаданными и распределённый движок исполнения Daft, разработчики могут интерактивно обрабатывать огромные наборы данных Iceberg и передавать результаты в нижестоящие системы или визуализации, оставаясь в единой среде Python.

B.5 Dremio

Dremio даёт полноценный SQL-движок для Apache Iceberg, а пользователи Python могут обращаться к его возможностям через API и коннекторы Dremio. С их помощью можно запрашивать, создавать и изменять таблицы Iceberg, пользуясь планировщиком запросов Dremio, рефлексиями и возможностями федерации данных. Два самых распространённых способа работы с Dremio из Python - через драйвер ODBC и интерфейс Apache Arrow Flight. Оба позволяют выполнять SQL-команды над таблицами Iceberg, как если бы это были обычные таблицы базы данных.

Для начала установите необходимые зависимости:

pip install pyodbc pyarrow

После установки можно подключиться к кластеру Dremio и выполнять SQL-команды прямо из Python. Следующий пример показывает запрос к таблице Iceberg через интерфейс ODBC:

import pyodbc

conn = pyodbc.connect(
   "Driver={Dremio Connector 64-bit};"
   "ConnectionType=Direct;"
   "HOST=demo.dremio.cloud;"
   "PORT=443;"
   "AuthenticationType=Plain;"
   "UID=my_username;"
   "PWD=my_token;"
   "SSL=1"
)  #1

cursor = conn.cursor()  #2

cursor.execute('SELECT * FROM "Samples"."NYC_Trips" LIMIT 5;')  #3

rows = cursor.fetchall()  #4
for row in rows:
   print(row)
  1. 1Создаёт прямое подключение к Dremio через драйвер ODBC и персональный токен доступа
  2. 2Создаёт курсор для отправки SQL-запросов движку Dremio
  3. 3Выполняет SQL-запрос к таблице Iceberg, хранящейся в каталоге Dremio
  4. 4Забирает и печатает результаты запроса, показывая, что к таблицам Iceberg можно обращаться как к любому реляционному источнику

Способ через ODBC хорош для лёгких приложений, а Arrow Flight даёт более быструю и эффективную передачу больших наборов данных в поколоночном формате Arrow. Следующий пример показывает запрос к таблице Iceberg через Arrow Flight:

from pyarrow import flight

client = flight.FlightClient("grpc+tls://data.dremio.cloud:443")  #1

headers = [(b"authorization", f"bearer {token}".encode())]  #2

sql = 'SELECT trip_distance, total_amount FROM "Samples"."NYC_Trips" WHERE total_amount > 100'  #3

flight_info = client.get_flight_info(
   flight.FlightDescriptor.for_command(sql),
   flight.FlightCallOptions(headers=headers)
)  #4

reader = client.do_get(flight_info.endpoints[0].ticket, flight.FlightCallOptions(headers=headers))  #5
arrow_table = reader.read_all()  #6

print(arrow_table.to_pandas().head())  #7
  1. 1Подключается к эндпоинту Arrow Flight в Dremio для высокопроизводительной передачи данных
  2. 2Задаёт заголовки аутентификации с bearer-токеном для защищённого доступа
  3. 3Определяет SQL-запрос, выбирающий и фильтрующий данные из таблицы Iceberg
  4. 4Запрашивает у Dremio сведения о выполнении запроса, получая метаданные о конечных точках результата
  5. 5Передаёт данные потоком через эффективный транспорт gRPC в Flight
  6. 6Читает весь набор данных в память как таблицу Arrow для дальнейшей обработки в Python
  7. 7Преобразует таблицу Arrow в pandas DataFrame для анализа или визуализации

Помимо запросов, через Dremio можно создавать и изменять таблицы Iceberg. Следующий пример показывает создание новой таблицы Iceberg и вставку в неё данных:

cursor.execute("""
   CREATE TABLE "Samples"."daily_summary" AS
   SELECT passenger_count, AVG(total_amount) AS avg_total
   FROM "Samples"."NYC_Trips"
   GROUP BY passenger_count;
""")  #1

cursor.execute("""
   INSERT INTO "Samples"."daily_summary"
   VALUES (4, 32.75), (5, 45.12);
""")  #2

conn.commit()  #3
  1. 1Создаёт новую таблицу Iceberg в Dremio, выбирая данные из другой таблицы Iceberg
  2. 2Вставляет новые строки в таблицу стандартным синтаксисом SQL
  3. 3Фиксирует транзакцию, чтобы изменения отразились в метаданных и файлах манифестов Iceberg

Через любой из интерфейсов Dremio поддерживает полный DDL и DML для таблиц Iceberg, включая операции CREATE, INSERT, DELETE, MERGE и UPDATE. Это делает Dremio удачной плоскостью управления data lake на Iceberg с сохранением совместимости со стеком data science в Python. Сочетая производительность PyArrow с интеграцией Iceberg в Dremio, разработчики могут вести масштабируемую аналитику на SQL поверх открытых данных, не перемещая и не дублируя наборы данных.

B.6 Bauplan

Bauplan - современная платформа данных, созданная для управления data lake на Iceberg с контролем версий. Она привносит в работу с данными процессы в духе Git: можно создавать ветки, экспериментировать с преобразованиями, проверять результаты и сливать изменения - и всё это опирается на метаданные и структуру хранения Apache Iceberg. Любая операция в Bauplan - создание ветки, импорт таблицы, эволюция схемы, слияние - в конечном счёте работает с таблицами Iceberg под капотом, что делает её мощным слоем абстракции для инженерии данных на Iceberg.

Для начала установите SDK Bauplan:

pip install bauplan

После установки можно пройти аутентификацию и инициализировать клиентское соединение. Клиент даёт доступ к вашим репозиториям на Iceberg и полный контроль над операциями с данными:

import bauplan

client = bauplan.Client(api_key="my_api_key")  #1

client.create_branch("dev", from_ref="main")  #2

client.create_namespace("default", branch="dev")  #3
  1. 1Инициализирует клиент Bauplan для связи с сервисом Bauplan или локальной средой
  2. 2Создаёт новую ветку «dev» на основе «main», обеспечивая изолированные эксперименты
  3. 3Определяет логическое пространство имён для группировки связанных таблиц Iceberg, подобное схеме или базе данных

Когда ветка и пространство имён готовы, можно создавать и загружать таблицы Iceberg. Внутри Bauplan использует PyIceberg для регистрации таблиц и управления их схемами и метаданными:

client.create_table(
   table="trips",
   search_uri="s3://data-lake/trips/*.parquet",
   namespace="default",
   branch="dev",
   replace=True
)  #1

client.import_data(
   table="trips",
   search_uri="s3://data-lake/trips/*.parquet",
   namespace="default",
   branch="dev"
)  #2

arrow_table = client.scan(
   table="trips",
   ref="dev",
   namespace="default",
   columns=["passenger_count", "total_amount"]
)  #3

print(arrow_table.to_pandas().head())  #4
  1. 1Создаёт новую таблицу Iceberg в ветке dev, сканируя и регистрируя данные Parquet как файлы Iceberg
  2. 2Импортирует или заменяет данные в таблице; каждая операция порождает новый снимок Iceberg
  3. 3Читает таблицу Iceberg как таблицу PyArrow, что позволяет вести аналитику в pandas или Polars
  4. 4Преобразует таблицу Arrow в pandas DataFrame для быстрого просмотра

Bauplan также поддерживает процессы проверки и публикации, повторяющие модель ветвления Git. После того как данные проверены или преобразованы на ветке, их можно слить обратно в main для промышленного использования.

assert arrow_table["passenger_count"].null_count == 0  #1

client.merge_branch(source_ref="dev", into_branch="main")  #2

client.delete_branch("dev")  #3
  1. 1Проверяет, что данные отвечают требованиям качества, перед публикацией
  2. 2Сливает ветку dev в main, атомарно фиксируя все изменения таблиц Iceberg
  3. 3Прибирается, удаляя временную ветку после успешного слияния

Поскольку операции Bauplan транзакционны и нативны для Iceberg, каждое изменение - импорт, перезапись или обновление схемы - создаёт новый снимок в нижележащей таблице Iceberg. Это обеспечивает полную прослеживаемость, проверяемость и возможность отката.

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

client.update_schema(
   table="trips",
   namespace="default",
   branch="dev",
   add_columns=[{"name": "pickup_zone", "type": "string"}]
)  #1
  1. 1Изменяет схему таблицы, добавляя новый столбец в контексте ветки и сохраняя историческую совместимость между снимками

Bauplan предлагает мощную модель управления lakehouse на Iceberg: она относится к данным как к коду. Вы можете ветвиться, фиксировать изменения, проверять и сливать с той же дисциплиной, что и в разработке ПО, сохраняя при этом все транзакционные гарантии Iceberg. Такой подход позволяет инженерам данных безопасно экспериментировать с преобразованиями, выполнять проверки качества и публиковать проверенные данные в продакшен, не ломая и не дублируя наборы данных.

B.7 Spice AI

Spice AI - современный движок запросов, обеспечивающий федеративный доступ ко множеству источников данных, включая Apache Iceberg. Он рассматривает таблицы Iceberg как полноценные наборы данных, которые можно запрашивать наравне с другими источниками вроде PostgreSQL, BigQuery или объектных хранилищ. Spice AI сочетает открытые стандарты - Arrow, Flight SQL и REST - с табличным форматом Iceberg, обеспечивая высокопроизводительные аналитические запросы через простой интерфейс Python.

Чтобы использовать Spice AI, установите его SDK для Python:

pip install spicepy

После установки можно настроить набор данных, указывающий на ваш каталог Iceberg. Spice AI использует REST-эндпоинт каталога, чтобы обнаруживать таблицы Iceberg и предоставлять их через единый слой SQL. После настройки с этими таблицами можно работать как с любым другим источником данных SQL.

Следующий пример показывает запрос к таблице Iceberg через клиент Spice AI для Python:

from spicepy import Client

client = Client()  #1

result = client.query("""
   SELECT pickup_borough, AVG(total_amount) AS avg_fare
   FROM iceberg.default.trips
   WHERE total_amount > 50
   GROUP BY pickup_borough
   ORDER BY avg_fare DESC
""")  #2

df = result.read_pandas()  #3

print(df.head())  #4
  1. 1Создаёт новый экземпляр клиента Spice AI, взаимодействующий со средой выполнения Spice или облачным сервисом
  2. 2Выполняет SQL-запрос к набору данных Iceberg, зарегистрированному в каталоге Spice. Запрос исполняется федеративным движком, использующим REST API Iceberg под капотом.
  3. 3Читает результаты запроса в pandas DataFrame для анализа или визуализации
  4. 4Показывает предварительный вид агрегированных результатов для быстрого просмотра

Помимо простых запросов, Spice AI поддерживает сложные аналитические операции между источниками. Таблицы Iceberg можно соединять с внешними наборами данных ради анализа сразу по нескольким доменам без дублирования данных:

joined = client.query("""
   SELECT i.trip_id, i.total_amount, p.driver_rating
   FROM iceberg.default.trips AS i
   JOIN postgres.analytics.driver_ratings AS p
   ON i.driver_id = p.driver_id
   WHERE i.total_amount > 100
""")  #1

joined_df = joined.read_pandas()  #2

print(joined_df.describe())  #3
  1. 1Выполняет федеративный запрос, соединяющий таблицу Iceberg с таблицей PostgreSQL, что показывает работу Spice AI с несколькими источниками
  2. 2Читает результат в pandas DataFrame, сохраняя типы столбцов и метаданные
  3. 3Печатает сводную статистику, показывая, насколько беспрепятственно данные Iceberg встраиваются в более широкий аналитический контекст

Таблицы Iceberg можно регистрировать динамически через конфигурацию или вызовы API, что упрощает подключение новых наборов данных по мере развития вашего lakehouse. Следующий пример регистрирует таблицу Iceberg через интерфейс конфигурации наборов данных:

datasets:
 - name: trips
   from: iceberg:https://catalog.mycompany.com/v1/namespaces/default/tables/trips  #1
  1. 1Объявляет набор данных Iceberg, ссылаясь на его REST-эндпоинт каталога, что позволяет Spice AI управлять им в рамках единого слоя запросов

Spice AI не изменяет данные Iceberg напрямую, но даёт мощный слой запросов и интеграции. Все операции записи и обслуживания - вставки, удаления, обновления схемы - выполняются инструментами каталога Iceberg или движками вроде Dremio и PyIceberg. Фокус Spice AI - быстрая федеративная аналитика и выводы на основе ИИ поверх ваших таблиц Iceberg и других структурированных источников.

Встроив Iceberg в свою федеративную архитектуру, Spice AI позволяет разработчикам на Python анализировать данные в открытом формате с минимальной настройкой, опираясь и на надёжность данных в Iceberg, и на единую модель запросов Spice. Это делает его удачным выбором для организаций, строящих готовые к ИИ lakehouse, которым нужен гибкий доступ к управляемым данным в Iceberg.

B.8 Итоги и лучшие практики

В этом приложении вы изучили, как Python взаимодействует с Apache Iceberg через ряд библиотек, у каждой из которых свои сильные стороны, абстракции и уровень интеграции. PyIceberg даёт прямой базовый интерфейс к операциям с таблицами Iceberg, а фреймворки более высокого уровня - Polars, DuckDB, Daft, Dremio, Bauplan и Spice AI - надстраиваются над этими принципами, предлагая более богатые аналитические и оркестрационные возможности. Понимание того, как эффективно сочетать эти инструменты, - ключ к поддержанию чистой, масштабируемой и эффективной платформы данных на Iceberg.

В основе любого процесса работы с Iceberg на Python лежит один и тот же набор операций: создание, чтение, запись, эволюция схем и эффективное сканирование данных. Следующий код иллюстрирует лучшие практики и паттерны, проступающие во всех этих интеграциях:

from pyiceberg.catalog import load_catalog
from pyarrow import parquet as pq

# Initialize a local Iceberg catalog for testing
catalog = load_catalog(
   "default",
   type="sql",
   uri="sqlite:///demo.db",
   warehouse="file:///tmp/warehouse"
)  #1

# Create a new Iceberg table from a Parquet dataset
table = catalog.create_table(
   "default.sales",
   schema=pq.read_schema("sales_data.parquet")
)  #2

# Read Parquet file as Arrow Table and append to Iceberg
data = pq.read_table("sales_data.parquet")
table.append(data)  #3

# Perform incremental overwrite for a given partition
from pyiceberg.expressions import EqualTo
overwrite = table.overwrite(where=EqualTo("region", "US"))  #4
overwrite.add_file("new_sales_data.parquet").commit()  #5

# Query the table with a filter for analytics
from pyiceberg.expressions import GreaterThan
results = table.scan(row_filter=GreaterThan("amount", 1000)).to_arrow()  #6
print(results.to_pandas().head())  #7
  1. 1Подключается к локальному каталогу Iceberg, управляющему метаданными и регистрацией таблиц
  2. 2Создаёт новую таблицу Iceberg по схеме, выведенной из файла Parquet
  3. 3Добавляет данные новыми файлами, оставляя атомарные коммиты на усмотрение Iceberg
  4. 4Определяет условную перезапись, нацеленную только на партицию «US»
  5. 5Фиксирует транзакцию перезаписи, порождая новый снимок
  6. 6Выполняет сканирование с фильтром, эффективно извлекая крупные транзакции
  7. 7Преобразует таблицу Arrow в pandas DataFrame для просмотра

Этот пример охватывает жизненный цикл, которому следуют все инструменты Iceberg на Python, - выполняются ли операции напрямую через PyIceberg или абстрагированы другими движками.

Во всех рассмотренных библиотеках несколько практик стабильно ведут к более надёжным и эффективным процессам работы с Iceberg:

  • Последовательно используйте каталоги. Всегда регистрируйте и запрашивайте таблицы Iceberg через каталог. Это обеспечивает синхронность метаданных и сохранение транзакционных гарантий даже при работе нескольких вычислительных движков.
  • Разумно применяйте эволюцию схемы. Добавлять столбцы в Iceberg легко, но планируйте обратную совместимость. Всегда развивайте схемы на тестовой ветке (в Bauplan) или в каталоге разработки, прежде чем сливать в продакшен.
  • Оптимизируйте запись заранее. Пишите данные файлами подходящего размера и применяйте стратегии партиционирования, отражающие типичные фильтры запросов. Сделанное на раннем этапе снижает потребность в тяжёлых задачах оптимизации позже.
  • Предпочитайте ленивые сканирования с проталкиванием предикатов. Библиотеки вроде Polars, Daft и DuckDB проталкивают фильтры в слой метаданных Iceberg, что минимизирует ввод-вывод и ускоряет запросы. Всегда проектируйте запросы так, чтобы пользоваться этими оптимизациями.
  • Принимайте эксперименты на ветках. Модель ветвления Bauplan идеальна для изоляции экспериментальных изменений. Относитесь к веткам как к промежуточным средам для преобразований данных перед слиянием чистых данных в промышленные таблицы.
  • Объедините доступ для аналитики и ИИ. Используйте федеративные движки вроде Dremio или Spice AI, когда нескольким командам нужен общий доступ к таблицам Iceberg в облаке и локальной инфраструктуре. Такой подход централизует управление данными, сохраняя открытый доступ.

Устройство Apache Iceberg позволяет этим разнородным инструментам беспрепятственно взаимодействовать. Управляете ли вы приёмом данных, оптимизируете запросы или обеспечиваете работу ИИ-приложений - сочетание Python и Iceberg даёт и гибкость, и контроль. Вместе они образуют мощный фундамент открытой, версионируемой и полностью управляемой экосистемы lakehouse.

Обновлено 26.07.2026